- migration 6: manual_records, ledger_event_revisions chain, current projections, source claims, subject suggestions, eligible_position_events - ledger_events.py: bank-event reconciliation, reversal/adjustment/reopen, append-only revision chain and rebuildable current projection - subjects.py: fixed subject mirror, draft suggestion dictionary, explicit administrator subject confirmation with expected_revision + idempotency - manual_records.py: submit, approve new/link, return/exception/reverse, candidate hints, idempotent replay and concurrency-safe claims - positions.py: Decimal aggregation, both-perspective conservation asserts, cutoff window, unresolved gross buckets, keyset pagination, evidence visibility (visible/masked/missing) - server.py: admin + company intercompany APIs with tenant isolation (404 on cross-tenant reads, 403 on company writes) and auto reconcile wiring - admin/company portals: balance directory, pair drill-down drawer, evidence drawer, subject/manual audit queue, company balance summary - tests: ledger events, subjects, manual records, positions, HTTP API and migration persistence (233 total, all green)
685 lines
25 KiB
Python
685 lines
25 KiB
Python
"""Canonical intercompany ledger events and the revision chain (B-44).
|
|
|
|
Bank source rows and approved manual records are immutable evidence. This
|
|
module turns them into one canonical ledger event each through ``reconcile_*``
|
|
functions and an append-only revision chain. Corrections are never in-place
|
|
edits: a reversal or adjustment is a new ledger event with its own effective
|
|
date, and the original event keeps its history so no earlier cutoff is
|
|
rewritten. Current revisions and source claims are rebuildable projections.
|
|
|
|
Only B-43 ``eligible_intercompany_events`` feeds bank facts here; same-company
|
|
transfers, external transactions, unresolved rows and unlocked single
|
|
observations never reach the confirmed balance.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from decimal import Decimal, InvalidOperation
|
|
import json
|
|
import sqlite3
|
|
|
|
from .db import utc_now
|
|
from .subjects import MIRROR, SUBJECTS, mirror_subject
|
|
|
|
|
|
class LedgerConflictError(ValueError):
|
|
"""A revision/claim/idempotency conflict (mapped to HTTP 409)."""
|
|
|
|
|
|
class LedgerInputError(ValueError):
|
|
"""Invalid input for a ledger operation (mapped to HTTP 400/422)."""
|
|
|
|
|
|
SUBJECT_RULE_VERSION = "subject-suggest-draft-v1"
|
|
|
|
|
|
def amount_scale(amount: object) -> int:
|
|
"""Decimal places of a decimal-string amount, never negative."""
|
|
try:
|
|
exponent = Decimal(str(amount)).as_tuple().exponent
|
|
except InvalidOperation:
|
|
return 0
|
|
return max(0, -int(exponent))
|
|
|
|
|
|
def parse_amount(amount: object) -> Decimal:
|
|
"""Parse a positive, valid decimal-string amount."""
|
|
try:
|
|
value = Decimal(str(amount))
|
|
except InvalidOperation:
|
|
raise LedgerInputError("金额不是有效的十进制数。") from None
|
|
if not value.is_finite() or value <= 0:
|
|
raise LedgerInputError("金额必须大于零。")
|
|
return value
|
|
|
|
|
|
def _company_exists(connection: sqlite3.Connection, company_id: int, label: str) -> None:
|
|
row = connection.execute(
|
|
"SELECT id FROM companies WHERE id = ?", (company_id,)
|
|
).fetchone()
|
|
if row is None:
|
|
raise LedgerInputError(f"{label}指向的公司不存在。")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Revision helpers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _ensure_transaction(connection: sqlite3.Connection) -> bool:
|
|
"""Begin an immediate transaction unless one is already open.
|
|
|
|
Write helpers may run standalone (they own the transaction) or nested
|
|
inside a caller's transaction (e.g. the sheet-confirm flow); nested calls
|
|
never start their own commit.
|
|
"""
|
|
began = False
|
|
if not connection.in_transaction:
|
|
connection.execute("BEGIN IMMEDIATE")
|
|
began = True
|
|
return began
|
|
|
|
|
|
def current_revision(connection: sqlite3.Connection, ledger_event_id: int) -> sqlite3.Row | None:
|
|
return connection.execute(
|
|
"""
|
|
SELECT r.* FROM current_ledger_event_revisions c
|
|
JOIN ledger_event_revisions r ON r.id = c.revision_id
|
|
WHERE c.ledger_event_id = ?
|
|
""",
|
|
(ledger_event_id,),
|
|
).fetchone()
|
|
|
|
|
|
def _event_lifecycle(connection: sqlite3.Connection, ledger_event_id: int) -> str | None:
|
|
row = connection.execute(
|
|
"SELECT lifecycle FROM ledger_events WHERE id = ?", (ledger_event_id,)
|
|
).fetchone()
|
|
return row["lifecycle"] if row is not None else None
|
|
|
|
|
|
def _next_revision_number(connection: sqlite3.Connection, ledger_event_id: int) -> int:
|
|
row = connection.execute(
|
|
"SELECT COALESCE(MAX(revision), 0) AS m FROM ledger_event_revisions WHERE ledger_event_id = ?",
|
|
(ledger_event_id,),
|
|
).fetchone()
|
|
return int(row["m"]) + 1
|
|
|
|
|
|
def create_event(
|
|
connection: sqlite3.Connection,
|
|
*,
|
|
state: str,
|
|
effective_at: str,
|
|
amount: str,
|
|
currency: str,
|
|
payer_company_id: int,
|
|
payee_company_id: int,
|
|
perspective_company_id: int | None,
|
|
subject_code: str | None,
|
|
source_kind: str,
|
|
source_revision_token: str | None,
|
|
posting_kind: str,
|
|
reverses_ledger_event_id: int | None = None,
|
|
adjusts_ledger_event_id: int | None = None,
|
|
rule_version: str | None = None,
|
|
evidence_json: str | None = None,
|
|
idempotency_key: str | None = None,
|
|
actor: sqlite3.Row | None = None,
|
|
reason: str | None = None,
|
|
supersedes_revision_id: int | None = None,
|
|
) -> tuple[int, int]:
|
|
"""Insert a new ledger event with one revision. Returns ``(event_id, revision_id)``."""
|
|
if state == "confirmed":
|
|
if perspective_company_id is None or subject_code is None:
|
|
raise LedgerInputError("已确认事件必须提供视角公司与科目。")
|
|
if subject_code not in SUBJECTS:
|
|
raise LedgerInputError("科目必须是应收/应付/其他应收/其他应付之一。")
|
|
if int(payer_company_id) == int(payee_company_id):
|
|
raise LedgerInputError("付款公司与收款公司不能相同。")
|
|
if perspective_company_id is not None and perspective_company_id not in (
|
|
int(payer_company_id), int(payee_company_id),
|
|
):
|
|
raise LedgerInputError("视角公司必须是事件参与方。")
|
|
began = _ensure_transaction(connection)
|
|
try:
|
|
now = utc_now()
|
|
cursor = connection.execute(
|
|
"INSERT INTO ledger_events (lifecycle, created_at) VALUES ('active', ?)",
|
|
(now,),
|
|
)
|
|
event_id = int(cursor.lastrowid)
|
|
revision_id = append_revision(
|
|
connection,
|
|
event_id,
|
|
state=state,
|
|
effective_at=effective_at,
|
|
amount=amount,
|
|
currency=currency,
|
|
payer_company_id=payer_company_id,
|
|
payee_company_id=payee_company_id,
|
|
perspective_company_id=perspective_company_id,
|
|
subject_code=subject_code,
|
|
source_kind=source_kind,
|
|
source_revision_token=source_revision_token,
|
|
posting_kind=posting_kind,
|
|
reverses_ledger_event_id=reverses_ledger_event_id,
|
|
adjusts_ledger_event_id=adjusts_ledger_event_id,
|
|
rule_version=rule_version,
|
|
evidence_json=evidence_json,
|
|
idempotency_key=idempotency_key,
|
|
actor=actor,
|
|
reason=reason,
|
|
supersedes_revision_id=supersedes_revision_id,
|
|
)
|
|
except Exception:
|
|
if began:
|
|
connection.rollback()
|
|
raise
|
|
else:
|
|
if began:
|
|
connection.commit()
|
|
return event_id, revision_id
|
|
|
|
|
|
def append_revision(
|
|
connection: sqlite3.Connection,
|
|
ledger_event_id: int,
|
|
*,
|
|
state: str,
|
|
effective_at: str,
|
|
amount: str,
|
|
currency: str,
|
|
payer_company_id: int,
|
|
payee_company_id: int,
|
|
perspective_company_id: int | None,
|
|
subject_code: str | None,
|
|
source_kind: str,
|
|
source_revision_token: str | None,
|
|
posting_kind: str,
|
|
reverses_ledger_event_id: int | None = None,
|
|
adjusts_ledger_event_id: int | None = None,
|
|
rule_version: str | None = None,
|
|
evidence_json: str | None = None,
|
|
idempotency_key: str | None = None,
|
|
actor: sqlite3.Row | None = None,
|
|
reason: str | None = None,
|
|
supersedes_revision_id: int | None = None,
|
|
) -> int:
|
|
"""Append one immutable revision and repoint the current projection."""
|
|
if _event_lifecycle(connection, ledger_event_id) != "active":
|
|
raise LedgerConflictError("该事件已停用,不能追加修订。")
|
|
if state == "confirmed":
|
|
if perspective_company_id is None or subject_code is None:
|
|
raise LedgerInputError("已确认事件必须提供视角公司与科目。")
|
|
if subject_code not in SUBJECTS:
|
|
raise LedgerInputError("科目必须是应收/应付/其他应收/其他应付之一。")
|
|
if perspective_company_id not in (int(payer_company_id), int(payee_company_id)):
|
|
raise LedgerInputError("视角公司必须是事件参与方。")
|
|
elif state != "pending_subject":
|
|
raise LedgerInputError("事件状态必须是 pending_subject 或 confirmed。")
|
|
revision = _next_revision_number(connection, ledger_event_id)
|
|
now = utc_now()
|
|
cursor = connection.execute(
|
|
"""
|
|
INSERT INTO ledger_event_revisions (
|
|
ledger_event_id, revision, state, effective_at, amount,
|
|
amount_scale, currency, payer_company_id, payee_company_id,
|
|
perspective_company_id, subject_code, source_kind,
|
|
source_revision_token, posting_kind, reverses_ledger_event_id,
|
|
adjusts_ledger_event_id, rule_version, evidence_json,
|
|
idempotency_key, actor_user_id, actor_username, reason,
|
|
supersedes_revision_id, created_at
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
""",
|
|
(
|
|
ledger_event_id, revision, state, effective_at, amount,
|
|
amount_scale(amount), currency, payer_company_id, payee_company_id,
|
|
perspective_company_id, subject_code, source_kind,
|
|
source_revision_token, posting_kind, reverses_ledger_event_id,
|
|
adjusts_ledger_event_id, rule_version, evidence_json,
|
|
idempotency_key,
|
|
actor["id"] if actor is not None else None,
|
|
actor["username"] if actor is not None else None,
|
|
reason, supersedes_revision_id, now,
|
|
),
|
|
)
|
|
connection.execute(
|
|
"""
|
|
INSERT OR REPLACE INTO current_ledger_event_revisions (ledger_event_id, revision_id)
|
|
VALUES (?, ?)
|
|
""",
|
|
(ledger_event_id, cursor.lastrowid),
|
|
)
|
|
return int(cursor.lastrowid)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Source claims
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def bank_source_claim(connection: sqlite3.Connection, bank_event_id: int) -> sqlite3.Row | None:
|
|
return connection.execute(
|
|
"SELECT * FROM ledger_event_bank_sources WHERE bank_event_id = ?",
|
|
(bank_event_id,),
|
|
).fetchone()
|
|
|
|
|
|
def manual_source_claim(connection: sqlite3.Connection, manual_record_id: int) -> sqlite3.Row | None:
|
|
return connection.execute(
|
|
"SELECT * FROM ledger_event_manual_sources WHERE manual_record_id = ?",
|
|
(manual_record_id,),
|
|
).fetchone()
|
|
|
|
|
|
def _facts_of(connection: sqlite3.Connection, ledger_event_id: int) -> dict[str, object]:
|
|
revision = current_revision(connection, ledger_event_id)
|
|
if revision is None:
|
|
return {}
|
|
return {
|
|
"amount": revision["amount"],
|
|
"currency": revision["currency"],
|
|
"effective_at": revision["effective_at"],
|
|
"payer_company_id": revision["payer_company_id"],
|
|
"payee_company_id": revision["payee_company_id"],
|
|
}
|
|
|
|
|
|
def _eligible_facts(event: sqlite3.Row) -> dict[str, object]:
|
|
return {
|
|
"amount": event["amount"],
|
|
"currency": event["currency"],
|
|
"effective_at": event["effective_at"],
|
|
"payer_company_id": event["payer_company_id"],
|
|
"payee_company_id": event["payee_company_id"],
|
|
}
|
|
|
|
|
|
def _has_reversal(connection: sqlite3.Connection, original_event_id: int) -> bool:
|
|
row = connection.execute(
|
|
"""
|
|
SELECT 1 FROM ledger_event_revisions r
|
|
JOIN current_ledger_event_revisions c ON c.revision_id = r.id
|
|
WHERE r.reverses_ledger_event_id = ? AND r.posting_kind = 'reversal'
|
|
LIMIT 1
|
|
""",
|
|
(original_event_id,),
|
|
).fetchone()
|
|
return row is not None
|
|
|
|
|
|
def create_reversal(
|
|
connection: sqlite3.Connection,
|
|
original_event_id: int,
|
|
*,
|
|
source_kind: str,
|
|
source_revision_token: str | None = None,
|
|
effective_at: str | None = None,
|
|
reason: str,
|
|
actor: sqlite3.Row | None,
|
|
idempotency_key: str | None = None,
|
|
rule_version: str | None = None,
|
|
) -> tuple[int, int]:
|
|
"""Create an equal-amount, opposite-direction reversal as a new ledger event.
|
|
|
|
The subject mirrors the original (应收<->应付, 其他应收<->其他应付). ``effective_at``
|
|
defaults to the original event's effective date so an earlier cutoff keeps
|
|
the original impact and later cutoffs see the net zero. The original event
|
|
is never modified or deleted.
|
|
"""
|
|
original = current_revision(connection, original_event_id)
|
|
if original is None:
|
|
raise LedgerConflictError("原事件不存在或没有当前修订。")
|
|
if original["state"] != "confirmed":
|
|
raise LedgerInputError("只有已确认事件才能生成冲销。")
|
|
perspective = mirror_perspective(original)
|
|
if effective_at is None:
|
|
effective_at = original["effective_at"]
|
|
return create_event(
|
|
connection,
|
|
state="confirmed",
|
|
effective_at=effective_at,
|
|
amount=original["amount"],
|
|
currency=original["currency"],
|
|
payer_company_id=original["payee_company_id"],
|
|
payee_company_id=original["payer_company_id"],
|
|
perspective_company_id=perspective,
|
|
subject_code=mirror_subject(original["subject_code"]),
|
|
source_kind=source_kind,
|
|
source_revision_token=source_revision_token,
|
|
posting_kind="reversal",
|
|
reverses_ledger_event_id=original_event_id,
|
|
rule_version=rule_version or original["rule_version"],
|
|
idempotency_key=idempotency_key,
|
|
actor=actor,
|
|
reason=reason,
|
|
)
|
|
|
|
|
|
def mirror_perspective(revision: sqlite3.Row) -> int:
|
|
"""The counterparty company from ``revision``'s perspective."""
|
|
perspective = int(revision["perspective_company_id"])
|
|
if perspective == int(revision["payer_company_id"]):
|
|
return int(revision["payee_company_id"])
|
|
return int(revision["payer_company_id"])
|
|
|
|
|
|
def create_adjustment(
|
|
connection: sqlite3.Connection,
|
|
ledger_event_id: int,
|
|
*,
|
|
effective_at: str,
|
|
amount: str,
|
|
currency: str,
|
|
payer_company_id: int,
|
|
payee_company_id: int,
|
|
perspective_company_id: int,
|
|
subject_code: str,
|
|
reason: str,
|
|
actor: sqlite3.Row,
|
|
idempotency_key: str | None = None,
|
|
rule_version: str | None = None,
|
|
) -> tuple[int, int]:
|
|
"""Create an audit adjustment event; the original event stays unchanged."""
|
|
return create_event(
|
|
connection,
|
|
state="confirmed",
|
|
effective_at=effective_at,
|
|
amount=str(parse_amount(amount)),
|
|
currency=currency,
|
|
payer_company_id=payer_company_id,
|
|
payee_company_id=payee_company_id,
|
|
perspective_company_id=perspective_company_id,
|
|
subject_code=subject_code,
|
|
source_kind="adjustment",
|
|
source_revision_token=None,
|
|
posting_kind="adjustment",
|
|
adjusts_ledger_event_id=ledger_event_id,
|
|
rule_version=rule_version or SUBJECT_RULE_VERSION,
|
|
idempotency_key=idempotency_key,
|
|
actor=actor,
|
|
reason=reason,
|
|
)
|
|
|
|
|
|
def reopen_subject(
|
|
connection: sqlite3.Connection,
|
|
ledger_event_id: int,
|
|
*,
|
|
reason: str,
|
|
actor: sqlite3.Row,
|
|
idempotency_key: str | None = None,
|
|
) -> tuple[int, int]:
|
|
"""Reverse a confirmed event and re-open it for subject re-review.
|
|
|
|
Creates an equal-amount reversal plus a fresh ``pending_subject`` event
|
|
that re-claims the original bank source, so the administrator can confirm
|
|
a corrected subject. The original event and its reversal keep history.
|
|
"""
|
|
current = current_revision(connection, ledger_event_id)
|
|
if current is None or current["state"] != "confirmed":
|
|
raise LedgerConflictError("只有已确认事件可以重新进入科目审核。")
|
|
bank_claim = connection.execute(
|
|
"SELECT * FROM ledger_event_bank_sources WHERE ledger_event_id = ?",
|
|
(ledger_event_id,),
|
|
).fetchone()
|
|
if bank_claim is None:
|
|
raise LedgerInputError(
|
|
"该事件没有银行来源,无法重新进入科目审核;请改用调整或冲销。"
|
|
)
|
|
if not _has_reversal(connection, ledger_event_id):
|
|
create_reversal(
|
|
connection, ledger_event_id,
|
|
source_kind=current["source_kind"],
|
|
source_revision_token=current["source_revision_token"],
|
|
reason="科目复核:原确认事件冲销",
|
|
actor=actor,
|
|
idempotency_key=(idempotency_key + ":rev" if idempotency_key else None),
|
|
rule_version=current["rule_version"],
|
|
)
|
|
ev = connection.execute(
|
|
"SELECT * FROM eligible_intercompany_events WHERE event_id = ?",
|
|
(bank_claim["bank_event_id"],),
|
|
).fetchone()
|
|
if ev is None:
|
|
raise LedgerInputError("银行事件已不再纳入往来,无法重新入账。")
|
|
event_id, revision_id = _create_bank_event(
|
|
connection, ev, actor, reason="科目复核后重新入账,待确认科目",
|
|
replacing_claim=bank_claim,
|
|
)
|
|
return event_id, revision_id
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Bank event reconciliation
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def reconcile_bank_events(
|
|
connection: sqlite3.Connection, actor: sqlite3.Row | None = None
|
|
) -> dict[str, object]:
|
|
"""Reconcile the current eligible intercompany events into ledger events.
|
|
|
|
Idempotent: first sight creates a ``pending_subject`` event; a changed B-43
|
|
decision updates a still-pending event's revision, or (for a confirmed
|
|
event) creates a reversal plus a fresh pending event. A source that left
|
|
the eligible set with a confirmed impact gets one reversal. Runs inside the
|
|
caller's transaction when one is open, otherwise in its own transaction.
|
|
"""
|
|
began = _ensure_transaction(connection)
|
|
try:
|
|
eligible = {
|
|
row["event_id"]: row
|
|
for row in connection.execute(
|
|
"SELECT * FROM eligible_intercompany_events"
|
|
).fetchall()
|
|
}
|
|
claims = {
|
|
row["bank_event_id"]: row
|
|
for row in connection.execute(
|
|
"SELECT * FROM ledger_event_bank_sources"
|
|
).fetchall()
|
|
}
|
|
stats = {
|
|
"created": 0, "updated_pending": 0, "reversal": 0,
|
|
"reopened": 0, "unchanged": 0, "sources": len(eligible),
|
|
}
|
|
for bank_event_id, event in sorted(eligible.items()):
|
|
claim = claims.get(bank_event_id)
|
|
if claim is None:
|
|
_create_bank_event(
|
|
connection, event, actor, reason="B-43 事件首次入账,待确认科目"
|
|
)
|
|
stats["created"] += 1
|
|
continue
|
|
current = current_revision(connection, claim["ledger_event_id"])
|
|
if current is None or _facts_of(connection, claim["ledger_event_id"]) != _eligible_facts(event):
|
|
if current is not None and current["state"] == "confirmed":
|
|
if not _has_reversal(connection, claim["ledger_event_id"]):
|
|
create_reversal(
|
|
connection, claim["ledger_event_id"],
|
|
source_kind="bank",
|
|
source_revision_token=event["decision_id"],
|
|
reason="B-43 事件事实变更,原确认事件冲销",
|
|
actor=actor,
|
|
)
|
|
stats["reversal"] += 1
|
|
_create_bank_event(
|
|
connection, event, actor,
|
|
reason="B-43 事件事实变更后重新入账,待确认科目",
|
|
replacing_claim=claim,
|
|
)
|
|
stats["reopened"] += 1
|
|
elif current is None or current["state"] == "pending_subject":
|
|
_append_bank_pending_revision(
|
|
connection, claim["ledger_event_id"], event, actor
|
|
)
|
|
stats["updated_pending"] += 1
|
|
else:
|
|
stats["unchanged"] += 1
|
|
else:
|
|
stats["unchanged"] += 1
|
|
|
|
for bank_event_id, claim in sorted(claims.items()):
|
|
if bank_event_id in eligible:
|
|
continue
|
|
current = current_revision(connection, claim["ledger_event_id"])
|
|
if current is not None and current["state"] == "confirmed":
|
|
if not _has_reversal(connection, claim["ledger_event_id"]):
|
|
create_reversal(
|
|
connection, claim["ledger_event_id"],
|
|
source_kind="bank",
|
|
source_revision_token=None,
|
|
reason="B-43 事件不再纳入往来,原确认事件冲销",
|
|
actor=actor,
|
|
)
|
|
stats["reversal"] += 1
|
|
except Exception:
|
|
if began:
|
|
connection.rollback()
|
|
raise
|
|
else:
|
|
if began:
|
|
connection.commit()
|
|
return stats
|
|
|
|
|
|
def _create_bank_event(
|
|
connection: sqlite3.Connection,
|
|
event: sqlite3.Row,
|
|
actor: sqlite3.Row | None,
|
|
*,
|
|
reason: str,
|
|
replacing_claim: sqlite3.Row | None = None,
|
|
) -> tuple[int, int]:
|
|
event_id, revision_id = create_event(
|
|
connection,
|
|
state="pending_subject",
|
|
effective_at=event["effective_at"],
|
|
amount=event["amount"],
|
|
currency=event["currency"],
|
|
payer_company_id=event["payer_company_id"],
|
|
payee_company_id=event["payee_company_id"],
|
|
perspective_company_id=None,
|
|
subject_code=None,
|
|
source_kind="bank",
|
|
source_revision_token=event["decision_id"],
|
|
posting_kind="normal",
|
|
rule_version=SUBJECT_RULE_VERSION,
|
|
evidence_json=json.dumps(
|
|
{
|
|
"bank_event_id": event["event_id"],
|
|
"decision_id": event["decision_id"],
|
|
"pairing": event["pairing"],
|
|
"evidence_count": event["evidence_count"],
|
|
},
|
|
ensure_ascii=False,
|
|
),
|
|
actor=actor,
|
|
reason=reason,
|
|
)
|
|
if replacing_claim is not None:
|
|
connection.execute(
|
|
"""
|
|
UPDATE ledger_event_bank_sources SET ledger_event_id = ?
|
|
WHERE bank_event_id = ?
|
|
""",
|
|
(event_id, event["event_id"]),
|
|
)
|
|
else:
|
|
connection.execute(
|
|
"""
|
|
INSERT INTO ledger_event_bank_sources (bank_event_id, ledger_event_id)
|
|
VALUES (?, ?)
|
|
""",
|
|
(event["event_id"], event_id),
|
|
)
|
|
from .subjects import store_suggestions
|
|
|
|
store_suggestions(connection, event_id)
|
|
return event_id, revision_id
|
|
|
|
|
|
def _append_bank_pending_revision(
|
|
connection: sqlite3.Connection,
|
|
ledger_event_id: int,
|
|
event: sqlite3.Row,
|
|
actor: sqlite3.Row | None,
|
|
) -> int:
|
|
current = current_revision(connection, ledger_event_id)
|
|
revision_id = append_revision(
|
|
connection,
|
|
ledger_event_id,
|
|
state="pending_subject",
|
|
effective_at=event["effective_at"],
|
|
amount=event["amount"],
|
|
currency=event["currency"],
|
|
payer_company_id=event["payer_company_id"],
|
|
payee_company_id=event["payee_company_id"],
|
|
perspective_company_id=None,
|
|
subject_code=None,
|
|
source_kind="bank",
|
|
source_revision_token=event["decision_id"],
|
|
posting_kind="normal",
|
|
rule_version=SUBJECT_RULE_VERSION,
|
|
evidence_json=json.dumps(
|
|
{
|
|
"bank_event_id": event["event_id"],
|
|
"decision_id": event["decision_id"],
|
|
"pairing": event["pairing"],
|
|
"evidence_count": event["evidence_count"],
|
|
},
|
|
ensure_ascii=False,
|
|
),
|
|
actor=actor,
|
|
reason="B-43 事件事实更新,追加待审修订",
|
|
supersedes_revision_id=current["id"] if current is not None else None,
|
|
)
|
|
from .subjects import store_suggestions
|
|
|
|
store_suggestions(connection, ledger_event_id)
|
|
return revision_id
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Projection rebuild
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def rebuild_current_ledger_projection(connection: sqlite3.Connection) -> int:
|
|
"""Rebuild current ledger revisions from the append-only log."""
|
|
began = _ensure_transaction(connection)
|
|
try:
|
|
connection.execute("DELETE FROM current_ledger_event_revisions")
|
|
rows = connection.execute(
|
|
"""
|
|
SELECT e.id AS ledger_event_id,
|
|
(SELECT r2.id FROM ledger_event_revisions r2
|
|
WHERE r2.ledger_event_id = e.id
|
|
ORDER BY r2.revision DESC LIMIT 1) AS latest_id
|
|
FROM ledger_events e
|
|
WHERE e.lifecycle = 'active'
|
|
""",
|
|
).fetchall()
|
|
rebuilt = 0
|
|
for row in rows:
|
|
if row["latest_id"] is None:
|
|
continue
|
|
connection.execute(
|
|
"""
|
|
INSERT OR REPLACE INTO current_ledger_event_revisions (ledger_event_id, revision_id)
|
|
VALUES (?, ?)
|
|
""",
|
|
(row["ledger_event_id"], row["latest_id"]),
|
|
)
|
|
rebuilt += 1
|
|
except Exception:
|
|
if began:
|
|
connection.rollback()
|
|
raise
|
|
else:
|
|
if began:
|
|
connection.commit()
|
|
return rebuilt
|