diff --git a/app/document_reconciliation_service.py b/app/document_reconciliation_service.py index 4cb5f9d..15068ac 100644 --- a/app/document_reconciliation_service.py +++ b/app/document_reconciliation_service.py @@ -372,6 +372,55 @@ def _dual_write(conn: Any, document_id: str, opportunity_id: str, relationship: "role": role, "primary": primary, "active": active}) +def _sync_reconciliation_item_for_active_link( + conn: Any, document_id: str, opportunity_id: str, relationship: str, +) -> int: + """Project an active document link onto its exact reconciliation item. + + Document identity is intentionally limited to source system plus the stored + document number/external id. Customer and payload inference do not belong in + this consistency boundary. + """ + relationship = str(relationship or "").upper() + if relationship not in {"PRIMARY", "SECONDARY"}: + return 0 + result = conn.execute(text(""" + UPDATE reconciliation_items ri + SET status = 'linked', + opportunity_id = CAST(:opportunity_id AS UUID), + resolved_at = COALESCE(ri.resolved_at, now()), + updated_at = now() + FROM commercial_documents d + WHERE d.id = CAST(:document_id AS UUID) + AND ri.source_system = d.system + AND ri.status NOT IN ('ignored', 'historical') + AND ( + (NULLIF(d.document_number, '') IS NOT NULL AND ri.document_number = d.document_number) + OR (NULLIF(d.external_id, '') IS NOT NULL AND ri.external_id = d.external_id) + ) + """), {"document_id": document_id, "opportunity_id": opportunity_id}) + return int(result.rowcount or 0) + + +def active_document_link_exclusion_sql(item_alias: str = "ri") -> str: + """SQL predicate hiding stale open items already owned by an active link.""" + if not item_alias.replace("_", "").isalnum(): + raise ValueError("invalid SQL alias") + return f"""NOT EXISTS ( + SELECT 1 + FROM commercial_documents linked_cd + JOIN opportunity_document_links linked_odl + ON linked_odl.document_id = linked_cd.id + AND linked_odl.ended_at IS NULL + AND linked_odl.relationship IN ('PRIMARY','SECONDARY') + WHERE linked_cd.system = {item_alias}.source_system + AND ( + (NULLIF({item_alias}.document_number, '') IS NOT NULL AND linked_cd.document_number = {item_alias}.document_number) + OR (NULLIF({item_alias}.external_id, '') IS NOT NULL AND linked_cd.external_id = {item_alias}.external_id) + ) + )""" + + def _command_fingerprint(*, opportunity_id: str, document_id: str, target_relationship: str, origin_opportunity_id: Optional[str] = None, destination_opportunity_id: Optional[str] = None, @@ -468,6 +517,9 @@ def set_document_relationship(opportunity_id: str, document_id: str, relationshi existing = _validate_membership(conn, opportunity_id, document_id) old = existing.get("relationship") if existing else None if old == relationship: + _sync_reconciliation_item_for_active_link( + conn, document_id, opportunity_id, relationship, + ) result = dict(existing) _complete_idempotency(conn, idempotency_key, result) return result @@ -505,6 +557,9 @@ def set_document_relationship(opportunity_id: str, document_id: str, relationshi VALUES(CAST(:opportunity_id AS UUID),CAST(:document_id AS UUID),:kind,:relationship,:manual, :reason,:actor,now(),:source,:correlation_id,:request_id,CAST(:metadata AS JSONB)) RETURNING *"""), params).mappings().first() _dual_write(conn, document_id, opportunity_id, relationship) + _sync_reconciliation_item_for_active_link( + conn, document_id, opportunity_id, relationship, + ) conn.execute(text("UPDATE commercial_document_lines SET opportunity_document_link_id=:link_id WHERE commercial_document_id=CAST(:document_id AS UUID)"), {"link_id": link["id"], "document_id": document_id}) _event(conn, link_id=str(link["id"]), opportunity_id=opportunity_id, document_id=document_id, event_type=event_type, actor=actor, reason=reason, old_relationship=old, @@ -576,6 +631,10 @@ def reassign_document(origin_opportunity_id: str, destination_opportunity_id: st existing_destination = conn.execute(text("""SELECT * FROM opportunity_document_links WHERE opportunity_id=CAST(:oid AS UUID) AND document_id=CAST(:did AS UUID) AND ended_at IS NULL"""), {"oid": destination_opportunity_id, "did": document_id}).mappings().first() if existing_destination: + _sync_reconciliation_item_for_active_link( + conn, document_id, destination_opportunity_id, + str(existing_destination.get("relationship") or ""), + ) result = dict(existing_destination) _complete_idempotency(conn, idempotency_key, result) return result @@ -594,6 +653,9 @@ def reassign_document(origin_opportunity_id: str, destination_opportunity_id: st "reason": reason, "actor": actor, "origin": origin_opportunity_id, "correlation_id": correlation_id, "request_id": request_id, "old_link": str(old["id"])}).mappings().first() _dual_write(conn, document_id, destination_opportunity_id, "SECONDARY") + _sync_reconciliation_item_for_active_link( + conn, document_id, destination_opportunity_id, "SECONDARY", + ) conn.execute(text("UPDATE commercial_document_lines SET opportunity_document_link_id=:link WHERE commercial_document_id=CAST(:doc AS UUID)"), {"link": link["id"], "doc": document_id}) _event(conn, link_id=str(old["id"]), opportunity_id=origin_opportunity_id, document_id=document_id, event_type="REASSIGNED", actor=actor, reason=reason, old_relationship=old["relationship"], @@ -644,6 +706,10 @@ def classify_sync_document(conn: Any, opportunity_id: str, document_id: str, doc correlation_id=None, request_id=None, idempotency_key=f"jasmin-primary-cancelled:{opportunity_id}:{document_id}") result["relationship"] = "REVIEW_REQUIRED" + else: + _sync_reconciliation_item_for_active_link( + conn, document_id, opportunity_id, str(link.get("relationship") or ""), + ) return result candidates = conn.execute(text("""SELECT * FROM opportunity_document_links WHERE opportunity_id=CAST(:oid AS UUID) AND document_kind=:kind AND ended_at IS NULL FOR UPDATE"""), @@ -671,6 +737,9 @@ def classify_sync_document(conn: Any, opportunity_id: str, document_id: str, doc RETURNING *"""), {"oid": opportunity_id, "did": document_id, "kind": document_kind, "relationship": relationship, "reason": reason, "actor": actor}).mappings().first() _dual_write(conn, document_id, opportunity_id, relationship) + _sync_reconciliation_item_for_active_link( + conn, document_id, opportunity_id, relationship, + ) _event(conn, link_id=str(new_link["id"]), opportunity_id=opportunity_id, document_id=document_id, event_type="SYNC_LINK_CLASSIFIED", actor=actor, reason=reason, old_relationship=None, new_relationship=relationship, correlation_id=None, request_id=None, diff --git a/app/jasmin_backfill_service.py b/app/jasmin_backfill_service.py index d2318f0..a4599df 100644 --- a/app/jasmin_backfill_service.py +++ b/app/jasmin_backfill_service.py @@ -14,6 +14,7 @@ from sqlalchemy import text from app.company_opportunity_linking import is_public_email_domain from app.db import engine +from app.document_reconciliation_service import active_document_link_exclusion_sql from app.reconciliation_service import ( _apply_jasmin_documents_to_opportunity, # noqa: PLC2701 - deliberate operator maintenance helper _jasmin_document_lines_from_item, # noqa: PLC2701 @@ -540,6 +541,7 @@ def find_jasmin_document_candidates_for_opportunity( AND COALESCE(ri.customer_tax_id, '') <> '' AND ri.customer_tax_id <> CAST(:fiscal_tax_id AS TEXT) ) + AND {active_link_exclusion} AND NOT EXISTS ( SELECT 1 FROM commercial_documents cd WHERE cd.opportunity_id = CAST(:opportunity_id AS UUID) @@ -552,7 +554,7 @@ def find_jasmin_document_candidates_for_opportunity( AND ri.external_id IS NOT NULL AND cd.external_id = ri.external_id ) - """ + """.format(active_link_exclusion=active_document_link_exclusion_sql("ri")) params = { "opportunity_id": opportunity_id, "fiscal_customer_id": fiscal_customer_id, diff --git a/app/reconciliation_service.py b/app/reconciliation_service.py index e3f6351..5c2a0a0 100644 --- a/app/reconciliation_service.py +++ b/app/reconciliation_service.py @@ -575,6 +575,9 @@ def list_reconciliation_items(*, status: str = "open", external_type: Optional[s if status != "all": where.append("ri.status = :status") params["status"] = status + from app.document_reconciliation_service import active_document_link_exclusion_sql + where.append("(ri.status NOT IN ('open','needs_review','conflict') OR " + + active_document_link_exclusion_sql("ri") + ")") if external_type: where.append("ri.external_type = :external_type") params["external_type"] = external_type @@ -2718,6 +2721,7 @@ def _get_reconciliation_items_by_ids(item_ids: List[str]) -> List[Dict[str, Any] cleaned = [str(x).strip() for x in item_ids if str(x).strip()] if not cleaned: return [] + from app.document_reconciliation_service import active_document_link_exclusion_sql with engine.begin() as conn: rows = conn.execute(text(""" SELECT id::text, source_system, external_type, external_id, title, description, @@ -2728,6 +2732,7 @@ def _get_reconciliation_items_by_ids(item_ids: List[str]) -> List[Dict[str, Any] FROM reconciliation_items WHERE id = ANY(CAST(:ids AS UUID[])) AND status IN ('open','needs_review','conflict') + AND """ + active_document_link_exclusion_sql("reconciliation_items") + """ ORDER BY document_date NULLS FIRST, created_at """), {"ids": cleaned}).mappings().all() return [dict(row) for row in rows] diff --git a/scripts/repair_document_reconciliation_item_links.py b/scripts/repair_document_reconciliation_item_links.py new file mode 100644 index 0000000..465f613 --- /dev/null +++ b/scripts/repair_document_reconciliation_item_links.py @@ -0,0 +1,79 @@ +"""Repair stale reconciliation_items for active document links. + +Dry-run by default. Only exact document identity is used, and ambiguous items +matching active links in more than one opportunity are deliberately skipped. +""" +from __future__ import annotations + +import argparse +import sys +from pathlib import Path + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +if str(PROJECT_ROOT) not in sys.path: + sys.path.insert(0, str(PROJECT_ROOT)) + +from sqlalchemy import text + +from app.db import engine + + +MATCHES_CTE = """ +WITH exact_matches AS ( + SELECT ri.id AS reconciliation_item_id, odl.opportunity_id + FROM reconciliation_items ri + JOIN commercial_documents cd + ON cd.system = ri.source_system + AND ( + (NULLIF(ri.document_number, '') IS NOT NULL AND cd.document_number = ri.document_number) + OR (NULLIF(ri.external_id, '') IS NOT NULL AND cd.external_id = ri.external_id) + ) + JOIN opportunity_document_links odl + ON odl.document_id = cd.id + AND odl.ended_at IS NULL + AND odl.relationship IN ('PRIMARY','SECONDARY') + WHERE ri.status NOT IN ('ignored','historical') +), repairable AS ( + SELECT reconciliation_item_id, + (array_agg(DISTINCT opportunity_id))[1] AS opportunity_id + FROM exact_matches + GROUP BY reconciliation_item_id + HAVING COUNT(DISTINCT opportunity_id) = 1 +) +""" + + +def repair(*, apply: bool = False) -> dict[str, int]: + with engine.begin() as conn: + candidates = int(conn.execute(text(MATCHES_CTE + """ + SELECT COUNT(*) FROM repairable r + JOIN reconciliation_items ri ON ri.id = r.reconciliation_item_id + WHERE ri.status <> 'linked' + OR ri.opportunity_id IS DISTINCT FROM r.opportunity_id + OR ri.resolved_at IS NULL + """)).scalar() or 0) + updated = 0 + if apply and candidates: + updated = int(conn.execute(text(MATCHES_CTE + """ + UPDATE reconciliation_items ri + SET status='linked', opportunity_id=r.opportunity_id, + resolved_at=COALESCE(ri.resolved_at, now()), updated_at=now() + FROM repairable r + WHERE ri.id=r.reconciliation_item_id + AND (ri.status <> 'linked' + OR ri.opportunity_id IS DISTINCT FROM r.opportunity_id + OR ri.resolved_at IS NULL) + """)).rowcount or 0) + return {"candidates": candidates, "updated": updated} + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--apply", action="store_true", help="Apply exact, unambiguous repairs") + args = parser.parse_args() + result = repair(apply=args.apply) + print(f"candidates={result['candidates']} updated={result['updated']} apply={args.apply}") + + +if __name__ == "__main__": + main() diff --git a/tests/test_document_reconciliation_item_consistency.py b/tests/test_document_reconciliation_item_consistency.py new file mode 100644 index 0000000..5d16623 --- /dev/null +++ b/tests/test_document_reconciliation_item_consistency.py @@ -0,0 +1,100 @@ +from inspect import getsource +from pathlib import Path + +import app.document_reconciliation_service as service +import app.reconciliation_service as reconciliation_service + + +class _Result: + def __init__(self, rowcount=1): + self.rowcount = rowcount + + +class _Connection: + def __init__(self): + self.calls = [] + + def execute(self, statement, params=None): + self.calls.append((str(statement), params or {})) + return _Result() + + +def test_primary_projects_exact_document_identity_to_linked(): + connection = _Connection() + assert service._sync_reconciliation_item_for_active_link( + connection, "doc-primary", "opp-primary", "PRIMARY", + ) == 1 + sql, params = connection.calls[0] + assert "status = 'linked'" in sql + assert "resolved_at = COALESCE(ri.resolved_at, now())" in sql + assert "ri.source_system = d.system" in sql + assert "ri.document_number = d.document_number" in sql + assert "ri.external_id = d.external_id" in sql + assert "payload" not in sql.lower() + assert "customer" not in sql.lower() + assert params == {"document_id": "doc-primary", "opportunity_id": "opp-primary"} + + +def test_secondary_and_reassign_project_destination(): + connection = _Connection() + service._sync_reconciliation_item_for_active_link( + connection, "doc-secondary", "opp-destination", "SECONDARY", + ) + assert connection.calls[0][1]["opportunity_id"] == "opp-destination" + reassign_source = getsource(service.reassign_document) + assert "destination_opportunity_id, \"SECONDARY\"" in reassign_source + assert "_sync_reconciliation_item_for_active_link" in reassign_source + + +def test_source_system_is_mandatory_and_cannot_cross_match(): + connection = _Connection() + service._sync_reconciliation_item_for_active_link(connection, "doc", "opp", "PRIMARY") + sql = connection.calls[0][0] + assert "ri.source_system = d.system" in sql + assert "source_system" in sql + + +def test_ignored_historical_and_removed_do_not_change_items(): + for relationship in ("IGNORED", "HISTORICAL", "REMOVED", "REVIEW_REQUIRED", "REASSIGNED"): + connection = _Connection() + assert service._sync_reconciliation_item_for_active_link( + connection, "doc", "opp", relationship, + ) == 0 + assert connection.calls == [] + + +def test_set_relationship_and_auto_classification_use_central_projection(): + set_source = getsource(service.set_document_relationship) + classify_source = getsource(service.classify_sync_document) + assert set_source.count("_sync_reconciliation_item_for_active_link") >= 2 + assert classify_source.count("_sync_reconciliation_item_for_active_link") >= 2 + assert "if old == relationship" in set_source + + +def test_read_safeguard_hides_stale_open_items_with_active_links(): + predicate = service.active_document_link_exclusion_sql("ri") + assert "ended_at IS NULL" in predicate + assert "relationship IN ('PRIMARY','SECONDARY')" in predicate + assert "linked_cd.system = ri.source_system" in predicate + list_source = getsource(reconciliation_service.list_reconciliation_items) + by_ids_source = getsource(reconciliation_service._get_reconciliation_items_by_ids) + assert "active_document_link_exclusion_sql" in list_source + assert "active_document_link_exclusion_sql" in by_ids_source + + +def test_unlink_reopen_policy_is_preserved(): + source = Path("app/admin_ui/pages/opportunities.py").read_text() + assert "status = CASE WHEN status IN ('resolved','linked','applied','open','needs_review','conflict') THEN 'needs_review' ELSE status END" in source + assert "resolved_at = NULL" in source + assert "reconciliation_items_unlinked" in source + + +def test_repair_is_dry_run_safe_exact_and_unambiguous(): + source = Path("scripts/repair_document_reconciliation_item_links.py").read_text() + assert "apply: bool = False" in source + assert "--apply" in source + assert "COUNT(DISTINCT opportunity_id) = 1" in source + assert "cd.system = ri.source_system" in source + assert "payload::text" not in source + assert "customer_id" not in source + assert "ri.status NOT IN ('ignored','historical')" in source