Files
clientflow_backend/tests/test_document_reconciliation_v2.py
2026-08-05 23:55:33 +00:00

371 lines
18 KiB
Python

import os
from pathlib import Path
import subprocess
import sys
import pytest
from app.document_reconciliation_backfill import classify_group, select_complete_group_batch
from app.document_reconciliation_service import (
_command_fingerprint,
resolve_document_links,
select_valid_primary,
)
from scripts.migrate_document_reconciliation_v2 import actions_for_only_unambiguous
from scripts.migrate_document_reconciliation_v2 import apply_actions
from scripts.preflight_document_reconciliation_v2 import run_preflight
def _doc(doc_id: str, *, role="related", primary=False, status="OPEN", external_id=None, payload=None):
return {"id": doc_id, "opportunity_id": "00000000-0000-0000-0000-000000000001",
"document_kind": "quotation", "role": role, "is_primary": primary,
"status": status, "external_id": external_id or doc_id, "payload": payload or {}}
class _Rows:
def __init__(self, rows=None, scalar=None):
self._rows, self._scalar = rows or [], scalar
def scalar(self): return self._scalar
def mappings(self): return self
def all(self): return self._rows
class _ResolverConnection:
def __init__(self, legacy, v2=None, schema=True):
self.legacy, self.v2, self.schema = legacy, v2 or [], schema
def execute(self, statement, params=None):
sql = str(statement)
if "to_regclass" in sql:
return _Rows(scalar=self.schema)
if "FROM commercial_documents d" in sql and "opportunity_document_links" not in sql:
return _Rows(self.legacy)
if "FROM opportunity_document_links l" in sql:
return _Rows(self.v2)
raise AssertionError(sql)
def _legacy_row(doc_id, kind="quotation", primary=True):
return {"id": doc_id, "opportunity_id": "o", "document_kind": kind,
"role": "current" if primary else "related", "is_primary": primary,
"is_active": True, "status": "OPEN", "updated_at": "1"}
def test_resolver_works_without_v2_schema_and_with_empty_schema():
legacy = [_legacy_row("a")]
without = resolve_document_links("o", conn=_ResolverConnection(legacy, schema=False))
empty = resolve_document_links("o", conn=_ResolverConnection(legacy, v2=[], schema=True))
assert without[0]["resolution_source"] == "legacy_rollout"
assert empty[0]["resolution_source"] == "legacy_rollout"
def test_resolver_uses_legacy_for_partial_group_and_v2_after_complete_group():
legacy = [_legacy_row("a"), _legacy_row("b", primary=False)]
partial_v2 = [{"document_id": "a", "document_kind": "quotation", "relationship": "PRIMARY",
"ended_at": None, "updated_at": "2"}]
partial = resolve_document_links("o", conn=_ResolverConnection(legacy, partial_v2))
assert len(partial) == 2 and {row["resolution_source"] for row in partial} == {"legacy_rollout"}
complete_v2 = partial_v2 + [{"document_id": "b", "document_kind": "quotation",
"relationship": "SECONDARY", "ended_at": None, "updated_at": "2"}]
complete = resolve_document_links("o", conn=_ResolverConnection(legacy, complete_v2))
assert len(complete) == 2 and {row["resolution_source"] for row in complete} == {"v2"}
def test_resolver_contract_keeps_document_id_separate_from_link_id():
legacy = [_legacy_row("document-a")]
v2 = [{"id": "document-a", "document_id": "document-a", "link_id": "link-z",
"document_kind": "quotation", "relationship": "PRIMARY",
"ended_at": None, "updated_at": "2"}]
resolved = resolve_document_links("o", conn=_ResolverConnection(legacy, v2))
assert resolved[0]["id"] == "document-a"
assert resolved[0]["document_id"] == "document-a"
assert resolved[0]["link_id"] == "link-z"
fallback = resolve_document_links("o", conn=_ResolverConnection(legacy, schema=False))
assert fallback[0]["id"] == fallback[0]["document_id"] == "document-a"
assert fallback[0]["link_id"] is None
def test_v2_projection_contains_ui_document_fields_and_explicit_ids():
source = Path("app/document_reconciliation_service.py").read_text()
projection = source[source.index('_LINK_SELECT = """'):source.index('def list_document_links')]
assert "d.id::text AS id" in projection
assert "d.id::text AS document_id" in projection
assert "l.id::text AS link_id" in projection
for field in ("customer_id", "company", "document_type", "serie", "series_number",
"document_number", "version_number", "external_id", "external_url", "role",
"is_primary", "is_active", "status", "amount", "tax_amount", "total_amount",
"currency", "document_date", "due_date", "payload"):
assert f"d.{field}" in projection, field
def test_backfill_single_primary_is_unambiguous():
result = classify_group([_doc("a", role="current", primary=True), _doc("b")])
assert [r["relationship"] for r in result] == ["PRIMARY", "SECONDARY"]
def test_backfill_two_primaries_requires_review():
result = classify_group([_doc("a", role="current", primary=True), _doc("b", role="accepted", primary=True)])
assert {r["relationship"] for r in result} == {"REVIEW_REQUIRED"}
def test_backfill_cancelled_primary_requires_review_and_never_ignored():
result = classify_group([_doc("a", role="current", primary=True, status="cancelled")])
assert result[0]["relationship"] == "REVIEW_REQUIRED"
assert result[0]["cancelled"] is True
def test_detached_without_evidence_is_not_invented_removed():
assert classify_group([_doc("a", role="detached")])[0]["relationship"] == "REVIEW_REQUIRED"
assert classify_group([_doc("a", role="detached", payload={"removed_reason": "manual"})])[0]["relationship"] == "REMOVED"
def test_classifier_import_is_isolated_from_application_configuration():
env = os.environ.copy()
env.pop("DATABASE_URL", None)
env.pop("OPENROUTER_API_KEY", None)
result = subprocess.run(
[
sys.executable,
"-c",
"import sys; from app.document_reconciliation_backfill import classify_group; "
"assert callable(classify_group); "
"assert 'app.config' not in sys.modules; "
"assert 'app.db' not in sys.modules; "
"assert 'sqlalchemy' not in sys.modules",
],
cwd=Path(__file__).resolve().parents[1],
env=env,
capture_output=True,
text=True,
check=False,
)
assert result.returncode == 0, result.stderr
def test_schema_has_database_guards_and_append_only_audit():
sql = Path("migrations/007_document_reconciliation_v2.sql").read_text()
assert "ux_document_link_primary_kind" in sql and "WHERE ended_at IS NULL AND relationship = 'PRIMARY'" in sql
assert "ux_document_link_current" in sql
assert "BEFORE UPDATE OR DELETE" in sql
assert "ux_document_link_events_idempotency" in sql
def test_active_router_owns_v2_handlers_and_htmx_panel_target():
router = Path("app/admin_ui/router.py").read_text()
page = Path("app/admin_ui/pages/opportunities.py").read_text()
assert "router.include_router(opportunities.router)" in router
assert "document_reconciliation_v2_action" in page
assert "HTMLResponse(jasmin_documents_html(opportunity_id" in page
assert "_authenticated_actor(request)" in page
def test_sync_does_not_assign_legacy_primary_by_import_order():
source = Path("app/reconciliation_service.py").read_text()
function = source[source.index("def _upsert_jasmin_document_from_item"):source.index("def _", source.index("def _upsert_jasmin_document_from_item") + 10)]
assert '"role": "related"' in function
assert '"is_primary": False' in function
assert "classify_sync_document" in function
canonical = Path("app/document_reconciliation_service.py").read_text()
assert "multiple Jasmin candidates require review" in canonical
assert "SYNC_AMBIGUITY_DEMOTED" in canonical
def test_total_never_falls_back_to_non_primary_or_cancelled_primary():
rows = [{"id": "secondary", "relationship": "SECONDARY", "total_amount": 999}]
assert select_valid_primary(rows) is None
rows.insert(0, {"id": "cancelled", "relationship": "PRIMARY", "status": "CANCELLED"})
assert select_valid_primary(rows) is None
def test_only_unambiguous_keeps_ambiguous_group_fully_represented():
actions = classify_group([_doc("a", role="current", primary=True),
_doc("b", role="accepted", primary=True)])
selected = actions_for_only_unambiguous(actions)
assert len(selected) == 2
assert {row["relationship"] for row in selected} == {"REVIEW_REQUIRED"}
def test_backfill_paginates_complete_groups_with_composite_cursor():
source = Path("scripts/migrate_document_reconciliation_v2.py").read_text()
assert "WITH selected_groups AS" in source
assert "GROUP BY d.opportunity_id, d.document_kind" in source
assert "LIMIT :limit" in source
assert "OPPORTUNITY_ID|DOCUMENT_KIND" in source
assert "JOIN commercial_documents d ON d.opportunity_id=g.opportunity_id" in source
def test_group_larger_than_batch_size_is_never_split():
rows = [{"id": str(index), "opportunity_id": "a", "document_kind": "quotation"}
for index in range(5)] + [{"id": "z", "opportunity_id": "b", "document_kind": "quotation"}]
batch = select_complete_group_batch(rows, batch_size=1)
assert len(batch) == 5 and {row["opportunity_id"] for row in batch} == {"a"}
def test_composite_resume_keeps_other_kind_on_opportunity_boundary():
rows = [
{"id": "1", "opportunity_id": "a", "document_kind": "invoice"},
{"id": "2", "opportunity_id": "a", "document_kind": "quotation"},
{"id": "3", "opportunity_id": "b", "document_kind": "quotation"},
]
first = select_complete_group_batch(rows, batch_size=1)
cursor = (first[-1]["opportunity_id"], first[-1]["document_kind"])
second = select_complete_group_batch(rows, batch_size=1, resume_from=cursor)
assert [(row["opportunity_id"], row["document_kind"]) for row in second] == [("a", "quotation")]
def test_idempotency_fingerprint_retry_and_conflicting_reuse_contract():
first, payload = _command_fingerprint(opportunity_id="o", document_id="d",
target_relationship="PRIMARY", reason="approved")
retry, _ = _command_fingerprint(opportunity_id="o", document_id="d",
target_relationship="PRIMARY", reason="approved")
conflict, _ = _command_fingerprint(opportunity_id="o", document_id="d",
target_relationship="SECONDARY", reason="approved")
assert first == retry and first != conflict
source = Path("app/document_reconciliation_service.py").read_text()
function = source[source.index("def set_document_relationship"):source.index("def set_primary_document")]
assert function.index("_claim_idempotency") < function.index("_lock_context")
assert "already used with a different command" in source
def test_idempotency_fingerprint_covers_every_material_command_field():
base = dict(opportunity_id="o", document_id="d", target_relationship="PRIMARY",
reason="approved", source="api", event_type="CHANGED", metadata={"a": 1},
is_manual=True, actor="alice", document_kind="invoice")
fingerprint, _ = _command_fingerprint(**base)
for field, value in (("source", "sync"), ("event_type", "BACKFILLED"),
("metadata", {"a": 2}), ("is_manual", False), ("actor", "bob"),
("correlation_id", "corr-1"), ("request_id", "req-1")):
changed, _ = _command_fingerprint(**(base | {field: value}))
assert changed != fingerprint, field
reordered, _ = _command_fingerprint(**(base | {"metadata": {"z": 2, "a": 1}}))
same_reordered, _ = _command_fingerprint(**(base | {"metadata": {"a": 1, "z": 2}}))
assert reordered == same_reordered
def test_generic_document_writer_registers_quotation_and_invoice_through_canonical_service():
source = Path("app/commercial_service.py").read_text()
function = source[source.index("def create_commercial_document"):source.index("def add_document_lines")]
assert "register_created_document(" in function
assert function.index("INSERT INTO commercial_documents") < function.index("register_created_document(")
assert "and not v2_available" in function
assert "UPDATE opportunity_document_links" not in function
assert "INSERT INTO opportunity_document_links" not in function
canonical = Path("app/document_reconciliation_service.py").read_text()
register = canonical[canonical.index("def register_created_document"):canonical.index("def record_jasmin_status_change")]
assert "document_reconciliation_v2_available(conn)" in register
assert "classify_sync_document(" in register
def test_new_group_member_preserves_manual_primary_and_uses_secondary_or_review():
source = Path("app/document_reconciliation_service.py").read_text()
classify = source[source.index("def classify_sync_document"):source.index("def register_created_document")]
assert 'relationship = "SECONDARY" if manual_primary' in classify
assert 'else "REVIEW_REQUIRED"' in classify
assert "multiple Jasmin candidates require review" in classify
assert 'and not previous["is_manual"]' in classify
def test_all_commercial_document_producers_register_v2_relationships():
commercial = Path("app/commercial_service.py").read_text()
reconciliation = Path("app/reconciliation_service.py").read_text()
assert commercial.count("INSERT INTO commercial_documents") == 1
assert "register_created_document(" in commercial
assert reconciliation.count("INSERT INTO commercial_documents") == 1
assert "classify_sync_document(" in reconciliation
def test_backfill_applies_one_transaction_per_complete_group():
source = Path("scripts/migrate_document_reconciliation_v2.py").read_text()
function = source[source.index("def apply_actions"):source.index("def write_reports")]
assert "groups.setdefault" in function
assert function.count("with engine.begin() as conn") == 1
assert "_conn=conn" in function
def test_backfill_second_document_failure_rolls_back_whole_group(monkeypatch):
import app.db
import app.document_reconciliation_service as service
committed = {"links": [], "events": [], "legacy": []}
class Transaction:
def __enter__(self):
self.pending = {key: [] for key in committed}
return self
def __exit__(self, exc_type, exc, tb):
if exc_type is None:
for key in committed:
committed[key].extend(self.pending[key])
return False
class FakeEngine:
def begin(self): return Transaction()
calls = 0
def fake_set(*args, _conn=None, **kwargs):
nonlocal calls
calls += 1
_conn.pending["links"].append(args[1])
_conn.pending["events"].append(args[1])
_conn.pending["legacy"].append(args[1])
if calls == 2:
raise RuntimeError("simulated second document failure")
monkeypatch.setattr(app.db, "engine", FakeEngine())
monkeypatch.setattr(service, "set_document_relationship", fake_set)
actions = [dict(_doc("a"), relationship="PRIMARY", decision_reason="one"),
dict(_doc("b"), relationship="SECONDARY", decision_reason="two")]
with pytest.raises(RuntimeError, match="second document"):
apply_actions(actions)
assert committed == {"links": [], "events": [], "legacy": []}
def test_generic_line_writer_is_schema_compatible_and_keeps_ids_coherent():
source = Path("app/commercial_service.py").read_text()
function = source[source.index("def add_document_lines"):source.index("def list_commercial_documents")]
assert "to_regclass('commercial_document_lines')" in function
assert "FROM pg_attribute" in function
assert "attname = 'commercial_document_id'" in function
assert "table_schema='public'" not in function
assert 'CAST(:document_id AS UUID), CAST(:document_id AS UUID)' in function
def test_startup_does_not_partially_apply_migration_007():
source = Path("app/commercial_service.py").read_text()
function = source[source.index("def ensure_commercial_schema"):source.index("def get_customer_by_tax_id")]
assert "ADD COLUMN IF NOT EXISTS commercial_document_id" not in function
assert "ADD COLUMN IF NOT EXISTS opportunity_document_link_id" not in function
def test_no_new_authoritative_v2_reader_outside_canonical_service():
allowlist = {"document_reconciliation_service.py"}
offenders = []
for path in Path("app").rglob("*.py"):
if path.name in allowlist:
continue
if "opportunity_document_links" in path.read_text(encoding="utf-8"):
offenders.append(str(path))
# Remaining consumers are deliberately tracked until migrated; additions
# outside this explicit inventory fail the review gate.
assert offenders == []
def test_membership_and_manual_reason_are_domain_guards():
source = Path("app/document_reconciliation_service.py").read_text()
assert "document does not belong to this opportunity" in source
assert "reason is required for every manual relationship change" in source
assert "commercial_documents WHERE id=CAST(:did AS UUID)" in source
def test_preflight_and_down_migration_cover_blockers_and_dependency_order():
assert callable(run_preflight)
preflight = Path("scripts/preflight_document_reconciliation_v2.py").read_text()
assert "null_document_id" in preflight and "orphan_document_id" in preflight
assert "multiple_legacy_primaries" in preflight and "duplicate_document_groups" in preflight
down = Path("migrations/007_document_reconciliation_v2_down.sql").read_text()
assert down.index("DROP COLUMN IF EXISTS opportunity_document_link_id") < down.index("DROP TABLE IF EXISTS opportunity_document_links")
assert down.index("DROP TRIGGER") < down.index("DROP FUNCTION") < down.index("DROP TABLE IF EXISTS opportunity_document_link_events")
assert "document_ledger_exported" in down and "v2_consumers_active" in down
def test_jasmin_items_keep_document_and_link_provenance():
source = Path("app/reconciliation_service.py").read_text()
assert '"commercial_document_id": origin.get("commercial_document_id")' in source
assert '"opportunity_document_link_id": origin.get("opportunity_document_link_id")' in source