往来查询
公司间往来余额目录
公司余额目录
每行余额都附带截止日、期初状态、本期借贷、结果与未决金额
diff --git a/docs/decisions/006-intercompany-positions.md b/docs/decisions/006-intercompany-positions.md new file mode 100644 index 0000000..9323c37 --- /dev/null +++ b/docs/decisions/006-intercompany-positions.md @@ -0,0 +1,87 @@ +# 006: Intercompany positions, subject review and drill-down evidence (B-44) + +Status: accepted (2026-08-19) + +## Decision + +The intercompany ledger adds one canonical layer above B-43's eligible events: + +- **Append-only ledger events.** A canonical `ledger_event` carries an + append-only `ledger_event_revisions` chain. A B-43 eligible bank event first + becomes a `pending_subject` revision; an administrator confirms the subject + into a `confirmed` revision. 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 earlier cutoffs are not + rewritten. +- **Manual records as immutable submitted facts.** Companies submit + `manual_records`; only an administrator-approved record becomes a confirmed + ledger event (`approve_new`) or joins one (`approve_link`). Returned, + exception and pending records never affect a balance and never leak to the + counterparty. Approved facts change only through `reverse` (a new opposite + event) — the original is never edited. +- **One perspective, fixed mirror.** Subjects are stored from one participating + company's perspective (`receivable/payable/other_receivable/other_payable`); + the other side is the fixed mirror (应收<->应付, 其他应收<->其他应付), so the + two companies can never book conflicting subjects. +- **Subjects are confirmed, never auto-posted.** Bank summary/purpose text only + feeds a deterministic *suggestion* dictionary (`subject-suggest-draft-v1`, + not group-approved). Without an approved trade dictionary every bank event + stays in subject review until an administrator confirms. +- **Decimal-only aggregation.** All money is stored as TEXT decimal strings and + aggregated with Python `Decimal`. SQLite `SUM`, JavaScript `Number` and + Python `float` never touch financial math. Different currencies are + aggregated and displayed separately; nothing is converted to a group total. +- **Conservation is asserted.** For every company pair and currency the two + perspectives must mirror exactly (`C_A == -C_B`); a violation raises a + calculation exception instead of rendering an unbalanced number. +- **Unresolved is absolute gross.** Unresolved amounts are summed by absolute + value per currency (never netted), broken down by reason + (`subject_review`, `unmatched_single`, `manual_pending`), with count and + gross amount exposed on every balance response. +- **B-45 boundary.** Until the B-45 opening balance exists, every response + returns `opening.status=unavailable`, `opening.amount=null` and + `result.kind=period_net_change`; the UI labels this "期间净变动", never + "期末余额". + +## Background + +B-43 produces `eligible_intercompany_events` as the only bank-event entry +point. Before B-44 there was no canonical financial event, no statutory +subject, no manual-record approval, and no server-side balance API. The +revision-chain design is inherited from `transfer_match_decisions` in +migration 5 and from the append-only audit posture of the rest of the system. + +## Consequences + +- **Positive:** balances are deterministic, conservable, auditable and + drillable from a group directory down to bank source rows; manual records + cannot double count; corrections never mutate evidence; company portals are + tenant-scoped on the server. +- **Negative:** pending bank events and returned/exception manual records are + intentionally invisible to counterparties, which can surprise cashiers who + expect symmetric disclosure; subject confirmation is manual until a + group-approved dictionary exists. +- **Operational:** migration 6 is forward-only for production once + approvals/revisions exist; pre-production it can be rolled back with + `--rollback-to 5`. Read aggregation runs against the rebuildable current + projection (no day snapshots yet); if B-45 monthly close needs them, immutable + monthly snapshots can be added behind the same API contract. + +## Files + +- `src/bank_importer/ledger_events.py` — event lifecycle, revision chain, + bank reconciliation, reversal/adjustment/reopen, projection rebuild. +- `src/bank_importer/subjects.py` — subject constants/mirror, suggestion + dictionary, `confirm_subject`. +- `src/bank_importer/manual_records.py` — submission, approval + (new/link), return/exception/reverse, candidate hints, idempotency. +- `src/bank_importer/positions.py` — Decimal aggregation, directories, pairs, + events, evidence visibility, unresolved buckets, keyset pagination. +- `server.py` — `/api/admin/intercompany/*` and `/api/company/intercompany/*` + plus `/api/admin/subject-reviews` and `/api/admin/manual-records`. +- `db.py` migration 6 — `manual_records`, `manual_record_decisions`, + `ledger_events`, `ledger_event_revisions`, current-pointer projections, + source claim tables, `ledger_subject_suggestions`, `eligible_position_events`. +- Tests: `test_ledger_events.py`, `test_manual_records.py`, + `test_positions.py`, `test_positions_api.py`, extended + `test_persistence.py`. diff --git a/server.py b/server.py index bf89568..71110e9 100644 --- a/server.py +++ b/server.py @@ -10,7 +10,10 @@ from http.cookies import SimpleCookie from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer from urllib.parse import parse_qs, urlparse -from bank_importer import auth, importing, master_data, matching, multipart, personal_transit +from bank_importer import ( + auth, importing, ledger_events, manual_records, master_data, matching, + multipart, personal_transit, positions, subjects, +) from bank_importer.db import connect, migrate, utc_now @@ -101,6 +104,68 @@ class AppHandler(SimpleHTTPRequestHandler): self._handle_company_transfer_event_detail(int(company_event_match.group(1))) return + # B-44 intercompany positions (admin) + if path == "/api/admin/intercompany/balances": + self._handle_admin_intercompany_balances(query) + return + admin_pair = re.fullmatch( + r"/api/admin/intercompany/pairs/(\d+)/(\d+)", path + ) + if admin_pair: + self._handle_admin_intercompany_pair( + int(admin_pair.group(1)), int(admin_pair.group(2)), query + ) + return + if path == "/api/admin/intercompany/events": + self._handle_admin_intercompany_events(query) + return + if path == "/api/admin/subject-reviews": + self._handle_admin_subject_reviews(query) + return + if path == "/api/admin/manual-records": + self._handle_admin_manual_records(query) + return + admin_event_match = re.fullmatch(r"/api/admin/intercompany/events/(\d+)", path) + if admin_event_match: + self._handle_admin_intercompany_event_detail(int(admin_event_match.group(1))) + return + admin_evidence = re.fullmatch( + r"/api/admin/intercompany/events/(\d+)/evidence", path + ) + if admin_evidence: + self._handle_admin_intercompany_evidence(int(admin_evidence.group(1))) + return + + # B-44 intercompany positions (company, own-company scope) + if path == "/api/company/intercompany/balances": + self._handle_company_intercompany_balances(query) + return + company_pair = re.fullmatch(r"/api/company/intercompany/pairs/(\d+)", path) + if company_pair: + self._handle_company_intercompany_pair(int(company_pair.group(1)), query) + return + if path == "/api/company/intercompany/events": + self._handle_company_intercompany_events(query) + return + if path == "/api/company/manual-records": + self._handle_company_manual_records(query) + return + if path == "/api/company/companies": + self._handle_company_companies() + return + company_event_match = re.fullmatch(r"/api/company/intercompany/events/(\d+)", path) + if company_event_match: + self._handle_company_intercompany_event_detail( + int(company_event_match.group(1)) + ) + return + company_evidence = re.fullmatch( + r"/api/company/intercompany/events/(\d+)/evidence", path + ) + if company_evidence: + self._handle_company_intercompany_evidence(int(company_evidence.group(1))) + return + if path == "/admin.html" and not self._guard_page("admin"): return if path == "/company.html" and not self._guard_page("company"): @@ -167,6 +232,29 @@ class AppHandler(SimpleHTTPRequestHandler): if mapping_review: self._handle_admin_review_personal_mapping(int(mapping_review.group(1))) return + + # B-44 intercompany positions (admin writes) + subject_decision = re.fullmatch( + r"/api/admin/intercompany/events/(\d+)/subject-decisions", path + ) + if subject_decision: + self._handle_admin_subject_decision(int(subject_decision.group(1))) + return + adjustment = re.fullmatch( + r"/api/admin/intercompany/events/(\d+)/adjustments", path + ) + if adjustment: + self._handle_admin_intercompany_adjustment(int(adjustment.group(1))) + return + manual_decision = re.fullmatch(r"/api/admin/manual-records/(\d+)/decisions", path) + if manual_decision: + self._handle_admin_manual_record_decision(int(manual_decision.group(1))) + return + + # B-44 intercompany positions (company writes) + if path == "/api/company/manual-records": + self._handle_company_manual_records_submit() + return self._send_json(404, {"status": "error", "message": "接口不存在。"}) # ------------------------------------------------------------------ @@ -1611,6 +1699,11 @@ class AppHandler(SimpleHTTPRequestHandler): except Exception as exc: self._send_json(500, {"status": "error", "message": f"重跑匹配失败:{exc}"}) return + try: + ledger_events.reconcile_bank_events(connection, actor=user) + except Exception as exc: + self._send_json(500, {"status": "error", "message": f"同步往来事件失败:{exc}"}) + return auth.audit( connection, "transfer_reconcile", actor=user, target=f"rows:{len(row_ids)}", @@ -1664,6 +1757,11 @@ class AppHandler(SimpleHTTPRequestHandler): except matching.MatchInputError as exc: self._send_json(400, {"status": "error", "message": str(exc)}) return + try: + ledger_events.reconcile_bank_events(connection, actor=user) + except Exception as exc: + self._send_json(500, {"status": "error", "message": f"同步往来事件失败:{exc}"}) + return self._send_json(200, {"status": "ok", "decision": payload}) finally: connection.close() @@ -1974,6 +2072,736 @@ class AppHandler(SimpleHTTPRequestHandler): return rows + # ------------------------------------------------------------------ + # B-44 intercompany positions (admin) + # ------------------------------------------------------------------ + + def _parse_window(self, query: dict[str, list[str]]) -> tuple[str, str]: + raw_from = (query.get("from") or [None])[0] + raw_cutoff = (query.get("cutoff") or [None])[0] + if raw_cutoff is None: + raw_cutoff = positions.today_shanghai() + try: + return positions.validate_window(raw_from, raw_cutoff) + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + raise + + def _parse_limit(self, query: dict[str, list[str]], default: int = 50) -> int: + raw = (query.get("limit") or [str(default)])[0] + try: + return max(1, min(int(raw), 200)) + except ValueError: + return default + + def _parse_int(self, query: dict[str, list[str]], key: str, label: str) -> int | None: + raw = (query.get(key) or [None])[0] + if not raw: + return None + try: + return int(raw) + except ValueError: + self._send_json(400, {"status": "error", "message": f"{label}参数无效。"}) + raise + + def _handle_admin_intercompany_balances(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + try: + from_, cutoff = self._parse_window(query) + company_id = self._parse_int(query, "company_id", "company_id") + currency = (query.get("currency") or [None])[0] + limit = self._parse_limit(query) + cursor = (query.get("cursor") or [None])[0] + except Exception: + return + try: + payload = positions.company_balances( + connection, from_=from_, cutoff=cutoff, currency=currency, + company_id=company_id, limit=limit, cursor=cursor, + ) + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_admin_intercompany_pair( + self, company_a: int, company_b: int, query: dict[str, list[str]] + ) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + if not self._companies_exist(connection, (company_a, company_b)): + self._send_json(404, {"status": "error", "message": "公司不存在。"}) + return + try: + from_, cutoff = self._parse_window(query) + currency = (query.get("currency") or [None])[0] + except Exception: + return + try: + payload = positions.pair_detail( + connection, company_a, company_b, + from_=from_, cutoff=cutoff, currency=currency, + ) + except positions.PositionError as exc: + self._send_json(422, {"status": "error", "message": str(exc)}) + return + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _companies_exist(self, connection, ids: tuple[int, ...]) -> bool: + for company_id in ids: + row = connection.execute( + "SELECT id FROM companies WHERE id = ?", (company_id,) + ).fetchone() + if row is None: + return False + return True + + def _handle_admin_intercompany_events(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + try: + from_, cutoff = self._parse_window(query) + currency = (query.get("currency") or [None])[0] + company_id = self._parse_int(query, "company_id", "company_id") + company_a = self._parse_int(query, "company_a", "company_a") + company_b = self._parse_int(query, "company_b", "company_b") + subject = (query.get("subject") or [None])[0] + state = (query.get("state") or [None])[0] + posting_kind = (query.get("posting_kind") or [None])[0] + source_kind = (query.get("source_kind") or [None])[0] + limit = self._parse_limit(query) + cursor = (query.get("cursor") or [None])[0] + except Exception: + return + try: + payload = positions.list_events( + connection, from_=from_, cutoff=cutoff, currency=currency, + company_id=company_id, company_a=company_a, company_b=company_b, + subject=subject, state=state, posting_kind=posting_kind, + source_kind=source_kind, limit=limit, cursor=cursor, + ) + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_admin_subject_reviews(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + try: + from_, cutoff = self._parse_window(query) + company_id = self._parse_int(query, "company_id", "company_id") + limit = self._parse_limit(query) + cursor = (query.get("cursor") or [None])[0] + except Exception: + return + try: + payload = positions.subject_review_queue( + connection, from_=from_, cutoff=cutoff, + company_id=company_id, limit=limit, cursor=cursor, + ) + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_admin_manual_records(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + company_id = self._parse_int(query, "company_id", "company_id") + if company_id is None and (query.get("company_id") or [None])[0]: + return + state = (query.get("state") or [None])[0] + limit = self._parse_limit(query, 100) + try: + rows = manual_records.list_records( + connection, company_id=company_id, state=state, limit=limit + ) + except manual_records.ManualInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + items = [] + for row in rows: + payload = manual_records._row_payload(connection, row) + items.append(payload) + self._send_json(200, {"status": "ok", "records": items}) + finally: + connection.close() + + def _handle_admin_intercompany_event_detail(self, event_id: int) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + payload = positions.event_detail(connection, event_id) + if payload is None: + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_admin_intercompany_evidence(self, event_id: int) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + payload = positions.event_evidence(connection, event_id) + if payload is None: + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_admin_subject_decision(self, event_id: int) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + data = self._read_json_body() + if data is None: + return + try: + perspective_company_id = int(data["perspective_company_id"]) + expected_revision = ( + int(data["expected_revision"]) + if data.get("expected_revision") is not None + else None + ) + except (KeyError, TypeError, ValueError): + self._send_json( + 400, {"status": "error", "message": "perspective_company_id 或 expected_revision 无效。"} + ) + return + try: + action = str(data.get("action") or "confirm") + if action in ("return", "exception"): + payload = subjects.park_subject( + connection, + event_id, + disposition=action, + reason=str(data.get("reason") or ""), + expected_revision=expected_revision, + request_key=str(data.get("request_key") or "") or None, + actor=user, + ) + else: + payload = subjects.confirm_subject( + connection, + event_id, + perspective_company_id=perspective_company_id, + subject_code=str(data.get("subject_code") or ""), + reason=str(data.get("reason") or ""), + expected_revision=expected_revision, + request_key=str(data.get("request_key") or "") or None, + actor=user, + ) + except subjects.SubjectConflictError as exc: + self._send_json(409, {"status": "error", "message": str(exc)}) + return + except subjects.SubjectInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + auth.audit( + connection, "subject_confirm", actor=user, + target=f"ledger_event:{event_id}", + detail=f"revision:{payload['revision']};subject:{payload['subject_code']}", + ip=self._client_ip, + ) + self._send_json(200, {"status": "ok", "revision": payload}) + finally: + connection.close() + + def _handle_admin_intercompany_adjustment(self, event_id: int) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + data = self._read_json_body() + if data is None: + return + action = str(data.get("action") or "") + reason = str(data.get("reason") or "") + request_key = str(data.get("request_key") or "") or None + if not reason: + self._send_json(400, {"status": "error", "message": "必须填写操作原因。"}) + return + try: + if action == "reverse": + event_id, _revision_id = ledger_events.create_reversal( + connection, event_id, + source_kind="adjustment", + source_revision_token=None, + effective_at=data.get("effective_at") or None, + reason=reason, actor=user, idempotency_key=request_key, + ) + outcome: dict[str, object] = { + "action": "reverse", "ledger_event_id": event_id, + } + elif action == "adjust": + try: + effective_at = str(data.get("effective_at") or "") + amount = str(data.get("amount") or "") + currency = str(data.get("currency") or "") + payer = int(data["payer_company_id"]) + payee = int(data["payee_company_id"]) + perspective = int(data["perspective_company_id"]) + subject_code = str(data.get("subject_code") or "") + except (KeyError, TypeError, ValueError): + self._send_json( + 400, {"status": "error", "message": "adjust 参数不完整或无效。"} + ) + return + event_id, _revision_id = ledger_events.create_adjustment( + connection, event_id, + effective_at=effective_at, amount=amount, currency=currency, + payer_company_id=payer, payee_company_id=payee, + perspective_company_id=perspective, subject_code=subject_code, + reason=reason, actor=user, idempotency_key=request_key, + ) + outcome = {"action": "adjust", "ledger_event_id": event_id} + elif action == "reopen": + event_id, _revision_id = ledger_events.reopen_subject( + connection, event_id, reason=reason, actor=user, + idempotency_key=request_key, + ) + outcome = {"action": "reopen", "ledger_event_id": event_id} + else: + self._send_json( + 400, + {"status": "error", "message": "action 必须是 reverse、adjust 或 reopen。"}, + ) + return + except ledger_events.LedgerConflictError as exc: + self._send_json(409, {"status": "error", "message": str(exc)}) + return + except ledger_events.LedgerInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + auth.audit( + connection, f"ledger_{action}", actor=user, + target=f"ledger_event:{event_id}", detail=reason, ip=self._client_ip, + ) + self._send_json(200, {"status": "ok", **outcome}) + finally: + connection.close() + + def _handle_admin_manual_record_decision(self, record_id: int) -> None: + connection = connect(DB_PATH) + try: + user = self._require_admin(connection) + if user is None: + return + data = self._read_json_body() + if data is None: + return + try: + expected_decision_id = ( + int(data["expected_decision_id"]) + if data.get("expected_decision_id") is not None + else None + ) + except (TypeError, ValueError): + self._send_json(400, {"status": "error", "message": "expected_decision_id 无效。"}) + return + try: + payload = manual_records.decide( + connection, + record_id, + str(data.get("action") or ""), + reason=str(data.get("reason") or ""), + expected_decision_id=expected_decision_id, + request_key=str(data.get("request_key") or "") or None, + actor=user, + subject_code=data.get("subject_code"), + target_ledger_event_id=data.get("target_ledger_event_id"), + effective_at=data.get("effective_at"), + ) + except manual_records.ManualConflictError as exc: + self._send_json(409, {"status": "error", "message": str(exc)}) + return + except manual_records.ManualInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", "decision": payload}) + finally: + connection.close() + + # ------------------------------------------------------------------ + # B-44 intercompany positions (company, own-company scope only) + # ------------------------------------------------------------------ + + def _company_intercompany_scope(self, connection): + """Reject company users with 403; return the session company id.""" + user = self._require_user(connection) + if user is None: + return None, None + if user["role"] != "company": + self._send_json(403, {"status": "error", "message": "该操作仅限公司用户。"}) + return None, None + return user, user["company_id"] + + def _handle_company_intercompany_balances(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + try: + from_, cutoff = self._parse_window(query) + currency = (query.get("currency") or [None])[0] + limit = self._parse_limit(query) + cursor = (query.get("cursor") or [None])[0] + except Exception: + return + try: + payload = positions.company_balances( + connection, from_=from_, cutoff=cutoff, currency=currency, + company_id=company_id, limit=limit, cursor=cursor, + ) + counterparties = self._company_counterparty_summary( + connection, company_id, from_, cutoff, currency + ) + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload, "counterparties": counterparties}) + finally: + connection.close() + + def _company_counterparty_summary( + self, connection, company_id: int, from_: str, cutoff: str, currency: str | None + ) -> list[dict[str, object]]: + events = positions.load_events( + connection, from_=from_, cutoff=cutoff, currency=currency, + company_id=company_id, + ) + # Buckets are keyed by ``(counterparty, currency)``: amounts never mix + # across currencies, so one counterparty with CNY and USD yields two rows. + buckets: dict[tuple[int, str], dict[str, object]] = {} + for event in events: + counterparty = ( + event["payee_company_id"] + if int(event["payer_company_id"]) == int(company_id) + else event["payer_company_id"] + ) + name = ( + event["payee_company_name"] + if int(event["payer_company_id"]) == int(company_id) + else event["payer_company_name"] + ) + key = (int(counterparty), event["currency"]) + bucket = buckets.setdefault( + key, + { + "counterparty_company_id": counterparty, + "counterparty_company_name": name, + "currency": event["currency"], + "signed": 0, + "event_count": 0, + }, + ) + bucket["signed"] += positions.signed_amount(event, company_id) + bucket["event_count"] += 1 + + # Counterparties with pending-subject exposure also show up so the + # company portal never hides unconfirmed balances. + pending = connection.execute( + """ + SELECT p.payer_company_id, p.payee_company_id, p.amount, p.currency, + cpayer.name AS payer_company_name, cpayee.name AS payee_company_name + FROM current_ledger_event_revisions cur + JOIN ledger_event_revisions p ON p.id = cur.revision_id + JOIN companies cpayer ON cpayer.id = p.payer_company_id + JOIN companies cpayee ON cpayee.id = p.payee_company_id + WHERE p.state = 'pending_subject' + AND p.effective_at >= ? AND p.effective_at <= ? + AND (p.payer_company_id = ? OR p.payee_company_id = ?) + """, + (from_, cutoff + "T23:59:59", company_id, company_id), + ).fetchall() + for row in pending: + if currency and row["currency"] != currency: + continue + counterparty = ( + row["payee_company_id"] + if int(row["payer_company_id"]) == int(company_id) + else row["payer_company_id"] + ) + name = ( + row["payee_company_name"] + if int(row["payer_company_id"]) == int(company_id) + else row["payer_company_name"] + ) + buckets.setdefault( + (int(counterparty), row["currency"]), + { + "counterparty_company_id": counterparty, + "counterparty_company_name": name, + "currency": row["currency"], + "signed": 0, + "event_count": 0, + }, + ) + items = [] + for (counterparty, cur), bucket in sorted(buckets.items()): + signed = bucket["signed"] + unresolved = positions.unresolved_for_company( + connection, company_id, cutoff, currency=cur, + counterparty_filter=(company_id, counterparty), + ) + direction = "receivable" if signed > 0 else ("payable" if signed < 0 else None) + items.append( + { + "counterparty_company_id": counterparty, + "counterparty_company_name": bucket["counterparty_company_name"], + "currency": cur, + "result": { + "kind": "period_net_change", + "signed_amount": str(signed), + "direction": direction, + "label": "期间净变动", + }, + "unresolved": unresolved, + "event_count": bucket["event_count"], + } + ) + return items + + def _handle_company_intercompany_pair( + self, counterparty_id: int, query: dict[str, list[str]] + ) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + if not self._companies_exist(connection, (counterparty_id,)): + self._send_json(404, {"status": "error", "message": "公司不存在。"}) + return + try: + from_, cutoff = self._parse_window(query) + currency = (query.get("currency") or [None])[0] + except Exception: + return + try: + payload = positions.pair_detail( + connection, company_id, counterparty_id, + from_=from_, cutoff=cutoff, currency=currency, + ) + except positions.PositionError as exc: + self._send_json(422, {"status": "error", "message": str(exc)}) + return + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_company_intercompany_events(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + try: + from_, cutoff = self._parse_window(query) + currency = (query.get("currency") or [None])[0] + subject = (query.get("subject") or [None])[0] + state = (query.get("state") or [None])[0] + posting_kind = (query.get("posting_kind") or [None])[0] + source_kind = (query.get("source_kind") or [None])[0] + limit = self._parse_limit(query) + cursor = (query.get("cursor") or [None])[0] + except Exception: + return + try: + payload = positions.list_events( + connection, from_=from_, cutoff=cutoff, currency=currency, + subject=subject, state=state, posting_kind=posting_kind, + source_kind=source_kind, viewer_company_id=company_id, + limit=limit, cursor=cursor, + ) + except positions.PositionInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_company_intercompany_event_detail(self, event_id: int) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + payload = positions.event_detail(connection, event_id) + if payload is None: + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + event = payload["event"] + if company_id not in (event["payer_company_id"], event["payee_company_id"]): + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + event.update( + positions.event_payload( + connection, + connection.execute( + positions._DETAIL_SELECT + " WHERE p.ledger_event_id = ?", + (event_id,), + ).fetchone(), + viewer_company_id=company_id, + ) + ) + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_company_intercompany_evidence(self, event_id: int) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + detail = positions.event_detail(connection, event_id) + if detail is None: + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + event = detail["event"] + if company_id not in (event["payer_company_id"], event["payee_company_id"]): + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + payload = positions.event_evidence( + connection, event_id, viewer_company_id=company_id + ) + if payload is None: + self._send_json(404, {"status": "error", "message": "事件不存在。"}) + return + self._send_json(200, {"status": "ok", **payload}) + finally: + connection.close() + + def _handle_company_companies(self) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + rows = connection.execute( + "SELECT id, name FROM companies ORDER BY id" + ).fetchall() + self._send_json( + 200, + {"status": "ok", "companies": [dict(row) for row in rows]}, + ) + finally: + connection.close() + + def _handle_company_manual_records(self, query: dict[str, list[str]]) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + state = (query.get("state") or [None])[0] + limit = self._parse_limit(query, 100) + try: + rows = manual_records.list_records( + connection, company_id=company_id, state=state, limit=limit + ) + except manual_records.ManualInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + # A company only ever sees its own submissions, never declarations + # where it is merely the counterparty. + own = [ + manual_records._row_payload(connection, row) + for row in rows + if row["company_id"] == company_id + ] + self._send_json(200, {"status": "ok", "records": own}) + finally: + connection.close() + + def _handle_company_manual_records_submit(self) -> None: + connection = connect(DB_PATH) + try: + user, company_id = self._company_intercompany_scope(connection) + if company_id is None: + return + data = self._read_json_body() + if data is None: + return + try: + payload = manual_records.submit( + connection, + company_id=company_id, + counterparty_company_id=int(data["counterparty_company_id"]), + occurred_at=str(data.get("occurred_at") or ""), + direction=str(data.get("direction") or ""), + amount=str(data.get("amount") or ""), + currency=str(data.get("currency") or ""), + funding_source=str(data.get("funding_source") or ""), + requested_subject=str(data.get("requested_subject") or ""), + request_key=str(data.get("request_key") or ""), + actor=user, + bank_account_id=data.get("bank_account_id"), + personal_transit_mapping_id=data.get("personal_transit_mapping_id"), + related_source_row_id=data.get("related_source_row_id"), + summary=data.get("summary"), + reason=data.get("reason"), + evidence=data.get("evidence"), + ) + except manual_records.ManualConflictError as exc: + self._send_json(409, {"status": "error", "message": str(exc)}) + return + except manual_records.ManualInputError as exc: + self._send_json(400, {"status": "error", "message": str(exc)}) + return + auth.audit( + connection, "manual_submit", actor=user, + target=f"manual_record:{payload['id']}", + detail=f"amount:{payload['amount']};currency:{payload['currency']}", + ip=self._client_ip, + ) + self._send_json(200, {"status": "ok", "record": payload}) + finally: + connection.close() + + # ------------------------------------------------------------------ # Request/response plumbing # ------------------------------------------------------------------ diff --git a/src/bank_importer/db.py b/src/bank_importer/db.py index b83a618..02775bc 100644 --- a/src/bank_importer/db.py +++ b/src/bank_importer/db.py @@ -569,6 +569,235 @@ MIGRATIONS: tuple[Migration, ...] = ( ALTER TABLE import_batches DROP COLUMN upload_bank_account_id; """, ), + Migration( + version=6, + name="0006_intercompany_ledger_events", + # B-44 canonical intercompany ledger layer. Manual records and the + # ledger event revision chain are append-only facts; current pointers + # (current revision per ledger event / manual decision, source claims) + # are rebuildable projections. Bank events enter only through + # ``eligible_intercompany_events``; nothing here rewrites bank rows. + up=""" + CREATE TABLE manual_records ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + company_id INTEGER NOT NULL REFERENCES companies (id), + counterparty_company_id INTEGER NOT NULL REFERENCES companies (id), + occurred_at TEXT NOT NULL, + direction TEXT NOT NULL CHECK (direction IN ('outgoing', 'incoming')), + amount TEXT NOT NULL, + amount_scale INTEGER NOT NULL, + currency TEXT NOT NULL, + funding_source TEXT NOT NULL CHECK (funding_source IN ( + 'approved_bank_account', 'personal_transit', 'other' + )), + bank_account_id INTEGER REFERENCES bank_accounts (id), + personal_transit_mapping_id INTEGER REFERENCES personal_transit_mappings (id), + related_source_row_id INTEGER REFERENCES source_rows (id), + requested_subject TEXT NOT NULL CHECK (requested_subject IN ( + 'receivable', 'payable', 'other_receivable', 'other_payable' + )), + summary TEXT, + reason TEXT, + evidence_json TEXT, + request_key TEXT NOT NULL, + supersedes_record_id INTEGER REFERENCES manual_records (id), + submitted_by INTEGER REFERENCES users (id), + created_at TEXT NOT NULL, + UNIQUE (company_id, request_key), + CHECK (counterparty_company_id != company_id) + ); + + CREATE TABLE manual_record_decisions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + record_id INTEGER NOT NULL REFERENCES manual_records (id), + revision INTEGER NOT NULL, + state TEXT NOT NULL CHECK (state IN ( + 'pending', 'approved', 'returned', 'exception', 'reversed' + )), + action TEXT NOT NULL, + reason TEXT, + actor_user_id INTEGER REFERENCES users (id), + actor_username TEXT, + idempotency_key TEXT, + supersedes_decision_id INTEGER REFERENCES manual_record_decisions (id), + created_at TEXT NOT NULL, + UNIQUE (record_id, revision) + ); + + CREATE TABLE current_manual_record_decisions ( + record_id INTEGER PRIMARY KEY REFERENCES manual_records (id), + decision_id INTEGER NOT NULL UNIQUE REFERENCES manual_record_decisions (id) + ); + + CREATE TABLE ledger_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + lifecycle TEXT NOT NULL DEFAULT 'active' + CHECK (lifecycle IN ('active', 'superseded')), + created_at TEXT NOT NULL + ); + + CREATE TABLE ledger_event_revisions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ledger_event_id INTEGER NOT NULL REFERENCES ledger_events (id), + revision INTEGER NOT NULL, + state TEXT NOT NULL CHECK (state IN ('pending_subject', 'confirmed')), + effective_at TEXT NOT NULL, + amount TEXT NOT NULL, + amount_scale INTEGER NOT NULL, + currency TEXT NOT NULL, + payer_company_id INTEGER NOT NULL REFERENCES companies (id), + payee_company_id INTEGER NOT NULL REFERENCES companies (id), + perspective_company_id INTEGER REFERENCES companies (id), + subject_code TEXT CHECK (subject_code IN ( + 'receivable', 'payable', 'other_receivable', 'other_payable' + )), + source_kind TEXT NOT NULL CHECK (source_kind IN ('bank', 'manual', 'adjustment')), + source_revision_token TEXT, + posting_kind TEXT NOT NULL CHECK (posting_kind IN ( + 'normal', 'reversal', 'adjustment' + )), + reverses_ledger_event_id INTEGER REFERENCES ledger_events (id), + adjusts_ledger_event_id INTEGER REFERENCES ledger_events (id), + rule_version TEXT, + evidence_json TEXT, + idempotency_key TEXT, + actor_user_id INTEGER REFERENCES users (id), + actor_username TEXT, + reason TEXT, + supersedes_revision_id INTEGER REFERENCES ledger_event_revisions (id), + created_at TEXT NOT NULL, + UNIQUE (ledger_event_id, revision), + CHECK (payer_company_id != payee_company_id), + CHECK (state = 'confirmed' OR subject_code IS NULL), + CHECK (state = 'confirmed' OR perspective_company_id IS NULL), + CHECK ( + state != 'confirmed' + OR (perspective_company_id IS NOT NULL AND subject_code IS NOT NULL) + ) + ); + + CREATE TABLE current_ledger_event_revisions ( + ledger_event_id INTEGER PRIMARY KEY REFERENCES ledger_events (id), + revision_id INTEGER NOT NULL UNIQUE REFERENCES ledger_event_revisions (id) + ); + + CREATE TABLE ledger_event_bank_sources ( + bank_event_id INTEGER PRIMARY KEY REFERENCES canonical_transfer_events (id), + ledger_event_id INTEGER NOT NULL REFERENCES ledger_events (id), + UNIQUE (ledger_event_id, bank_event_id) + ); + + CREATE TABLE ledger_event_manual_sources ( + manual_record_id INTEGER PRIMARY KEY REFERENCES manual_records (id), + ledger_event_id INTEGER NOT NULL REFERENCES ledger_events (id), + UNIQUE (ledger_event_id, manual_record_id) + ); + + CREATE TABLE ledger_subject_suggestions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ledger_event_id INTEGER NOT NULL REFERENCES ledger_events (id), + source_revision_id INTEGER NOT NULL REFERENCES ledger_event_revisions (id), + suggested_perspective_company_id INTEGER NOT NULL REFERENCES companies (id), + suggested_subject_code TEXT NOT NULL CHECK (suggested_subject_code IN ( + 'receivable', 'payable', 'other_receivable', 'other_payable' + )), + rule_version TEXT, + evidence_json TEXT, + created_at TEXT NOT NULL + ); + + CREATE INDEX idx_ledger_revisions_event ON ledger_event_revisions (ledger_event_id, revision); + CREATE INDEX idx_ledger_revisions_effective ON ledger_event_revisions (effective_at); + CREATE INDEX idx_ledger_revisions_pair_currency + ON ledger_event_revisions (payer_company_id, payee_company_id, currency); + CREATE INDEX idx_ledger_revisions_state + ON ledger_event_revisions (state, effective_at); + CREATE INDEX idx_manual_records_company ON manual_records (company_id, occurred_at); + CREATE INDEX idx_manual_decisions_record ON manual_record_decisions (record_id, revision); + CREATE INDEX idx_manual_decisions_state ON manual_record_decisions (state); + CREATE INDEX idx_subject_suggestions_event ON ledger_subject_suggestions (ledger_event_id); + + CREATE VIEW eligible_position_events AS + SELECT le.id AS ledger_event_id, + cur.revision_id AS ledger_revision_id, + r.effective_at AS effective_at, + r.amount AS amount, r.amount_scale AS amount_scale, + r.currency AS currency, + r.payer_company_id AS payer_company_id, + r.payee_company_id AS payee_company_id, + r.perspective_company_id AS perspective_company_id, + r.subject_code AS subject_code, + r.source_kind AS source_kind, r.posting_kind AS posting_kind, + r.reverses_ledger_event_id AS reverses_ledger_event_id, + r.adjusts_ledger_event_id AS adjusts_ledger_event_id, + COALESCE( + (SELECT bs.bank_event_id FROM ledger_event_bank_sources bs + WHERE bs.ledger_event_id = le.id LIMIT 1), + (SELECT ms.manual_record_id FROM ledger_event_manual_sources ms + WHERE ms.ledger_event_id = le.id LIMIT 1) + ) AS source_id, + ((SELECT COUNT(*) FROM ledger_event_bank_sources bs + WHERE bs.ledger_event_id = le.id) + + (SELECT COUNT(*) FROM ledger_event_manual_sources ms + WHERE ms.ledger_event_id = le.id)) AS evidence_count + FROM ledger_events le + JOIN current_ledger_event_revisions cur ON cur.ledger_event_id = le.id + JOIN ledger_event_revisions r ON r.id = cur.revision_id + WHERE le.lifecycle = 'active' AND r.state = 'confirmed'; + + CREATE TRIGGER ledger_events_no_delete BEFORE DELETE ON ledger_events + BEGIN SELECT RAISE (ABORT, 'ledger_events rows are immutable'); END; + CREATE TRIGGER ledger_events_no_update BEFORE UPDATE ON ledger_events + BEGIN + SELECT RAISE (ABORT, 'ledger_events only allow lifecycle changes') + WHERE OLD.lifecycle = NEW.lifecycle + OR OLD.id IS NOT NEW.id + OR OLD.created_at IS NOT NEW.created_at; + END; + + CREATE TRIGGER ledger_event_revisions_no_update BEFORE UPDATE ON ledger_event_revisions + BEGIN SELECT RAISE (ABORT, 'ledger_event_revisions rows are immutable'); END; + CREATE TRIGGER ledger_event_revisions_no_delete BEFORE DELETE ON ledger_event_revisions + BEGIN SELECT RAISE (ABORT, 'ledger_event_revisions rows are immutable'); END; + + CREATE TRIGGER ledger_subject_suggestions_no_update BEFORE UPDATE ON ledger_subject_suggestions + BEGIN SELECT RAISE (ABORT, 'ledger_subject_suggestions rows are immutable'); END; + CREATE TRIGGER ledger_subject_suggestions_no_delete BEFORE DELETE ON ledger_subject_suggestions + BEGIN SELECT RAISE (ABORT, 'ledger_subject_suggestions rows are immutable'); END; + + CREATE TRIGGER manual_records_no_update BEFORE UPDATE ON manual_records + BEGIN SELECT RAISE (ABORT, 'manual_records rows are immutable'); END; + CREATE TRIGGER manual_records_no_delete BEFORE DELETE ON manual_records + BEGIN SELECT RAISE (ABORT, 'manual_records rows are immutable'); END; + + CREATE TRIGGER manual_record_decisions_no_update BEFORE UPDATE ON manual_record_decisions + BEGIN SELECT RAISE (ABORT, 'manual_record_decisions rows are immutable'); END; + CREATE TRIGGER manual_record_decisions_no_delete BEFORE DELETE ON manual_record_decisions + BEGIN SELECT RAISE (ABORT, 'manual_record_decisions rows are immutable'); END; + """, + down=""" + DROP VIEW IF EXISTS eligible_position_events; + DROP TRIGGER IF EXISTS manual_record_decisions_no_delete; + DROP TRIGGER IF EXISTS manual_record_decisions_no_update; + DROP TRIGGER IF EXISTS manual_records_no_delete; + DROP TRIGGER IF EXISTS manual_records_no_update; + DROP TRIGGER IF EXISTS ledger_subject_suggestions_no_delete; + DROP TRIGGER IF EXISTS ledger_subject_suggestions_no_update; + DROP TRIGGER IF EXISTS ledger_event_revisions_no_delete; + DROP TRIGGER IF EXISTS ledger_event_revisions_no_update; + DROP TRIGGER IF EXISTS ledger_events_no_update; + DROP TRIGGER IF EXISTS ledger_events_no_delete; + DROP TABLE IF EXISTS ledger_subject_suggestions; + DROP TABLE IF EXISTS ledger_event_manual_sources; + DROP TABLE IF EXISTS ledger_event_bank_sources; + DROP TABLE IF EXISTS current_ledger_event_revisions; + DROP TABLE IF EXISTS ledger_event_revisions; + DROP TABLE IF EXISTS ledger_events; + DROP TABLE IF EXISTS current_manual_record_decisions; + DROP TABLE IF EXISTS manual_record_decisions; + DROP TABLE IF EXISTS manual_records; + """, + ), ) diff --git a/src/bank_importer/importing.py b/src/bank_importer/importing.py index 6c5fc0b..291fe9d 100644 --- a/src/bank_importer/importing.py +++ b/src/bank_importer/importing.py @@ -25,6 +25,7 @@ import sqlite3 import tempfile from . import auth +from . import ledger_events from . import matching from .db import utc_now from .models import SheetResult, StatementBatch @@ -648,6 +649,7 @@ def review_sheets( [item["id"] for item in confirmed_rows], actor=actor, ) + ledger_events.reconcile_bank_events(connection, actor=actor) if began: connection.commit() except Exception: diff --git a/src/bank_importer/ledger_events.py b/src/bank_importer/ledger_events.py new file mode 100644 index 0000000..694facc --- /dev/null +++ b/src/bank_importer/ledger_events.py @@ -0,0 +1,684 @@ +"""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 diff --git a/src/bank_importer/manual_records.py b/src/bank_importer/manual_records.py new file mode 100644 index 0000000..bb81059 --- /dev/null +++ b/src/bank_importer/manual_records.py @@ -0,0 +1,818 @@ +"""Manual evidence records, administrator approval and audit-safe reversal. + +Manual records are immutable submitted facts. Only an approved record becomes +a canonical ledger event (``approve_new``) or joins one (``approve_link``); +returned/exception/pending records never affect a balance and never leak to +the counterparty. Approved facts change only through a ``reverse`` decision +that creates an opposite new event (or detaches a linked claim) — the original +is never edited. Idempotency keys and UNIQUE claims prevent double counting. +""" + +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from decimal import Decimal, InvalidOperation +import json +import sqlite3 + +from .db import utc_now +from .ledger_events import ( + LedgerConflictError, + LedgerInputError, + create_event, + create_reversal, + current_revision, + manual_source_claim, +) +from .subjects import SUBJECTS + +MANUAL_STATES = ("pending", "approved", "returned", "exception", "reversed") +FUNDING_SOURCES = ("approved_bank_account", "personal_transit", "other") +DATE_KEYS = ("occurred_at",) + + +class ManualConflictError(ValueError): + """A claim/idempotency/revision conflict (mapped to HTTP 409).""" + + +class ManualInputError(ValueError): + """Invalid manual record input (mapped to HTTP 400/422).""" + + +def _parse_amount(amount: object) -> Decimal: + try: + value = Decimal(str(amount)) + except InvalidOperation: + raise ManualInputError("金额不是有效的十进制数。") from None + if not value.is_finite() or value <= 0: + raise ManualInputError("金额必须大于零。") + return value + + +def _validate_date(value: object, field: str) -> str: + text = str(value or "").strip() + if len(text) < 10: + raise ManualInputError(f"{field}必须是 YYYY-MM-DD 或完整时间。") + try: + datetime.fromisoformat(text[:10]) + except ValueError: + raise ManualInputError(f"{field}必须是 YYYY-MM-DD 或完整时间。") from None + return text + + +def _business_today() -> str: + """Shanghai business date (the default reversal effective date).""" + return datetime.now(timezone(timedelta(hours=8))).date().isoformat() + + +# --------------------------------------------------------------------------- +# Submit +# --------------------------------------------------------------------------- + + +def submit( + connection: sqlite3.Connection, + *, + company_id: int, + counterparty_company_id: int, + occurred_at: str, + direction: str, + amount: str, + currency: str, + funding_source: str, + requested_subject: str, + request_key: str, + actor: sqlite3.Row, + bank_account_id: object = None, + personal_transit_mapping_id: object = None, + related_source_row_id: object = None, + summary: object = None, + reason: object = None, + evidence: object = None, + supersedes_record_id: object = None, +) -> dict[str, object]: + """Submit one manual record for review. Idempotent on ``(company_id, request_key)``.""" + request_key = str(request_key or "").strip() + if not request_key: + raise ManualInputError("必须提供提交幂等键 request_key。") + if int(company_id) == int(counterparty_company_id): + raise ManualInputError("对方公司不能与本公司相同。") + if direction not in ("outgoing", "incoming"): + raise ManualInputError("方向必须是 outgoing 或 incoming。") + if funding_source not in FUNDING_SOURCES: + raise ManualInputError(f"资金来源必须是:{'、'.join(FUNDING_SOURCES)}。") + if requested_subject not in SUBJECTS: + raise ManualInputError("科目必须是应收/应付/其他应收/其他应付之一。") + _validate_date(occurred_at, "业务日期") + currency = str(currency or "").strip() + if not currency: + raise ManualInputError("币种不能为空。") + amount = str(_parse_amount(amount)) + + for label, raw in ( + ("company_id", company_id), ("counterparty_company_id", counterparty_company_id), + ): + row = connection.execute("SELECT id FROM companies WHERE id = ?", (int(raw),)).fetchone() + if row is None: + raise ManualInputError(f"{label} 指向的公司不存在。") + + bank_account_id = _resolve_account_ref( + connection, bank_account_id, company_id, "银行账户" + ) + mapping_id = _resolve_account_ref( + connection, personal_transit_mapping_id, company_id, "个人过账映射" + ) + related_row = None + if related_source_row_id not in (None, ""): + related_row = connection.execute( + "SELECT r.id, b.company_id FROM source_rows r " + "JOIN sheet_batches s ON s.id = r.sheet_batch_id " + "JOIN import_batches b ON b.id = s.import_batch_id " + "WHERE r.id = ?", + (int(related_source_row_id),), + ).fetchone() + if related_row is None: + raise ManualInputError("关联银行源行不存在。") + + if funding_source == "approved_bank_account" and bank_account_id is None: + raise ManualInputError("资金来源为已批准账户时必须指定银行账户。") + if funding_source == "personal_transit" and mapping_id is None: + raise ManualInputError("资金来源为个人过账时必须指定个人过账映射。") + + supersedes_id = None + if supersedes_record_id not in (None, ""): + parent = connection.execute( + "SELECT id, company_id FROM manual_records WHERE id = ?", + (int(supersedes_record_id),), + ).fetchone() + if parent is None or parent["company_id"] != int(company_id): + raise ManualInputError("supersedes_record_id 无效。") + supersedes_id = int(supersedes_record_id) + + now = utc_now() + began = False + if not connection.in_transaction: + connection.execute("BEGIN IMMEDIATE") + began = True + try: + # Re-check inside the write transaction: concurrent identical submits + # serialize here, so a replay is found before any INSERT. + existing = connection.execute( + "SELECT id FROM manual_records WHERE company_id = ? AND request_key = ?", + (int(company_id), request_key), + ).fetchone() + if existing is not None: + if began: + connection.commit() + return _record_payload(connection, existing["id"], idempotent_replay=True) + try: + cursor = connection.execute( + """ + INSERT INTO manual_records ( + company_id, counterparty_company_id, occurred_at, direction, + amount, amount_scale, currency, funding_source, bank_account_id, + personal_transit_mapping_id, related_source_row_id, + requested_subject, summary, reason, evidence_json, request_key, + supersedes_record_id, submitted_by, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + int(company_id), int(counterparty_company_id), occurred_at, direction, + amount, _scale_of(amount), currency, funding_source, bank_account_id, + mapping_id, related_row["id"] if related_row is not None else None, + requested_subject, str(summary or "") or None, + str(reason or "") or None, + json.dumps(evidence, ensure_ascii=False) if evidence else None, + request_key, supersedes_id, + actor["id"], now, + ), + ) + except sqlite3.IntegrityError: + # A concurrent identical submit won the race and committed first; + # surface the existing record idempotently instead of a UNIQUE 500. + if began: + connection.rollback() + existing = connection.execute( + "SELECT id FROM manual_records WHERE company_id = ? AND request_key = ?", + (int(company_id), request_key), + ).fetchone() + if existing is not None: + return _record_payload(connection, existing["id"], idempotent_replay=True) + raise + record_id = int(cursor.lastrowid) + _append_decision( + connection, record_id, state="pending", action="submit", + reason=str(reason or "") or None, actor=actor, + ) + except Exception: + if began: + connection.rollback() + raise + else: + if began: + connection.commit() + return _record_payload(connection, record_id) + + +def _resolve_account_ref(connection, raw, company_id: int, label: str) -> int | None: + if raw in (None, ""): + return None + row = connection.execute( + "SELECT id, company_id FROM bank_accounts WHERE id = ?", (int(raw),) + ).fetchone() + if row is None: + raise ManualInputError(f"{label}不存在。") + if row["company_id"] != int(company_id): + raise ManualInputError(f"{label}必须属于提交公司。") + return int(raw) + + +def _scale_of(amount: str) -> int: + exponent = Decimal(amount).as_tuple().exponent + return max(0, -int(exponent)) + + +# --------------------------------------------------------------------------- +# Decisions +# --------------------------------------------------------------------------- + + +def _append_decision( + connection: sqlite3.Connection, + record_id: int, + *, + state: str, + action: str, + reason: str | None, + actor: sqlite3.Row, + idempotency_key: str | None = None, + supersedes_decision_id: int | None = None, +) -> int: + row = connection.execute( + "SELECT COALESCE(MAX(revision), 0) AS m FROM manual_record_decisions WHERE record_id = ?", + (record_id,), + ).fetchone() + revision = int(row["m"]) + 1 + now = utc_now() + cursor = connection.execute( + """ + INSERT INTO manual_record_decisions ( + record_id, revision, state, action, reason, actor_user_id, + actor_username, idempotency_key, supersedes_decision_id, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + record_id, revision, state, action, reason, + actor["id"], actor["username"], idempotency_key, + supersedes_decision_id, now, + ), + ) + decision_id = int(cursor.lastrowid) + connection.execute( + """ + INSERT OR REPLACE INTO current_manual_record_decisions (record_id, decision_id) + VALUES (?, ?) + """, + (record_id, decision_id), + ) + return decision_id + + +def _current_decision(connection: sqlite3.Connection, record_id: int) -> sqlite3.Row | None: + return connection.execute( + """ + SELECT d.* FROM current_manual_record_decisions c + JOIN manual_record_decisions d ON d.id = c.decision_id + WHERE c.record_id = ? + """, + (record_id,), + ).fetchone() + + +def decide( + connection: sqlite3.Connection, + record_id: int, + action: str, + *, + reason: str, + expected_decision_id: int | None, + request_key: str | None, + actor: sqlite3.Row, + subject_code: object = None, + target_ledger_event_id: object = None, + effective_at: object = None, +) -> dict[str, object]: + """Apply an administrator decision to a manual record. + + ``approve_new`` creates a confirmed ledger event; ``approve_link`` joins an + existing ledger event without adding a second economic impact; ``return`` + and ``exception`` never produce a balance; ``reverse`` creates an opposite + reversal event (or detaches a linked claim) with an independent business + effective date — explicit ``effective_at`` or the approval business day. + Replays return the earlier outcome via ``idempotency_key``. + """ + reason = (reason or "").strip() + if not reason: + raise ManualInputError("必须填写审核原因。") + if action not in ("approve_new", "approve_link", "return", "exception", "reverse"): + raise ManualInputError("未知的审核决定类型。") + + began = False + if not connection.in_transaction: + connection.execute("BEGIN IMMEDIATE") + began = True + try: + record = connection.execute( + "SELECT * FROM manual_records WHERE id = ?", (record_id,) + ).fetchone() + if record is None: + raise ManualConflictError("手工记录不存在。") + if request_key: + existing = connection.execute( + "SELECT * FROM manual_record_decisions WHERE record_id = ? AND idempotency_key = ?", + (record_id, request_key), + ).fetchone() + if existing is not None: + if began: + connection.commit() + return _decision_payload(connection, record_id, existing["id"]) + + current = _current_decision(connection, record_id) + if current is None: + raise ManualConflictError("该记录没有当前状态。") + if expected_decision_id is not None and int(expected_decision_id) != current["id"]: + raise ManualConflictError("记录已发生变更,请刷新后重试。") + + if action == "approve_new": + outcome = _approve_new( + connection, record, current, actor, subject_code, reason, request_key + ) + elif action == "approve_link": + outcome = _approve_link( + connection, record, current, actor, target_ledger_event_id, reason, + request_key, + ) + elif action == "return": + if current["state"] != "pending": + raise ManualConflictError("只有待复核的记录可以退回。") + decision_id = _append_decision( + connection, record_id, state="returned", action=action, + reason=reason, actor=actor, idempotency_key=request_key, + supersedes_decision_id=current["id"], + ) + outcome = {"decision_id": decision_id, "ledger_event_id": None} + elif action == "exception": + if current["state"] != "pending": + raise ManualConflictError("只有待复核的记录可以转为异常。") + decision_id = _append_decision( + connection, record_id, state="exception", action=action, + reason=reason, actor=actor, idempotency_key=request_key, + supersedes_decision_id=current["id"], + ) + outcome = {"decision_id": decision_id, "ledger_event_id": None} + else: # reverse + if current["state"] != "approved": + raise ManualConflictError("只有已批准记录可以冲销。") + outcome = _reverse( + connection, record, current, actor, request_key, reason, + effective_at=effective_at, + ) + + _store_audit(connection, record, current, action, outcome, reason, actor) + except Exception: + if began: + connection.rollback() + raise + else: + if began: + connection.commit() + return _decision_payload(connection, record_id, outcome["decision_id"]) + + +def _approve_new( + connection: sqlite3.Connection, + record: sqlite3.Row, + current: sqlite3.Row, + actor: sqlite3.Row, + subject_code: object, + reason: str, + idempotency_key: str | None, +) -> dict[str, object]: + subject = str(subject_code or record["requested_subject"] or "") + if subject not in SUBJECTS: + raise ManualInputError("科目必须是应收/应付/其他应收/其他应付之一。") + if record["direction"] == "outgoing": + payer, payee = record["company_id"], record["counterparty_company_id"] + else: + payer, payee = record["counterparty_company_id"], record["company_id"] + event_id, _revision_id = create_event( + connection, + state="confirmed", + effective_at=record["occurred_at"], + amount=record["amount"], + currency=record["currency"], + payer_company_id=payer, + payee_company_id=payee, + perspective_company_id=record["company_id"], + subject_code=subject, + source_kind="manual", + source_revision_token=None, + posting_kind="normal", + rule_version="manual-record-v1", + evidence_json=json.dumps({"manual_record_id": record["id"]}, ensure_ascii=False), + actor=actor, + reason=reason, + ) + decision_id = _append_decision( + connection, record["id"], state="approved", action="approve_new", + reason=reason, actor=actor, + idempotency_key=idempotency_key, + supersedes_decision_id=current["id"], + ) + connection.execute( + """ + INSERT INTO ledger_event_manual_sources (manual_record_id, ledger_event_id) + VALUES (?, ?) + """, + (record["id"], event_id), + ) + return {"decision_id": decision_id, "ledger_event_id": event_id} + + +def _approve_link( + connection: sqlite3.Connection, + record: sqlite3.Row, + current: sqlite3.Row, + actor: sqlite3.Row, + target_ledger_event_id: object, + reason: str, + idempotency_key: str | None, +) -> dict[str, object]: + if target_ledger_event_id in (None, ""): + raise ManualInputError("approve_link 必须指定目标往来事件。") + target = connection.execute( + "SELECT id, lifecycle FROM ledger_events WHERE id = ?", + (int(target_ledger_event_id),), + ).fetchone() + if target is None or target["lifecycle"] != "active": + raise ManualConflictError("目标往来事件不存在。") + if manual_source_claim(connection, record["id"]) is not None: + raise ManualConflictError("该手工记录已关联往来事件。") + decision_id = _append_decision( + connection, record["id"], state="approved", action="approve_link", + reason=reason, actor=actor, + idempotency_key=idempotency_key, + supersedes_decision_id=current["id"], + ) + connection.execute( + """ + INSERT INTO ledger_event_manual_sources (manual_record_id, ledger_event_id) + VALUES (?, ?) + """, + (record["id"], int(target_ledger_event_id)), + ) + return {"decision_id": decision_id, "ledger_event_id": int(target_ledger_event_id)} + + +def _reverse( + connection: sqlite3.Connection, + record: sqlite3.Row, + current: sqlite3.Row, + actor: sqlite3.Row, + request_key: str | None, + reason: str, + effective_at: object = None, +) -> dict[str, object]: + claim = manual_source_claim(connection, record["id"]) + if claim is None: + raise ManualConflictError("该记录尚未关联往来事件,无法冲销。") + event_id = claim["ledger_event_id"] + revision = current_revision(connection, event_id) + if revision is None: + raise ManualConflictError("关联的往来事件没有当前修订。") + + if _is_manual_creation_source(revision, record["id"]): + # This record created the event (approve_new): it added the economic + # impact, so reversing it must always produce an equal-amount reversal + # event — even when other manual evidence was later linked onto the + # same event. The original impact must not survive in balances. + if effective_at is not None and str(effective_at).strip(): + effective_at = _validate_date(effective_at, "冲销生效日") + else: + effective_at = _business_today() + create_reversal( + connection, + event_id, + source_kind="manual", + source_revision_token=str(record["id"]), + effective_at=effective_at, + reason=reason, + actor=actor, + idempotency_key=request_key, + rule_version="manual-record-v1", + ) + else: + # The manual was linked evidence (approve_link) on an event it never + # created: it added no second impact, so reversing detaches the claim + # and the underlying economic impact stays. + connection.execute( + "DELETE FROM ledger_event_manual_sources WHERE manual_record_id = ?", + (record["id"],), + ) + decision_id = _append_decision( + connection, record["id"], state="reversed", action="reverse", + reason=reason, actor=actor, idempotency_key=request_key, + supersedes_decision_id=current["id"], + ) + return {"decision_id": decision_id, "ledger_event_id": None} + + +def _is_manual_creation_source(revision: sqlite3.Row, manual_record_id: int) -> bool: + """True when ``manual_record_id`` created the event via ``approve_new``. + + The creation source is recorded in the event's immutable revision + ``evidence_json``; a linked evidence record is never the creation source + and carries no second economic impact. + """ + if revision["source_kind"] != "manual": + return False + try: + evidence = json.loads(revision["evidence_json"] or "{}") + except (TypeError, ValueError): + return False + return evidence.get("manual_record_id") == manual_record_id + + +def _store_audit( + connection, record, current, action, outcome, reason, actor +) -> None: + from .auth import audit + + audit( + connection, + f"manual_{action}", + actor=actor, + target=f"manual_record:{record['id']}", + detail=( + f"decision:{outcome['decision_id']};" + f"ledger_event:{outcome.get('ledger_event_id')};reason:{reason}" + ), + ) + + +# --------------------------------------------------------------------------- +# Candidates and queries +# --------------------------------------------------------------------------- + + +def find_candidates(connection: sqlite3.Connection, record_id: int) -> list[dict[str, object]]: + """Deterministic hints shown before approval; never auto-merged.""" + record = connection.execute( + "SELECT * FROM manual_records WHERE id = ?", (record_id,) + ).fetchone() + if record is None: + return [] + wanted_direction = "incoming" if record["direction"] == "outgoing" else "outgoing" + date_prefix = str(record["occurred_at"])[:10] + candidates: list[dict[str, object]] = [] + + bank_rows = connection.execute( + """ + SELECT e.event_id, e.amount, e.currency, e.effective_at, e.pairing, + e.payer_company_id, e.payee_company_id, e.decision_id + FROM eligible_intercompany_events e + WHERE (e.payer_company_id = ? AND e.payee_company_id = ?) + OR (e.payer_company_id = ? AND e.payee_company_id = ?) + ORDER BY e.event_id + """, + ( + record["company_id"], record["counterparty_company_id"], + record["counterparty_company_id"], record["company_id"], + ), + ).fetchall() + for row in bank_rows: + if row["amount"] != record["amount"] or row["currency"] != record["currency"]: + continue + event_direction = ( + "outgoing" if row["payer_company_id"] == record["company_id"] else "incoming" + ) + if event_direction != wanted_direction: + continue + candidates.append( + { + "kind": "bank_event", + "ledger_event_id": _ledger_event_of_bank(connection, row["event_id"]), + "bank_event_id": row["event_id"], + "decision_id": row["decision_id"], + "amount": row["amount"], + "currency": row["currency"], + "effective_at": row["effective_at"], + "pairing": row["pairing"], + "hint": "已存在匹配的银行规范事件,建议关联", + } + ) + + manual_rows = connection.execute( + """ + SELECT m.id, m.company_id, m.counterparty_company_id, m.direction, + m.amount, m.currency, m.occurred_at, d.state + FROM manual_records m + JOIN current_manual_record_decisions c ON c.record_id = m.id + JOIN manual_record_decisions d ON d.id = c.decision_id + WHERE m.id != ? AND m.amount = ? AND m.currency = ? + AND ( + (m.company_id = ? AND m.counterparty_company_id = ?) + OR (m.company_id = ? AND m.counterparty_company_id = ?) + ) + ORDER BY m.id + """, + ( + record["id"], record["amount"], record["currency"], + record["company_id"], record["counterparty_company_id"], + record["counterparty_company_id"], record["company_id"], + ), + ).fetchall() + for row in manual_rows: + if row["direction"] != wanted_direction: + continue + if row["state"] not in ("approved", "pending"): + continue + candidates.append( + { + "kind": "manual_record", + "ledger_event_id": None, + "manual_record_id": row["id"], + "amount": row["amount"], + "currency": row["currency"], + "occurred_at": row["occurred_at"], + "state": row["state"], + "hint": "存在方向相反的同额手工记录,建议核对后关联", + } + ) + return candidates + + +def _ledger_event_of_bank(connection, bank_event_id: int) -> int | None: + claim = connection.execute( + "SELECT ledger_event_id FROM ledger_event_bank_sources WHERE bank_event_id = ?", + (bank_event_id,), + ).fetchone() + return claim["ledger_event_id"] if claim is not None else None + + +def list_records( + connection: sqlite3.Connection, + *, + company_id: int | None = None, + state: str | None = None, + limit: int = 100, +) -> list[sqlite3.Row]: + conditions: list[str] = [] + params: list[object] = [] + if company_id is not None: + conditions.append("(m.company_id = ? OR m.counterparty_company_id = ?)") + params.extend([company_id, company_id]) + if state is not None: + if state not in MANUAL_STATES: + raise ManualInputError("无效的记录状态。") + conditions.append("d.state = ?") + params.append(state) + where = f"WHERE {' AND '.join(conditions)}" if conditions else "" + return connection.execute( + f""" + SELECT m.*, d.id AS decision_id, d.state AS state, d.revision AS decision_revision, + d.action AS action, d.reason AS decision_reason, + d.actor_username AS decision_actor, d.created_at AS decision_at, + c.name AS company_name, cc.name AS counterparty_company_name, + u.username AS submitted_by_username + FROM manual_records m + JOIN current_manual_record_decisions c ON c.record_id = m.id + JOIN manual_record_decisions d ON d.id = c.decision_id + LEFT JOIN companies c ON c.id = m.company_id + LEFT JOIN companies cc ON cc.id = m.counterparty_company_id + LEFT JOIN users u ON u.id = m.submitted_by + {where} + ORDER BY m.id DESC + LIMIT ? + """, + (*params, max(1, int(limit))), + ).fetchall() + + +def rebuild_current_manual_projection(connection: sqlite3.Connection) -> int: + began = False + if not connection.in_transaction: + connection.execute("BEGIN IMMEDIATE") + began = True + try: + connection.execute("DELETE FROM current_manual_record_decisions") + rows = connection.execute( + """ + SELECT m.id AS record_id, + (SELECT d2.id FROM manual_record_decisions d2 + WHERE d2.record_id = m.id + ORDER BY d2.revision DESC LIMIT 1) AS latest_id + FROM manual_records m + """ + ).fetchall() + rebuilt = 0 + for row in rows: + if row["latest_id"] is None: + continue + connection.execute( + """ + INSERT OR REPLACE INTO current_manual_record_decisions (record_id, decision_id) + VALUES (?, ?) + """, + (row["record_id"], row["latest_id"]), + ) + rebuilt += 1 + except Exception: + if began: + connection.rollback() + raise + else: + if began: + connection.commit() + return rebuilt + + +# --------------------------------------------------------------------------- +# Payloads +# --------------------------------------------------------------------------- + + +def _record_payload(connection: sqlite3.Connection, record_id: int, *, idempotent_replay: bool = False) -> dict[str, object]: + rows = list_records(connection, limit=1000) + row = next((item for item in rows if item["id"] == record_id), None) + if row is None: + raise ManualInputError("手工记录不存在。") + payload = _row_payload(connection, row) + if idempotent_replay: + payload["idempotent_replay"] = True + return payload + + +def _row_payload(connection: sqlite3.Connection, row: sqlite3.Row) -> dict[str, object]: + evidence = json.loads(row["evidence_json"] or "{}") if row["evidence_json"] else {} + return { + "id": row["id"], + "company_id": row["company_id"], + "company_name": row["company_name"], + "counterparty_company_id": row["counterparty_company_id"], + "counterparty_company_name": row["counterparty_company_name"], + "occurred_at": row["occurred_at"], + "direction": row["direction"], + "amount": row["amount"], + "currency": row["currency"], + "funding_source": row["funding_source"], + "bank_account_id": row["bank_account_id"], + "personal_transit_mapping_id": row["personal_transit_mapping_id"], + "related_source_row_id": row["related_source_row_id"], + "requested_subject": row["requested_subject"], + "summary": row["summary"], + "reason": row["reason"], + "request_key": row["request_key"], + "supersedes_record_id": row["supersedes_record_id"], + "submitted_by": row["submitted_by"], + "submitted_by_username": row["submitted_by_username"] if "submitted_by_username" in row.keys() else None, + "attachment_name": evidence.get("attachment_name"), + "created_at": row["created_at"], + "state": row["state"], + "decision_id": row["decision_id"], + "decision_revision": row["decision_revision"], + "decision_action": row["action"], + "decision_reason": row["decision_reason"], + "decision_actor": row["decision_actor"], + "decision_at": row["decision_at"], + "candidates": find_candidates(connection, row["id"]), + } + + +def _decision_payload(connection: sqlite3.Connection, record_id: int, decision_id: int) -> dict[str, object]: + row = connection.execute( + """ + SELECT d.* FROM manual_record_decisions d + WHERE d.id = ? + """, + (decision_id,), + ).fetchone() + record = connection.execute( + "SELECT * FROM manual_records WHERE id = ?", (record_id,) + ).fetchone() + claim = manual_source_claim(connection, record_id) + return { + "record_id": record_id, + "decision_id": decision_id, + "revision": row["revision"], + "state": row["state"], + "action": row["action"], + "reason": row["reason"], + "actor_username": row["actor_username"], + "created_at": row["created_at"], + "ledger_event_id": claim["ledger_event_id"] if claim is not None else None, + "requested_subject": record["requested_subject"], + "amount": record["amount"], + "currency": record["currency"], + "occurred_at": record["occurred_at"], + } diff --git a/src/bank_importer/positions.py b/src/bank_importer/positions.py new file mode 100644 index 0000000..9c484b9 --- /dev/null +++ b/src/bank_importer/positions.py @@ -0,0 +1,1314 @@ +"""Intercompany position aggregation: Decimal math, conservation and drill-down. + +Every confirmed ledger event produces a signed claim: the payer books a debit +(claim ``+amount``) and the payee books a credit (claim ``-amount``). For one +company pair and currency the two perspectives must mirror exactly +(``C_A == -C_B``) — the module asserts that invariant after every pair +aggregation and refuses to render a non-conserving number. Subjects are +mirrored from the stored perspective; all money math uses ``Decimal`` on the +stored decimal strings, never SQLite ``SUM`` or floating point. +""" + +from __future__ import annotations + +import base64 +from datetime import datetime, timedelta, timezone +from decimal import Decimal +import json +import sqlite3 + +from .ledger_events import current_revision +from .manual_records import list_records +from .master_data import mask_account_number +from .subjects import MIRROR, SUBJECTS, mirror_subject, subject_label +from . import matching + + +class PositionError(ValueError): + """Calculation failure: conservation violated or bad parameters.""" + + +class PositionInputError(ValueError): + """Invalid query parameters (mapped to HTTP 400).""" + + +SUBJECT_LABEL_MAP = { + "receivable": "应收", "payable": "应付", + "other_receivable": "其他应收", "other_payable": "其他应付", +} + + +def today_shanghai() -> str: + return datetime.now(timezone(timedelta(hours=8))).date().isoformat() + + +def validate_window(from_: object, cutoff: object) -> tuple[str, str]: + pattern = r"^\d{4}-\d{2}-\d{2}$" + import re + + def check(value, label): + text = str(value or "") + if not text or not re.fullmatch(pattern, text): + raise PositionInputError(f"{label}必须是 YYYY-MM-DD 格式。") + try: + datetime.strptime(text, "%Y-%m-%d") + except ValueError: + raise PositionInputError(f"{label}不是有效日期。") from None + return text + + start = check(from_, "from") or "0001-01-01" + end = check(cutoff, "cutoff") + if start > end: + raise PositionInputError("from 不能晚于 cutoff。") + return start, end + + +def encode_cursor(parts: tuple[str, ...]) -> str: + raw = "|".join(str(part) for part in parts) + return base64.urlsafe_b64encode(raw.encode("utf-8")).decode("ascii") + + +def decode_cursor(cursor: str | None, parts: int) -> tuple[str, ...] | None: + if not cursor: + return None + try: + raw = base64.urlsafe_b64decode(cursor.encode("ascii")).decode("utf-8") + except Exception: + raise PositionInputError("分页游标无效。") from None + values = tuple(raw.split("|")) + if len(values) != parts: + raise PositionInputError("分页游标无效。") + return values + + +def _row(r) -> dict[str, object]: + return {key: r[key] for key in r.keys()} + + +# --------------------------------------------------------------------------- +# Event loading +# --------------------------------------------------------------------------- + +_EVENT_SELECT = """ + SELECT p.ledger_event_id, p.ledger_revision_id, p.effective_at, p.amount, + p.amount_scale, p.currency, p.payer_company_id, p.payee_company_id, + p.perspective_company_id, p.subject_code, p.source_kind, + p.posting_kind, p.reverses_ledger_event_id, p.adjusts_ledger_event_id, + p.source_id, p.evidence_count, + cpayer.name AS payer_company_name, cpayee.name AS payee_company_name + FROM eligible_position_events p + JOIN companies cpayer ON cpayer.id = p.payer_company_id + JOIN companies cpayee ON cpayee.id = p.payee_company_id +""" + + +def _event_filters( + *, + from_: str, + cutoff: str, + currency: str | None = None, + company_id: int | None = None, + pair: tuple[int, int] | None = None, + subject: str | None = None, + posting_kind: str | None = None, + source_kind: str | None = None, + viewer_company_id: int | None = None, + state: str | None = None, +) -> tuple[str, list[object]]: + conditions = ["p.effective_at >= ?", "p.effective_at <= ?"] + params: list[object] = [from_, cutoff + "T23:59:59"] + if currency: + conditions.append("p.currency = ?") + params.append(currency) + if company_id is not None or pair is not None: + if pair is not None: + a, b = int(pair[0]), int(pair[1]) + conditions.append( + "(p.payer_company_id IN (?, ?) AND p.payee_company_id IN (?, ?))" + ) + params.extend([a, b, a, b]) + else: + conditions.append("(p.payer_company_id = ? OR p.payee_company_id = ?)") + params.extend([company_id, company_id]) + if subject: + if subject not in SUBJECTS: + raise PositionInputError("科目筛选无效。") + # The stored subject lives on one company's perspective; the viewer on + # the other side sees its mirror. Match both so mirror events are never + # dropped from the filter. + conditions.append("(p.subject_code = ? OR p.subject_code = ?)") + params.extend([subject, MIRROR[subject]]) + if posting_kind: + conditions.append("p.posting_kind = ?") + params.append(posting_kind) + if source_kind: + conditions.append("p.source_kind = ?") + params.append(source_kind) + if viewer_company_id is not None: + conditions.append("(p.payer_company_id = ? OR p.payee_company_id = ?)") + params.extend([viewer_company_id, viewer_company_id]) + if state is not None: + if state not in ("confirmed", "pending_subject"): + raise PositionInputError("无效的事件状态。") + conditions.append("p.state = ?") + params.append(state) + return "WHERE " + " AND ".join(conditions), params + + +def load_events( + connection: sqlite3.Connection, + *, + from_: str, + cutoff: str, + currency: str | None = None, + company_id: int | None = None, + pair: tuple[int, int] | None = None, + subject: str | None = None, + posting_kind: str | None = None, + source_kind: str | None = None, + viewer_company_id: int | None = None, +) -> list[sqlite3.Row]: + where, params = _event_filters( + from_=from_, cutoff=cutoff, currency=currency, company_id=company_id, + pair=pair, subject=subject, posting_kind=posting_kind, + source_kind=source_kind, viewer_company_id=viewer_company_id, + ) + return connection.execute( + _EVENT_SELECT + where + " ORDER BY p.ledger_event_id", params + ).fetchall() + + +def load_all_events( + connection: sqlite3.Connection, + *, + from_: str, + cutoff: str, + currency: str | None = None, + company_id: int | None = None, + pair: tuple[int, int] | None = None, + subject: str | None = None, + state: str | None = None, + posting_kind: str | None = None, + source_kind: str | None = None, + viewer_company_id: int | None = None, +) -> list[sqlite3.Row]: + """Load every current ledger revision (confirmed and pending_subject).""" + where, params = _event_filters( + from_=from_, cutoff=cutoff, currency=currency, company_id=company_id, + pair=pair, subject=subject, posting_kind=posting_kind, + source_kind=source_kind, viewer_company_id=viewer_company_id, + state=state, + ) + return connection.execute( + _DETAIL_SELECT + where + " ORDER BY p.ledger_event_id", params + ).fetchall() + + +def signed_amount(event: sqlite3.Row, company_id: int) -> Decimal: + """+amount when the company is the payer, -amount when it is the payee.""" + value = Decimal(event["amount"]) + if int(event["payer_company_id"]) == int(company_id): + return value + return -value + + +def viewer_direction(event: sqlite3.Row, company_id: int) -> str: + return "outgoing" if int(event["payer_company_id"]) == int(company_id) else "incoming" + + +def viewer_subject(event: sqlite3.Row, company_id: int) -> str | None: + if event["perspective_company_id"] is None: + return None + if int(event["perspective_company_id"]) == int(company_id): + return event["subject_code"] + return mirror_subject(event["subject_code"]) + + +def _subject_side(event: sqlite3.Row, company_id: int) -> tuple[str, str, str]: + """``(subject_code, side)`` for ``company_id``: side is 'debit' or 'credit'.""" + perspective = int(event["perspective_company_id"]) + is_payer = int(event["payer_company_id"]) == int(company_id) + if perspective == int(company_id): + subject = event["subject_code"] + side = "debit" if is_payer else "credit" + else: + subject = mirror_subject(event["subject_code"]) + side = "debit" if is_payer else "credit" + return subject, side, perspective + + +def _balance_payload( + connection: sqlite3.Connection, + events: list[sqlite3.Row], + *, + company_id: int, + from_: str, + cutoff: str, + currency: str, +) -> dict[str, object]: + """One company/currency balance row; zero events still yields the row.""" + debit = credit = signed = Decimal("0") + event_count = 0 + for event in events: + if event["currency"] != currency: + continue + if int(event["payer_company_id"]) == int(company_id): + debit += Decimal(event["amount"]) + signed += Decimal(event["amount"]) + else: + credit += Decimal(event["amount"]) + signed -= Decimal(event["amount"]) + event_count += 1 + unresolved = unresolved_for_company( + connection, company_id, cutoff, currency=currency + ) + direction = "receivable" if signed > 0 else ("payable" if signed < 0 else None) + return { + "company_id": company_id, + "company_name": _company_name(connection, company_id), + "window": {"from": from_, "cutoff": cutoff, "cutoff_inclusive": True}, + "currency": currency, + "opening": {"status": "unavailable", "amount": None}, + "period": {"debit": str(debit), "credit": str(credit)}, + "result": { + "kind": "period_net_change", + "signed_amount": str(signed), + "direction": direction, + "label": "期间净变动", + }, + "unresolved": unresolved, + "trace": { + "event_count": event_count, + "events_url": ( + f"/api/admin/intercompany/events?company_id={company_id}" + f"&from={from_}&cutoff={cutoff}¤cy={currency}" + ), + }, + } + + +def _company_name(connection: sqlite3.Connection, company_id: int) -> str: + row = connection.execute( + "SELECT name FROM companies WHERE id = ?", (company_id,) + ).fetchone() + return row["name"] if row is not None else "" + + +# --------------------------------------------------------------------------- +# Unresolved amounts +# --------------------------------------------------------------------------- + + +def _pending_subject_events( + connection: sqlite3.Connection, cutoff: str +) -> list[sqlite3.Row]: + return connection.execute( + """ + SELECT r.ledger_event_id, r.effective_at, r.amount, r.currency, + r.payer_company_id, r.payee_company_id + FROM current_ledger_event_revisions c + JOIN ledger_event_revisions r ON r.id = c.revision_id + WHERE r.state = 'pending_subject' AND r.effective_at <= ? + ORDER BY r.ledger_event_id + """, + (cutoff + "T23:59:59",), + ).fetchall() + + +def _manual_pending( + connection: sqlite3.Connection, cutoff: str +) -> list[sqlite3.Row]: + return connection.execute( + """ + SELECT m.id, m.company_id, m.counterparty_company_id, m.occurred_at, + m.amount, m.currency + FROM manual_records m + JOIN current_manual_record_decisions c ON c.record_id = m.id + JOIN manual_record_decisions d ON d.id = c.decision_id + WHERE d.state = 'pending' AND m.occurred_at <= ? + ORDER BY m.id + """, + (cutoff + "T23:59:59",), + ).fetchall() + + +def _unmatched_singles( + connection: sqlite3.Connection, company_id: int, cutoff: str +) -> list[sqlite3.Row]: + return matching.unresolved_amounts(connection, company_id, cutoff + "T23:59:59") + + +def unresolved_for_company( + connection: sqlite3.Connection, + company_id: int, + cutoff: str, + currency: str | None = None, + *, + include_own_manual_only: bool = True, + counterparty_filter: tuple[int, int] | None = None, +) -> dict[str, object]: + """Unresolved absolute gross per currency (no plus/minus netting). + + ``by_reason`` buckets: subject_review (bank events waiting for a subject), + unmatched_single (B-43 open observations) and manual_pending (unapproved + manual records; only the submitting company / the admin may see them). + """ + gross: dict[str, Decimal] = {} + counts: dict[str, int] = {} + reasons: dict[str, dict[str, object]] = {} + + def add(reason: str, cur: str, amount: Decimal) -> None: + if currency and cur != currency: + return + if reason not in reasons: + reasons[reason] = {"gross_amount": Decimal("0"), "count": 0} + reasons[reason]["gross_amount"] += amount + reasons[reason]["count"] += 1 + gross[cur] = gross.get(cur, Decimal("0")) + amount + counts[cur] = counts.get(cur, 0) + 1 + + for event in _pending_subject_events(connection, cutoff): + if counterparty_filter is not None: + a, b = counterparty_filter + if not ({event["payer_company_id"], event["payee_company_id"]} == {a, b}): + continue + elif company_id not in (event["payer_company_id"], event["payee_company_id"]): + continue + add("subject_review", event["currency"], Decimal(event["amount"])) + + for row in _unmatched_singles(connection, company_id, cutoff): + if counterparty_filter is not None: + participants = _single_participants(connection, row["decision_id"]) + a, b = counterparty_filter + if participants != {a, b}: + continue + add("unmatched_single", row["currency"], abs(Decimal(row["amount"]))) + + for row in _manual_pending(connection, cutoff): + if counterparty_filter is not None: + a, b = counterparty_filter + if not ({row["company_id"], row["counterparty_company_id"]} == {a, b}): + continue + elif include_own_manual_only and row["company_id"] != int(company_id): + continue + elif not include_own_manual_only and company_id not in ( + row["company_id"], row["counterparty_company_id"], + ): + continue + add("manual_pending", row["currency"], Decimal(row["amount"])) + + items = sorted(gross.items()) + return { + "gross_amount": str(sum((gross[k] for k, _ in items), Decimal("0"))), + "count": sum(counts.values()), + "by_reason": { + reason: { + "gross_amount": str(payload["gross_amount"]), + "count": payload["count"], + } + for reason, payload in sorted(reasons.items()) + }, + } + + +def _single_participants(connection: sqlite3.Connection, decision_id: int) -> set[int]: + rows = connection.execute( + "SELECT company_id FROM transfer_decision_participants WHERE decision_id = ?", + (decision_id,), + ).fetchall() + return {row["company_id"] for row in rows if row["company_id"] is not None} + + +# --------------------------------------------------------------------------- +# Directory +# --------------------------------------------------------------------------- + + +def company_balances( + connection: sqlite3.Connection, + *, + from_: str, + cutoff: str, + currency: str | None = None, + company_id: int | None = None, + limit: int = 50, + cursor: str | None = None, +) -> dict[str, object]: + from_, cutoff = validate_window(from_, cutoff) + events = load_events( + connection, from_=from_, cutoff=cutoff, currency=currency, + company_id=company_id, + ) + all_companies = _event_companies(connection) + items_by_company: dict[int, list[sqlite3.Row]] = {} + for event in events: + for company in (event["payer_company_id"], event["payee_company_id"]): + items_by_company.setdefault(company, []).append(event) + + rows: list[tuple[int, str]] = [] + for company in all_companies: + if company_id is not None and company != company_id: + continue + if company in items_by_company: + rows.append((company, "has_events")) + else: + # Companies with no confirmed events still appear when they have + # unresolved exposure, so the directory never hides risk. + unresolved = unresolved_for_company(connection, company, cutoff, currency) + if unresolved["count"] > 0: + rows.append((company, "unresolved")) + + if cursor is not None: + decoded = decode_cursor(cursor, 2) + cursor_company = int(decoded[0]) + cursor_currency = decoded[1] + else: + cursor_company, cursor_currency = None, None + + flat: list[tuple[int, str]] = [] + pending_subject = _pending_subject_events(connection, cutoff) + for company, _marker in rows: + bucket_events = items_by_company.get(company, []) + pending_currencies = { + event["currency"] + for event in pending_subject + if company in (event["payer_company_id"], event["payee_company_id"]) + } + bucket_currencies = sorted( + {event["currency"] for event in bucket_events} + | pending_currencies + | { + item["currency"] + for item in _manual_pending(connection, cutoff) + if item["company_id"] == company or item["counterparty_company_id"] == company + } + | { + item["currency"] for item in _unmatched_singles(connection, company, cutoff) + } + ) + if currency: + bucket_currencies = [currency] if currency in bucket_currencies else [] + for cur in bucket_currencies: + flat.append((company, cur)) + + filtered: list[tuple[int, str]] = [] + for company, cur in flat: + if cursor_company is not None: + if company < cursor_company: + continue + if company == cursor_company and cur <= cursor_currency: + continue + filtered.append((company, cur)) + + filtered.sort(key=lambda item: (item[0], item[1])) + page = filtered[:limit] + has_more = len(filtered) > limit + next_cursor = None + if has_more: + last = page[-1] + next_cursor = encode_cursor((str(last[0]), last[1])) + + items = [] + for company, cur in page: + items.append( + _balance_payload( + connection, items_by_company.get(company, []), + company_id=company, from_=from_, cutoff=cutoff, currency=cur, + ) + ) + return { + "window": {"from": from_, "cutoff": cutoff, "cutoff_inclusive": True}, + "items": items, + "next_cursor": next_cursor, + "has_more": has_more, + } + + +def _event_companies(connection: sqlite3.Connection) -> list[int]: + """All master companies (plus any with ledger exposure) for the directory.""" + rows = connection.execute( + """ + SELECT id AS company_id FROM companies + UNION + SELECT payer_company_id AS company_id FROM eligible_position_events + UNION + SELECT payee_company_id AS company_id FROM eligible_position_events + UNION + SELECT payer_company_id AS company_id FROM ledger_event_revisions + WHERE state = 'pending_subject' + UNION + SELECT payee_company_id AS company_id FROM ledger_event_revisions + WHERE state = 'pending_subject' + ORDER BY company_id + """ + ).fetchall() + return [row["company_id"] for row in rows] + + +# --------------------------------------------------------------------------- +# Pair detail +# --------------------------------------------------------------------------- + + +def pair_detail( + connection: sqlite3.Connection, + company_a: int, + company_b: int, + *, + from_: str, + cutoff: str, + currency: str | None = None, +) -> dict[str, object]: + from_, cutoff = validate_window(from_, cutoff) + if int(company_a) == int(company_b): + raise PositionInputError("公司对的两个公司不能相同。") + pair = (int(company_a), int(company_b)) + events = load_events( + connection, from_=from_, cutoff=cutoff, currency=currency, pair=pair + ) + pending_currencies = { + event["currency"] + for event in _pending_subject_events(connection, cutoff) + if {event["payer_company_id"], event["payee_company_id"]} == set(pair) + } + manual_currencies = { + row["currency"] + for row in _manual_pending(connection, cutoff) + if {row["company_id"], row["counterparty_company_id"]} == set(pair) + } + currencies = sorted( + {event["currency"] for event in events} | pending_currencies | manual_currencies + ) + if currency: + currencies = [currency] if currency in currencies else [] + + outputs = [] + for cur in currencies: + subset = [event for event in events if event["currency"] == cur] + outputs.append( + _pair_currency( + connection, pair, subset, from_=from_, cutoff=cutoff, currency=cur + ) + ) + if not outputs: + outputs.append( + _pair_currency( + connection, pair, [], from_=from_, cutoff=cutoff, + currency=currency or "CNY", + ) + ) + return { + "window": {"from": from_, "cutoff": cutoff, "cutoff_inclusive": True}, + "companies": { + "a": {"company_id": pair[0], "name": _company_name(connection, pair[0])}, + "b": {"company_id": pair[1], "name": _company_name(connection, pair[1])}, + }, + "items": outputs, + } + + +def _pair_currency( + connection: sqlite3.Connection, + pair: tuple[int, int], + events: list[sqlite3.Row], + *, + from_: str, + cutoff: str, + currency: str, +) -> dict[str, object]: + a, b = pair + debit_a = credit_a = debit_b = credit_b = Decimal("0") + subject_sides: dict[str, dict[str, object]] = {} + for subject in SUBJECTS: + subject_sides[subject] = { + "subject_code": subject, + "label": SUBJECT_LABEL_MAP[subject], + "a_debit": Decimal("0"), "a_credit": Decimal("0"), + "b_debit": Decimal("0"), "b_credit": Decimal("0"), + "count": 0, + } + for event in events: + amount = Decimal(event["amount"]) + if int(event["payer_company_id"]) == a: + debit_a += amount + credit_b += amount + else: + debit_b += amount + credit_a += amount + for company in (a, b): + subject, side, _perspective = _subject_side(event, company) + bucket = subject_sides[subject] + bucket[f"{'a' if company == a else 'b'}_{side}"] += amount + bucket["count"] += 1 + + signed_a = debit_a - credit_a + signed_b = debit_b - credit_b + if signed_a != -signed_b: + raise PositionError( + "公司对余额不守恒:A 与 B 的净结果未镜像,拒绝展示该数据。" + ) + + def result_of(signed: Decimal) -> dict[str, object]: + direction = "receivable" if signed > 0 else ("payable" if signed < 0 else None) + return { + "kind": "period_net_change", + "signed_amount": str(signed), + "direction": direction, + "label": "期间净变动", + } + + exception_subjects = [ + subject for subject, bucket in subject_sides.items() + if _subject_abnormal(bucket["a_debit"], bucket["a_credit"], subject) + or _subject_abnormal(bucket["b_debit"], bucket["b_credit"], subject) + ] + return { + "currency": currency, + "opening": {"status": "unavailable", "amount": None}, + "a": { + "period": {"debit": str(debit_a), "credit": str(credit_a)}, + "result": result_of(signed_a), + }, + "b": { + "period": {"debit": str(debit_b), "credit": str(credit_b)}, + "result": result_of(signed_b), + }, + "conservation": { + "abs_equal": abs(signed_a) == abs(signed_b), + "opposite": signed_a == -signed_b, + }, + "subjects": { + subject: { + "subject_code": bucket["subject_code"], + "label": bucket["label"], + "a_debit": str(bucket["a_debit"]), + "a_credit": str(bucket["a_credit"]), + "b_debit": str(bucket["b_debit"]), + "b_credit": str(bucket["b_credit"]), + "count": bucket["count"], + } + for subject, bucket in subject_sides.items() + if bucket["count"] > 0 + }, + "normal_balance_exception": exception_subjects, + "unresolved": unresolved_for_company( + connection, a, cutoff, currency=currency, + counterparty_filter=pair, + ), + "trace": { + "event_count": len(events), + "events_url": ( + f"/api/admin/intercompany/events?company_a={a}&company_b={b}" + f"&from={from_}&cutoff={cutoff}¤cy={currency}" + ), + }, + } + + +def _subject_abnormal(debit: Decimal, credit: Decimal, subject: str) -> bool: + if subject in ("receivable", "other_receivable"): + return debit < credit + return credit < debit + + +# --------------------------------------------------------------------------- +# Events list and detail +# --------------------------------------------------------------------------- + + +def list_events( + connection: sqlite3.Connection, + *, + from_: str, + cutoff: str, + currency: str | None = None, + company_id: int | None = None, + company_a: int | None = None, + company_b: int | None = None, + subject: str | None = None, + state: str | None = None, + posting_kind: str | None = None, + source_kind: str | None = None, + viewer_company_id: int | None = None, + limit: int = 50, + cursor: str | None = None, +) -> dict[str, object]: + from_, cutoff = validate_window(from_, cutoff) + pair = None + if company_a is not None and company_b is not None: + pair = (int(company_a), int(company_b)) + elif company_a is not None: + company_id = int(company_a) + events = load_all_events( + connection, from_=from_, cutoff=cutoff, currency=currency, + company_id=company_id, pair=pair, subject=subject, state=state, + posting_kind=posting_kind, source_kind=source_kind, + viewer_company_id=viewer_company_id, + ) + + def sort_key(event): + return (event["effective_at"], event["ledger_event_id"]) + + events.sort(key=sort_key, reverse=True) + + if cursor is not None: + decoded = decode_cursor(cursor, 2) + cursor_date = decoded[0] + cursor_id = int(decoded[1]) + else: + cursor_date, cursor_id = None, None + filtered = [] + for event in events: + if cursor_date is not None: + key = sort_key(event) + if key[0] > cursor_date or (key[0] == cursor_date and key[1] >= cursor_id): + continue + filtered.append(event) + + page = filtered[:limit] + has_more = len(filtered) > limit + next_cursor = None + if has_more: + last = page[-1] + next_cursor = encode_cursor((last["effective_at"], str(last["ledger_event_id"]))) + + items = [] + for event in page: + item = event_payload( + connection, event, viewer_company_id=viewer_company_id + ) + item["state"] = event["state"] + items.append(item) + return { + "window": {"from": from_, "cutoff": cutoff, "cutoff_inclusive": True}, + "items": items, + "next_cursor": next_cursor, + "has_more": has_more, + } + + +def _current_state(connection: sqlite3.Connection, ledger_event_id: int) -> str: + revision = current_revision(connection, ledger_event_id) + return revision["state"] if revision is not None else "" + + +def event_payload( + connection: sqlite3.Connection, event: sqlite3.Row, *, viewer_company_id: int | None = None +) -> dict[str, object]: + item = { + "ledger_event_id": event["ledger_event_id"], + "ledger_revision_id": event["ledger_revision_id"], + "effective_at": event["effective_at"], + "amount": event["amount"], + "currency": event["currency"], + "payer_company_id": event["payer_company_id"], + "payer_company_name": event["payer_company_name"], + "payee_company_id": event["payee_company_id"], + "payee_company_name": event["payee_company_name"], + "perspective_company_id": event["perspective_company_id"], + "subject_code": event["subject_code"], + "subject_label": subject_label(event["subject_code"]), + "posting_kind": event["posting_kind"], + "source_kind": event["source_kind"], + "reverses_ledger_event_id": event["reverses_ledger_event_id"], + "adjusts_ledger_event_id": event["adjusts_ledger_event_id"], + "source_id": event["source_id"], + "evidence_count": event["evidence_count"], + } + if viewer_company_id is not None: + item["direction"] = viewer_direction(event, viewer_company_id) + own_subject = viewer_subject(event, viewer_company_id) + item["own_subject"] = own_subject + item["own_subject_label"] = ( + subject_label(own_subject) if own_subject is not None else None + ) + item["counterparty_company_id"] = ( + event["payee_company_id"] + if int(event["payer_company_id"]) == int(viewer_company_id) + else event["payer_company_id"] + ) + item["counterparty_company_name"] = ( + event["payee_company_name"] + if int(event["payer_company_id"]) == int(viewer_company_id) + else event["payer_company_name"] + ) + item.update(_event_line_display(connection, event, viewer_company_id)) + return item + + +_REPAY_MARKERS = ("还款", "归还借款", "归还往来款") + + +def bank_short_name(name: str | None) -> str: + text = str(name or "").strip() + if text.startswith("中国"): + text = text[2:] + if text.endswith("银行"): + text = text[:-2] + return text or "银行" + + +def _account_chip(visibility: str, bank_name: str | None, account: str | None) -> dict[str, object]: + if visibility == "missing": + return {"visibility": "missing", "label": None} + if visibility == "masked": + return {"visibility": "masked", "label": "按对方授权不可见"} + number = str(account or "") + tail = number[-4:] if number else "" + short = bank_short_name(bank_name) + label = f"{short} {tail}".strip() if tail else short + return {"visibility": "visible", "label": label} + + +def _event_line_display( + connection: sqlite3.Connection, + event: sqlite3.Row, + viewer_company_id: int | None, +) -> dict[str, object]: + """Account chips, summary and repayment flag for the event table.""" + missing = _account_chip("missing", None, None) + payer_chip, payee_chip = missing, missing + summary = None + texts: list[str] = [] + sides = _event_source_sides(connection, event["ledger_event_id"]) + for side in sides: + company_id = side["company_id"] + own = viewer_company_id is None or int(company_id) == int(viewer_company_id) + visibility = "visible" if own else "masked" + chip = _account_chip(visibility, side.get("bank_name"), side.get("own_account")) + if int(company_id) == int(event["payer_company_id"]): + payer_chip = chip + elif int(company_id) == int(event["payee_company_id"]): + payee_chip = chip + if own or viewer_company_id is None: + if side.get("summary"): + texts.append(str(side["summary"])) + if side.get("reason"): + texts.append(str(side["reason"])) + if texts: + summary = texts[0] + blob = " ".join(texts) + is_repayment = any(marker in blob for marker in _REPAY_MARKERS) + return { + "payer_account": payer_chip, + "payee_account": payee_chip, + "summary": summary, + "is_repayment": is_repayment, + } + + +def _event_source_sides( + connection: sqlite3.Connection, ledger_event_id: int +) -> list[dict[str, object]]: + rows = connection.execute( + """ + SELECT b.company_id AS company_id, s.bank_name AS bank_name, + r.own_account AS own_account, r.summary AS summary, + r.purpose AS purpose, NULL AS reason + FROM ledger_event_bank_sources bs + JOIN current_transfer_decisions cur ON cur.event_id = bs.bank_event_id + JOIN transfer_decision_observations o ON o.decision_id = cur.decision_id + JOIN source_rows r ON r.id = o.source_row_id + JOIN sheet_batches s ON s.id = r.sheet_batch_id + JOIN import_batches b ON b.id = s.import_batch_id + WHERE bs.ledger_event_id = ? + """, + (ledger_event_id,), + ).fetchall() + if rows: + return [_row(row) for row in rows] + manuals = connection.execute( + """ + SELECT m.company_id AS company_id, ba.bank_name AS bank_name, + ba.account_number AS own_account, m.summary AS summary, + NULL AS purpose, m.reason AS reason + FROM ledger_event_manual_sources ms + JOIN manual_records m ON m.id = ms.manual_record_id + LEFT JOIN bank_accounts ba ON ba.id = m.bank_account_id + WHERE ms.ledger_event_id = ? + """, + (ledger_event_id,), + ).fetchall() + return [_row(row) for row in manuals] + + +_DETAIL_SELECT = """ + SELECT p.ledger_event_id, p.id AS ledger_revision_id, p.effective_at, p.amount, + p.amount_scale, p.currency, p.payer_company_id, p.payee_company_id, + p.perspective_company_id, p.subject_code, p.source_kind, + p.posting_kind, p.reverses_ledger_event_id, p.adjusts_ledger_event_id, + p.state, p.revision AS revision_number, + cpayer.name AS payer_company_name, cpayee.name AS payee_company_name, + COALESCE( + (SELECT bs.bank_event_id FROM ledger_event_bank_sources bs + WHERE bs.ledger_event_id = p.ledger_event_id LIMIT 1), + (SELECT ms.manual_record_id FROM ledger_event_manual_sources ms + WHERE ms.ledger_event_id = p.ledger_event_id LIMIT 1) + ) AS source_id, + ((SELECT COUNT(*) FROM ledger_event_bank_sources bs + WHERE bs.ledger_event_id = p.ledger_event_id) + + (SELECT COUNT(*) FROM ledger_event_manual_sources ms + WHERE ms.ledger_event_id = p.ledger_event_id)) AS evidence_count + FROM current_ledger_event_revisions cur + JOIN ledger_event_revisions p ON p.id = cur.revision_id + JOIN companies cpayer ON cpayer.id = p.payer_company_id + JOIN companies cpayee ON cpayee.id = p.payee_company_id +""" + + +def event_detail(connection: sqlite3.Connection, ledger_event_id: int) -> dict[str, object] | None: + event = connection.execute( + _DETAIL_SELECT + " WHERE p.ledger_event_id = ?", (ledger_event_id,) + ).fetchone() + if event is None: + return None + history = connection.execute( + """ + SELECT r.* FROM ledger_event_revisions r + WHERE r.ledger_event_id = ? + ORDER BY r.revision + """, + (ledger_event_id,), + ).fetchall() + suggestions = connection.execute( + """ + SELECT s.*, c.name AS company_name + FROM ledger_subject_suggestions s + LEFT JOIN companies c ON c.id = s.suggested_perspective_company_id + WHERE s.ledger_event_id = ? + ORDER BY s.id + """, + (ledger_event_id,), + ).fetchall() + return { + "event": event_payload(connection, event), + "state": _current_state(connection, ledger_event_id), + "history": [ + { + "revision": rev["revision"], + "state": rev["state"], + "effective_at": rev["effective_at"], + "amount": rev["amount"], + "currency": rev["currency"], + "perspective_company_id": rev["perspective_company_id"], + "subject_code": rev["subject_code"], + "posting_kind": rev["posting_kind"], + "source_kind": rev["source_kind"], + "source_revision_token": rev["source_revision_token"], + "actor_username": rev["actor_username"], + "reason": rev["reason"], + "created_at": rev["created_at"], + } + for rev in history + ], + "suggestions": [ + { + "suggested_perspective_company_id": sug["suggested_perspective_company_id"], + "suggested_company_name": sug["company_name"], + "suggested_subject_code": sug["suggested_subject_code"], + "suggested_subject_label": subject_label(sug["suggested_subject_code"]), + "rule_version": sug["rule_version"], + "evidence": json.loads(sug["evidence_json"]) if sug["evidence_json"] else {}, + "created_at": sug["created_at"], + } + for sug in suggestions + ], + } + + +def subject_review_queue( + connection: sqlite3.Connection, + *, + from_: str, + cutoff: str, + company_id: int | None = None, + limit: int = 50, + cursor: str | None = None, +) -> dict[str, object]: + from_, cutoff = validate_window(from_, cutoff) + conditions = ["r.state = 'pending_subject'", "r.effective_at >= ?", "r.effective_at <= ?"] + params: list[object] = [from_, cutoff + "T23:59:59"] + if company_id is not None: + conditions.append("(r.payer_company_id = ? OR r.payee_company_id = ?)") + params.extend([company_id, company_id]) + rows = connection.execute( + f""" + SELECT r.ledger_event_id, r.effective_at, r.amount, r.currency, + r.payer_company_id, r.payee_company_id, + cpayer.name AS payer_company_name, cpayee.name AS payee_company_name, + r.id AS revision_id, r.evidence_json AS evidence_json + FROM current_ledger_event_revisions cur + JOIN ledger_event_revisions r ON r.id = cur.revision_id + JOIN companies cpayer ON cpayer.id = r.payer_company_id + JOIN companies cpayee ON cpayee.id = r.payee_company_id + WHERE {' AND '.join(conditions)} + ORDER BY r.ledger_event_id + """, + params, + ).fetchall() + rows.sort(key=lambda row: (row["effective_at"], row["ledger_event_id"]), reverse=True) + if cursor is not None: + decoded = decode_cursor(cursor, 2) + cursor_date, cursor_id = decoded[0], int(decoded[1]) + else: + cursor_date, cursor_id = None, None + filtered = [] + for row in rows: + evidence = json.loads(row["evidence_json"] or "{}") if row["evidence_json"] else {} + if evidence.get("admin_disposition") == "exception": + continue + if cursor_date is not None: + if row["effective_at"] > cursor_date or ( + row["effective_at"] == cursor_date and row["ledger_event_id"] >= cursor_id + ): + continue + filtered.append(row) + page = filtered[:limit] + has_more = len(filtered) > limit + next_cursor = None + if has_more: + last = page[-1] + next_cursor = encode_cursor((last["effective_at"], str(last["ledger_event_id"]))) + items = [] + for row in page: + suggestions = connection.execute( + """ + SELECT s.suggested_perspective_company_id, s.suggested_subject_code, + c.name AS company_name, s.rule_version, s.evidence_json + FROM ledger_subject_suggestions s + LEFT JOIN companies c ON c.id = s.suggested_perspective_company_id + WHERE s.ledger_event_id = ? + ORDER BY s.id + """, + (row["ledger_event_id"],), + ).fetchall() + display = _event_line_display(connection, row, None) + items.append( + { + "ledger_event_id": row["ledger_event_id"], + "revision_id": row["revision_id"], + "effective_at": row["effective_at"], + "amount": row["amount"], + "currency": row["currency"], + "payer_company_id": row["payer_company_id"], + "payer_company_name": row["payer_company_name"], + "payee_company_id": row["payee_company_id"], + "payee_company_name": row["payee_company_name"], + "summary": display.get("summary"), + "is_repayment": display.get("is_repayment"), + "suggestions": [ + { + "suggested_perspective_company_id": sug["suggested_perspective_company_id"], + "suggested_company_name": sug["company_name"], + "suggested_subject_code": sug["suggested_subject_code"], + "suggested_subject_label": subject_label(sug["suggested_subject_code"]), + "rule_version": sug["rule_version"], + "evidence": json.loads(sug["evidence_json"]) if sug["evidence_json"] else {}, + } + for sug in suggestions + ], + } + ) + return { + "window": {"from": from_, "cutoff": cutoff, "cutoff_inclusive": True}, + "items": items, + "next_cursor": next_cursor, + "has_more": has_more, + } + + +# --------------------------------------------------------------------------- +# Evidence drill-down +# --------------------------------------------------------------------------- + + +def event_evidence( + connection: sqlite3.Connection, + ledger_event_id: int, + *, + viewer_company_id: int | None = None, +) -> dict[str, object] | None: + """Evidence blocks behind one ledger event with explicit visibility. + + ``visibility`` is always present (``visible``/``masked``/``missing``); a + company user never guesses from absent fields. Company viewers see their + own bank/manual evidence in full, the counterparty side masked, and a + ``missing`` block when no source exists. + """ + event = connection.execute( + _EVENT_SELECT + " WHERE p.ledger_event_id = ?", (ledger_event_id,) + ).fetchone() + if event is None: + return None + blocks: list[dict[str, object]] = [] + bank_claims = connection.execute( + "SELECT * FROM ledger_event_bank_sources WHERE ledger_event_id = ? ORDER BY bank_event_id", + (ledger_event_id,), + ).fetchall() + manual_claims = connection.execute( + "SELECT * FROM ledger_event_manual_sources WHERE ledger_event_id = ? ORDER BY manual_record_id", + (ledger_event_id,), + ).fetchall() + + for claim in bank_claims: + blocks.extend( + _bank_evidence_blocks(connection, claim["bank_event_id"], viewer_company_id) + ) + for claim in manual_claims: + blocks.append( + _manual_evidence_block(connection, claim["manual_record_id"], viewer_company_id) + ) + + if not blocks: + blocks.append( + { + "side": "counterparty", + "source_kind": "bank", + "visibility": "missing", + "fields": {}, + } + ) + + return { + "ledger_event_id": ledger_event_id, + "summary": { + "effective_at": event["effective_at"], + "amount": event["amount"], + "currency": event["currency"], + "payer_company_id": event["payer_company_id"], + "payer_company_name": event["payer_company_name"], + "payee_company_id": event["payee_company_id"], + "payee_company_name": event["payee_company_name"], + "subject_code": event["subject_code"], + "subject_label": subject_label(event["subject_code"]), + "posting_kind": event["posting_kind"], + }, + "blocks": blocks, + } + + +def _bank_evidence_blocks( + connection: sqlite3.Connection, + bank_event_id: int, + viewer_company_id: int | None, +) -> list[dict[str, object]]: + decision = matching._current_decision_for_event(connection, bank_event_id) + if decision is None: + return [] + observations = matching._decision_observations(connection, decision["id"]) + blocks: list[dict[str, object]] = [] + for observation in observations: + source = connection.execute( + """ + SELECT r.id, r.source_row, r.transaction_at, r.income, r.expense, + r.own_account, r.own_name, r.counterparty_account, + r.counterparty_name, r.summary, r.purpose, r.reference, + r.currency, s.sheet_name, f.original_filename, + b.company_id AS batch_company_id, c.name AS company_name + FROM source_rows r + JOIN sheet_batches s ON s.id = r.sheet_batch_id + JOIN import_batches b ON b.id = s.import_batch_id + JOIN source_files f ON f.id = b.source_file_id + LEFT JOIN companies c ON c.id = b.company_id + WHERE r.id = ? + """, + (observation["source_row_id"],), + ).fetchone() + if source is None: + continue + own = int(source["batch_company_id"]) == int(viewer_company_id) if viewer_company_id is not None else True + side = "own" if own else "counterparty" + if viewer_company_id is None or own: + visibility = "visible" + fields = { + "original_filename": source["original_filename"], + "sheet_name": source["sheet_name"], + "source_row": source["source_row"], + "transaction_at": source["transaction_at"], + "income": source["income"], + "expense": source["expense"], + "own_account": source["own_account"] if viewer_company_id is None + else mask_account_number(source["own_account"]) if source["own_account"] else None, + "own_name": source["own_name"], + "counterparty_account_masked": ( + mask_account_number(source["counterparty_account"]) + if source["counterparty_account"] else None + ), + "counterparty_name": source["counterparty_name"], + "summary": source["summary"], + "purpose": source["purpose"], + "reference": source["reference"], + "currency": source["currency"], + "role": observation["role"], + } + else: + visibility = "masked" + fields = { + "company_id": source["batch_company_id"], + "company_name": source["company_name"], + "own_account_masked": ( + mask_account_number(source["own_account"]) + if source["own_account"] else None + ), + "note": "按对方授权不可见", + } + blocks.append( + { + "side": side, + "source_kind": "bank", + "visibility": visibility, + "fields": fields, + } + ) + return blocks + + +def _manual_evidence_block( + connection: sqlite3.Connection, + manual_record_id: int, + viewer_company_id: int | None, +) -> dict[str, object]: + rows = list_records(connection, limit=10000) + record = next((row for row in rows if row["id"] == manual_record_id), None) + if record is None: + return { + "side": "counterparty", "source_kind": "manual", + "visibility": "missing", "fields": {}, + } + own = ( + int(record["company_id"]) == int(viewer_company_id) + if viewer_company_id is not None + else True + ) + if viewer_company_id is None or own: + visibility = "visible" + fields = { + "manual_record_id": record["id"], + "company_id": record["company_id"], + "company_name": record["company_name"], + "counterparty_company_id": record["counterparty_company_id"], + "counterparty_company_name": record["counterparty_company_name"], + "occurred_at": record["occurred_at"], + "direction": record["direction"], + "amount": record["amount"], + "currency": record["currency"], + "funding_source": record["funding_source"], + "requested_subject": record["requested_subject"], + "summary": record["summary"], + "reason": record["reason"], + "state": record["state"], + } + else: + visibility = "masked" + fields = { + "company_id": record["company_id"], + "company_name": record["company_name"], + "state": record["state"], + "note": "按对方授权不可见", + } + return { + "side": "own" if own else "counterparty", + "source_kind": "manual", + "visibility": visibility, + "fields": fields, + } diff --git a/src/bank_importer/subjects.py b/src/bank_importer/subjects.py new file mode 100644 index 0000000..b891331 --- /dev/null +++ b/src/bank_importer/subjects.py @@ -0,0 +1,398 @@ +"""Statutory subject suggestions, mirror mapping and confirmation (B-44). + +Subjects are stored from one participating company's perspective and the +other side is the fixed mirror (应收<->应付, 其他应收<->其他应付), so the two +companies can never record conflicting subjects. Bank summary/purpose text +only ever produces a *suggestion*; nothing here confirms a subject +automatically. Confirmation is an explicit administrator decision carrying +``expected_revision`` and an idempotency key. +""" + +from __future__ import annotations + +import json +import re +import sqlite3 + +from . import matching + +SUBJECTS = ("receivable", "payable", "other_receivable", "other_payable") +SUBJECT_RULE_VERSION = "subject-suggest-draft-v1" + +MIRROR = { + "receivable": "payable", + "payable": "receivable", + "other_receivable": "other_payable", + "other_payable": "other_receivable", +} + +SUBJECT_LABELS = { + "receivable": "应收", + "payable": "应付", + "other_receivable": "其他应收", + "other_payable": "其他应付", +} + +_FULL_WIDTH = str.maketrans( + "ABCDEFGHIJKLMNOPQRSTUVWXYZ" + "abcdefghijklmnopqrstuvwxyz0123456789", + "ABCDEFGHIJKLMNOPQRSTUVWXYZ" + "abcdefghijklmnopqrstuvwxyz0123456789", +) + +# Draft v1 dictionary. Exact-keyword matching only; every hit is a suggestion +# and never an automatic posting. Trade-type keywords are deliberately absent +# until the group supplies an approved dictionary (they always go to review). +_LOAN_LIKE = ("借款", "往来款", "资金往来", "临时借款", "资金调拨", "代垫", "垫付") +_REPAY_LIKE = ("还款", "归还借款", "归还往来款") + + +class SubjectConflictError(ValueError): + """A stale revision or idempotency conflict (mapped to HTTP 409).""" + + +class SubjectInputError(ValueError): + """Invalid input for a subject decision (mapped to HTTP 400/422).""" + + +def mirror_subject(subject_code: str) -> str: + if subject_code not in MIRROR: + raise SubjectInputError("科目必须是应收/应付/其他应收/其他应付之一。") + return MIRROR[subject_code] + + +def subject_label(subject_code: str) -> str: + return SUBJECT_LABELS.get(subject_code, subject_code) + + +def _normalize(text: object) -> str: + return re.sub(r"[\s\ufeff]+", "", str(text or "").translate(_FULL_WIDTH)) + + +def _bank_evidence_texts(connection: sqlite3.Connection, ledger_event_id: int) -> dict[str, str]: + """Purpose/summary text of the B-43 source rows behind a bank event.""" + row = connection.execute( + """ + SELECT bs.bank_event_id FROM ledger_event_bank_sources bs + WHERE bs.ledger_event_id = ? + """, + (ledger_event_id,), + ).fetchone() + if row is None: + return {"purpose": "", "summary": ""} + decision = matching._current_decision_for_event(connection, row["bank_event_id"]) + if decision is None: + return {"purpose": "", "summary": ""} + observations = matching._decision_observations(connection, decision["id"]) + texts: dict[str, list[str]] = {"purpose": [], "summary": []} + for observation in observations: + source = connection.execute( + "SELECT purpose, summary FROM source_rows WHERE id = ?", + (observation["source_row_id"],), + ).fetchone() + if source is None: + continue + for key in ("purpose", "summary"): + value = str(source[key] or "").strip() + if value: + texts[key].append(value) + return { + "purpose": " ".join(texts["purpose"]), + "summary": " ".join(texts["summary"]), + } + + +def compute_suggestions( + connection: sqlite3.Connection, ledger_event_id: int +) -> list[dict[str, object]]: + """ Deterministic draft suggestions for a pending event, never confirmation. + + Purpose rules take precedence over summary rules. When both loan-like and + repay-like keywords match, both candidates are returned as a conflict for + the reviewer; no priority breaks the tie. + """ + from .ledger_events import current_revision + + revision = current_revision(connection, ledger_event_id) + if revision is None or revision["state"] != "pending_subject": + return [] + texts = _bank_evidence_texts(connection, ledger_event_id) + purpose = _normalize(texts["purpose"]) + summary = _normalize(texts["summary"]) + search = purpose or summary + payer = revision["payer_company_id"] + payee = revision["payee_company_id"] + + loan_hit = next((word for word in _LOAN_LIKE if word in search), None) + repay_hit = next((word for word in _REPAY_LIKE if word in search), None) + + suggestions: list[dict[str, object]] = [] + if loan_hit: + suggestions.append( + { + "suggested_perspective_company_id": payer, + "suggested_subject_code": "other_receivable", + "reason": f"匹配建议词典「{loan_hit}」,建议付款方其他应收", + "rule_version": SUBJECT_RULE_VERSION, + "evidence": { + "keyword": loan_hit, + "matched_text": search, + "approved": False, + }, + } + ) + if repay_hit: + suggestions.append( + { + "suggested_perspective_company_id": payee, + "suggested_subject_code": "other_receivable", + "reason": f"匹配建议词典「{repay_hit}」,建议收款方其他应收", + "rule_version": SUBJECT_RULE_VERSION, + "evidence": { + "keyword": repay_hit, + "matched_text": search, + "approved": False, + }, + } + ) + return suggestions + + +def store_suggestions(connection: sqlite3.Connection, ledger_event_id: int) -> int: + """Compute and append suggestions for a pending event. Returns count stored.""" + from .ledger_events import current_revision + + revision = current_revision(connection, ledger_event_id) + if revision is None or revision["state"] != "pending_subject": + return 0 + stored = 0 + for suggestion in compute_suggestions(connection, ledger_event_id): + from .db import utc_now + + connection.execute( + """ + INSERT INTO ledger_subject_suggestions ( + ledger_event_id, source_revision_id, + suggested_perspective_company_id, suggested_subject_code, + rule_version, evidence_json, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?) + """, + ( + ledger_event_id, revision["id"], + suggestion["suggested_perspective_company_id"], + suggestion["suggested_subject_code"], + suggestion["rule_version"], + json.dumps(suggestion.get("evidence", {}), ensure_ascii=False), + utc_now(), + ), + ) + stored += 1 + return stored + + +def confirm_subject( + connection: sqlite3.Connection, + ledger_event_id: int, + *, + perspective_company_id: int, + subject_code: str, + reason: str, + expected_revision: int | None, + request_key: str | None, + actor: sqlite3.Row, +) -> dict[str, object]: + """Confirm a subject, turning a pending event into a confirmed revision.""" + from .ledger_events import append_revision, current_revision + + reason = (reason or "").strip() + if not reason: + raise SubjectInputError("必须填写科目确认依据。") + if subject_code not in SUBJECTS: + raise SubjectInputError("科目必须是应收/应付/其他应收/其他应付之一。") + + began = False + if not connection.in_transaction: + connection.execute("BEGIN IMMEDIATE") + began = True + try: + if request_key: + existing = connection.execute( + """ + SELECT * FROM ledger_event_revisions + WHERE ledger_event_id = ? AND idempotency_key = ? + ORDER BY id LIMIT 1 + """, + (ledger_event_id, request_key), + ).fetchone() + if existing is not None: + if began: + connection.commit() + return _revision_payload(connection, existing) + + current = current_revision(connection, ledger_event_id) + if current is None: + raise SubjectConflictError("该事件不存在或没有当前修订。") + if current["state"] != "pending_subject": + raise SubjectConflictError("只有待确认科目的事件可以确认科目。") + # ``expected_revision`` may be the revision row id (what the API/UI + # sends as ``ledger_revision_id``) or the per-event sequence number; + # both identify the exact revision the client saw. + if expected_revision is not None and int(expected_revision) not in ( + current["id"], current["revision"], + ): + raise SubjectConflictError("事件已发生变更,请刷新后重试。") + participants = {current["payer_company_id"], current["payee_company_id"]} + if perspective_company_id not in participants: + raise SubjectInputError("视角公司必须是事件参与方。") + + revision_id = append_revision( + connection, + ledger_event_id, + state="confirmed", + effective_at=current["effective_at"], + amount=current["amount"], + currency=current["currency"], + payer_company_id=current["payer_company_id"], + payee_company_id=current["payee_company_id"], + perspective_company_id=perspective_company_id, + subject_code=subject_code, + source_kind=current["source_kind"], + source_revision_token=current["source_revision_token"], + posting_kind=current["posting_kind"], + reverses_ledger_event_id=current["reverses_ledger_event_id"], + adjusts_ledger_event_id=current["adjusts_ledger_event_id"], + rule_version=current["rule_version"] or SUBJECT_RULE_VERSION, + evidence_json=current["evidence_json"], + idempotency_key=request_key, + actor=actor, + reason=reason, + supersedes_revision_id=current["id"], + ) + row = connection.execute( + "SELECT * FROM ledger_event_revisions WHERE id = ?", (revision_id,) + ).fetchone() + except Exception: + if began: + connection.rollback() + raise + else: + if began: + connection.commit() + return _revision_payload(connection, row) + + +def park_subject( + connection: sqlite3.Connection, + ledger_event_id: int, + *, + disposition: str, + reason: str, + expected_revision: int | None, + request_key: str | None, + actor: sqlite3.Row, +) -> dict[str, object]: + """Record 退回/转异常 without confirming a statutory subject. + + The event stays ``pending_subject`` so it never enters confirmed balances. + ``exception`` is hidden from the active review queue; ``return`` remains + visible so the company can supplement materials. + """ + from .ledger_events import append_revision, current_revision + + if disposition not in ("return", "exception"): + raise SubjectInputError("科目处理只能是退回或转异常。") + reason = (reason or "").strip() + if not reason: + raise SubjectInputError("必须填写处理依据。") + + began = False + if not connection.in_transaction: + connection.execute("BEGIN IMMEDIATE") + began = True + try: + if request_key: + existing = connection.execute( + """ + SELECT * FROM ledger_event_revisions + WHERE ledger_event_id = ? AND idempotency_key = ? + ORDER BY id LIMIT 1 + """, + (ledger_event_id, request_key), + ).fetchone() + if existing is not None: + if began: + connection.commit() + return _revision_payload(connection, existing) + + current = current_revision(connection, ledger_event_id) + if current is None: + raise SubjectConflictError("该事件不存在或没有当前修订。") + if current["state"] != "pending_subject": + raise SubjectConflictError("只有待确认科目的事件可以退回或转异常。") + if expected_revision is not None and int(expected_revision) not in ( + current["id"], current["revision"], + ): + raise SubjectConflictError("事件已发生变更,请刷新后重试。") + evidence = json.loads(current["evidence_json"] or "{}") if current["evidence_json"] else {} + evidence["admin_disposition"] = disposition + revision_id = append_revision( + connection, + ledger_event_id, + state="pending_subject", + effective_at=current["effective_at"], + amount=current["amount"], + currency=current["currency"], + payer_company_id=current["payer_company_id"], + payee_company_id=current["payee_company_id"], + perspective_company_id=None, + subject_code=None, + source_kind=current["source_kind"], + source_revision_token=current["source_revision_token"], + posting_kind=current["posting_kind"], + reverses_ledger_event_id=current["reverses_ledger_event_id"], + adjusts_ledger_event_id=current["adjusts_ledger_event_id"], + rule_version=current["rule_version"] or SUBJECT_RULE_VERSION, + evidence_json=json.dumps(evidence, ensure_ascii=False), + idempotency_key=request_key, + actor=actor, + reason=reason, + supersedes_revision_id=current["id"], + ) + row = connection.execute( + "SELECT * FROM ledger_event_revisions WHERE id = ?", (revision_id,) + ).fetchone() + except Exception: + if began: + connection.rollback() + raise + else: + if began: + connection.commit() + return _revision_payload(connection, row) + + +def _revision_payload(connection: sqlite3.Connection, revision: sqlite3.Row) -> dict[str, object]: + company = connection.execute( + "SELECT name FROM companies WHERE id = ?", (revision["perspective_company_id"],) + ).fetchone() + return { + "ledger_event_id": revision["ledger_event_id"], + "revision_id": revision["id"], + "revision": revision["revision"], + "state": revision["state"], + "effective_at": revision["effective_at"], + "amount": revision["amount"], + "currency": revision["currency"], + "payer_company_id": revision["payer_company_id"], + "payee_company_id": revision["payee_company_id"], + "perspective_company_id": revision["perspective_company_id"], + "perspective_company_name": company["name"] if company else None, + "subject_code": revision["subject_code"], + "subject_label": subject_label(revision["subject_code"]) + if revision["subject_code"] + else None, + "posting_kind": revision["posting_kind"], + "source_kind": revision["source_kind"], + "reason": revision["reason"], + "created_at": revision["created_at"], + } diff --git a/tests/b44_ui_check.js b/tests/b44_ui_check.js new file mode 100644 index 0000000..b801e9a --- /dev/null +++ b/tests/b44_ui_check.js @@ -0,0 +1,86 @@ +"use strict"; + +const path = require("path"); +const assert = require("assert"); +const ui = require(path.join(__dirname, "..", "web", "app.js")); + +assert.strictEqual(ui.fmtAbsMoney(-1280), "1,280.00"); +assert.strictEqual(ui.fmtAbsMoney(1280), "1,280.00"); +assert.strictEqual(ui.fmtAbsMoney("60.5"), "60.50"); + +assert.strictEqual(ui.eventIsNegative({ posting_kind: "reversal" }), true); +assert.strictEqual(ui.eventIsNegative({ posting_kind: "normal", is_repayment: true }), true); +assert.strictEqual(ui.eventIsNegative({ posting_kind: "normal", is_repayment: false }), false); + +assert.strictEqual(ui.cycleTab(0, 4, false), 1); +assert.strictEqual(ui.cycleTab(3, 4, false), 0); +assert.strictEqual(ui.cycleTab(0, 4, true), 3); +assert.strictEqual(ui.cycleTab(2, 5, true), 1); +assert.strictEqual(ui.cycleTab(0, 0, false), 0); + +assert.strictEqual(ui.drawerEscAction(1), "close"); +assert.strictEqual(ui.drawerEscAction(0), "close"); +assert.strictEqual(ui.drawerEscAction(2), "back"); +assert.strictEqual(ui.drawerEscAction(3), "back"); + +const abs = ui.amountWithCurrency(1280, "CNY"); +assert.ok(abs.includes("CNY")); +assert.ok(abs.includes("1,280.00")); +assert.ok(!abs.includes("+")); +assert.ok(!abs.includes("−")); + +const repay = ui.amountWithCurrency(3200000, "CNY", { signed: true, negative: true }); +assert.ok(repay.includes("−")); +assert.ok(!repay.includes("+")); +assert.ok(repay.includes("3,200,000.00")); + +assert.strictEqual(ui.resultDirection(10).label, "应收"); +assert.strictEqual(ui.resultDirection(-10).label, "应付"); +assert.strictEqual(ui.resultDirection(0).label, "持平"); + +assert.strictEqual(ui.isCompactAmount(1280), false); +assert.strictEqual(ui.isCompactAmount(999999999), false); +assert.strictEqual(ui.isCompactAmount(1000000000), true); +assert.strictEqual(ui.isCompactAmount("123456789012345"), true); +assert.ok(ui.amountWithCurrency("123456789012345", "CNY").includes("is-compact")); +assert.ok(!ui.amountWithCurrency(1280, "CNY").includes("is-compact")); + +assert.strictEqual(ui.cashDirectionLabel("outgoing"), "转出"); +assert.strictEqual(ui.cashDirectionLabel("incoming"), "转入"); +assert.strictEqual(ui.relatedFlowLabel(null), "无关联流水"); +assert.strictEqual(ui.relatedFlowLabel(""), "无关联流水"); +assert.strictEqual(ui.relatedFlowLabel("42"), "42"); + +const subjectFields = ui.auditEvidenceFields("subject-review", { + direction: "outgoing", + effectiveAt: "2026-07-18T09:00:00", + amount: "600", + currency: "CNY", + summary: "资金调拨", +}); +assert.deepStrictEqual(subjectFields.map((item) => item[0]), ["方向", "日期", "金额", "摘要"]); +assert.strictEqual(subjectFields[0][1], "转出"); + +const manualFields = ui.auditEvidenceFields("manual-review", { + counterpartyName: "乙公司", + direction: "incoming", + relatedSourceRowId: "", + submittedBy: "出纳甲", + attachmentName: "", + amount: "80", + currency: "CNY", + summary: "补记", +}); +assert.deepStrictEqual( + manualFields.map((item) => item[0]), + ["对方公司名", "方向", "关联流水号", "提交人", "附件", "摘要", "金额"], +); +assert.strictEqual(manualFields[0][1], "乙公司"); +assert.strictEqual(manualFields[1][1], "转入"); +assert.strictEqual(manualFields[2][1], "无关联流水"); + +assert.ok(ui.accountCell({ visibility: "visible", label: "中信 5316" }).includes("中信 5316")); +assert.ok(ui.accountCell({ visibility: "masked" }).includes("按对方授权不可见")); +assert.ok(ui.accountCell({ visibility: "missing" }).includes("源行缺失")); + +console.log("b44_ui_check ok"); diff --git a/tests/fixtures/b44-balance-directory.html b/tests/fixtures/b44-balance-directory.html new file mode 100644 index 0000000..541c33d --- /dev/null +++ b/tests/fixtures/b44-balance-directory.html @@ -0,0 +1,138 @@ + + +
+ + +公司间往来余额目录
每行余额都附带截止日、期初状态、本期借贷、结果与未决金额
按公司对查询双方口径、科目和逐笔凭证
公司间往来余额目录、公司对明细与逐层追溯
统计口径 2026.01.01—2026.07.31
| 交易日期 | 方向 | 科目 | 本方账户 | 对方账户 | 摘要 | 匹配 | 金额(万元) |
|---|---|---|---|---|---|---|---|
| 2026.07.18 | 转出 | 应收 | 中信 · 5316 | 工行 · 9481 | 往来款 | 双边匹配 | 1,000.00 |
| 2026.07.06 | 转入 | 应收 | 中信 · 5316 | 工行 · 9481 | 归还往来款 | 双边匹配 | −320.00 |
| 2026.06.27 | 转出 | 其他应收 | 建行 · 0845 | 工行 · 9481 | 资金调拨 | 单边待核 | 600.00 |
每行余额都附带截止日、期初状态、本期借贷、结果与未决金额