fix: synchronize document links with reconciliation items

This commit is contained in:
plx
2026-08-14 13:44:37 +00:00
parent 7d138eccdc
commit 475ff15b9b
5 changed files with 256 additions and 1 deletions

View File

@@ -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,

View File

@@ -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,

View File

@@ -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]

View File

@@ -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()

View File

@@ -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