From 29382ec6c1881127dfcb3c759a3a8e58c279a60a Mon Sep 17 00:00:00 2001 From: plx Date: Wed, 5 Aug 2026 23:55:33 +0000 Subject: [PATCH] Implement document reconciliation v2 --- app/admin_auth.py | 2 + app/admin_dashboard.py | 29 +- app/admin_ui/pages/opportunities.py | 125 +++- app/commercial_service.py | 183 ++--- app/company_opportunity_linking.py | 11 +- app/document_reconciliation_backfill.py | 119 +++ app/document_reconciliation_service.py | 684 ++++++++++++++++++ app/domain/opportunity_flow/evidence.py | 10 +- app/followup_service.py | 13 +- app/jasmin_fiscal_sync_service.py | 37 +- app/operation_service.py | 18 +- app/opportunity_next_action_service.py | 21 +- app/reconciliation_service.py | 104 +-- app/reply_assistant_service.py | 3 +- app/revenue_forecast_service.py | 67 +- app/task_service.py | 48 +- app/workflow_guard.py | 46 +- docs/document_reconciliation_v2.md | 39 + migrations/007_document_reconciliation_v2.sql | 100 +++ .../007_document_reconciliation_v2_down.sql | 37 + scripts/apply_migrations.py | 11 +- scripts/migrate_document_reconciliation_v2.py | 137 ++++ .../preflight_document_reconciliation_v2.py | 107 +++ tests/test_document_reconciliation_v2.py | 370 ++++++++++ ...est_document_reconciliation_v2_postgres.py | 104 +++ 25 files changed, 2082 insertions(+), 343 deletions(-) create mode 100644 app/document_reconciliation_backfill.py create mode 100644 app/document_reconciliation_service.py create mode 100644 docs/document_reconciliation_v2.md create mode 100644 migrations/007_document_reconciliation_v2.sql create mode 100644 migrations/007_document_reconciliation_v2_down.sql create mode 100644 scripts/migrate_document_reconciliation_v2.py create mode 100644 scripts/preflight_document_reconciliation_v2.py create mode 100644 tests/test_document_reconciliation_v2.py create mode 100644 tests/test_document_reconciliation_v2_postgres.py diff --git a/app/admin_auth.py b/app/admin_auth.py index b5c181d..5312b3f 100644 --- a/app/admin_auth.py +++ b/app/admin_auth.py @@ -47,11 +47,13 @@ def require_admin_auth(request: Request, *, area: str) -> None: ).strip() if not received or not compare_digest(received, expected): raise HTTPException(status_code=401, detail=f"{detail_prefix} auth required") + request.state.clientflow_admin_user = "token-admin" return if mode == "local": if not _is_loopback_request(request): raise HTTPException(status_code=401, detail=f"{detail_prefix} local access required") + request.state.clientflow_admin_user = "local-admin" return # Settings validates the value, but fail closed if it is mutated at runtime. diff --git a/app/admin_dashboard.py b/app/admin_dashboard.py index cd40e13..904c060 100644 --- a/app/admin_dashboard.py +++ b/app/admin_dashboard.py @@ -665,8 +665,8 @@ def operation_cockpit_html(opportunity_id: str, opportunity: dict, snapshot: dic and str(doc.get("status") or "").lower() not in {"cancelled", "failed", "rejected"} and str(doc.get("role") or "current") not in {"superseded"} ] - invoice_document = next((doc for doc in invoice_candidates if bool(doc.get("is_primary", False))), invoice_candidates[0] if invoice_candidates else None) - quotation_document = next((doc for doc in quotation_candidates if bool(doc.get("is_primary", False))), quotation_candidates[0] if quotation_candidates else None) + invoice_document = next((doc for doc in invoice_candidates if doc.get("relationship") == "PRIMARY"), None) + quotation_document = next((doc for doc in quotation_candidates if doc.get("relationship") == "PRIMARY"), None) except Exception: invoice_document = None quotation_document = None @@ -1335,23 +1335,30 @@ def jasmin_documents_html(opportunity_id: str, *, notice: str = "", error_notice number = commercial_document_display_number(doc, fallback="número por atualizar") amount = doc.get("total_amount") if doc.get("total_amount") is not None else doc.get("amount") kind = kind_labels.get(str(doc.get("document_kind") or ""), doc.get("document_kind") or "Documento") - doc_id = str(doc.get("id") or "") - role = str(doc.get("role") or "current") - is_primary = bool(doc.get("is_primary")) - make_primary_action = "" if (role in {"current", "accepted"} and is_primary) else f'
' - historical_action = "" if role == "historical" else f'
' + # Every document action targets commercial_documents.id. link_id is + # audit provenance only and must never leak into document routes. + doc_id = str(doc.get("document_id") or doc.get("id") or "") + relationship = str(doc.get("relationship") or "").upper() + role = str(doc.get("role") or relationship.lower() or "review_required") + is_primary = relationship == "PRIMARY" + make_primary_action = "" if relationship == "PRIMARY" else f'
' + historical_action = "" if relationship == "HISTORICAL" else f'
' + secondary_action = "" if relationship == "SECONDARY" else f'
' unlink_action = f'
' refresh_action = f'
' pdf_action = f'PDF' - actions = f'
{make_primary_action}{historical_action}{unlink_action}{refresh_action}{pdf_action}
' + history_action = f'Ver histórico' + actions = f'
{make_primary_action}{secondary_action}{historical_action}{unlink_action}{history_action}{refresh_action}{pdf_action}
' role_label = { "current": "Atual", "accepted": "Aceite", "historical": "Histórico", "cancelled": "Cancelado", "related": "Relacionado", - }.get(str(doc.get("role") or "current"), str(doc.get("role") or "current")) - role_class = "text-bg-primary" if str(doc.get("role") or "current") in {"current", "accepted"} and doc.get("is_primary") else "text-bg-light" + "primary": "PRIMARY", "secondary": "SECONDARY", "review_required": "REVIEW_REQUIRED", + "ignored": "IGNORED", "removed": "REMOVED", + }.get(role, relationship or role) + role_class = "text-bg-danger" if relationship == "REVIEW_REQUIRED" else ("text-bg-primary" if relationship == "PRIMARY" else "text-bg-light") external_id_html = "" if not doc.get("document_number") and doc.get("external_id") and not is_uuid_text(doc.get("external_id")): external_id_html = f"
ref. externa {esc(doc.get('external_id') or '')}
" @@ -1361,7 +1368,7 @@ def jasmin_documents_html(opportunity_id: str, *, notice: str = "", error_notice "" f"{esc(kind)}
v{esc(doc.get('version_number') or '—')}
{esc(role_label)}" f"{esc(number)}{external_id_html}" - f"{operation_status_badge(str(doc.get('status') or 'created'))}" + f"{operation_status_badge(str(doc.get('status') or 'created'))}
Relação: {esc(relationship or role_label)}
{esc(doc.get('decision_reason') or '—')} · {esc(doc.get('decided_by') or '—')}
" f"{money_html(amount or 0)}" f"{esc(fmt_dt(doc.get('created_at')))}" f"{actions}" diff --git a/app/admin_ui/pages/opportunities.py b/app/admin_ui/pages/opportunities.py index a251b91..60ff0de 100644 --- a/app/admin_ui/pages/opportunities.py +++ b/app/admin_ui/pages/opportunities.py @@ -47,6 +47,13 @@ _opportunity_board_column_for_stage = legacy._opportunity_board_column_for_stage router = APIRouter() + +def _authenticated_actor(request: Request) -> str: + actor = str(getattr(request.state, "clientflow_admin_user", "") or "").strip() + if not actor: + raise ValueError("Identidade administrativa autenticada em falta.") + return actor + PAYMENT_TERM_LABELS = { "before_shipping": "Antes do envio", "after_delivery": "Após entrega", @@ -223,16 +230,21 @@ def _opportunity_invoice_sent_evidence(opportunity_id: str, invoice_number: str def _document_display_number(doc: dict | None) -> str: if not doc: return "—" - return str(doc.get("document_number") or doc.get("external_id") or doc.get("id") or "documento") + return str(doc.get("document_number") or doc.get("external_id") + or doc.get("document_id") or doc.get("id") or "documento") def _finance_quick_card_html(opportunity_id: str, linked_documents: list[dict], payment_term: str, payment_term_label: str) -> str: # BLIF default flow: quotation -> payment -> invoice -> prepare/ship. - quotation = next((d for d in linked_documents if str(d.get("document_kind") or "") in {"quotation", "proforma"} and str(d.get("role") or "current") in {"current", "accepted"}), None) - invoice = next((d for d in linked_documents if str(d.get("document_kind") or "") == "invoice" and str(d.get("role") or "current") in {"current", "accepted"}), None) + quotation = next((d for d in linked_documents if str(d.get("document_kind") or "") in {"quotation", "proforma"} + and str(d.get("relationship") or "").upper() == "PRIMARY" + and str(d.get("jasmin_status") or d.get("status") or "").upper() not in {"CANCELLED", "CANCELED"}), None) + invoice = next((d for d in linked_documents if str(d.get("document_kind") or "") == "invoice" + and str(d.get("relationship") or "").upper() == "PRIMARY" + and str(d.get("jasmin_status") or d.get("status") or "").upper() not in {"CANCELLED", "CANCELED"}), None) payment_confirmed = _opportunity_payment_confirmed(opportunity_id) base_doc = invoice or quotation - base_doc_label = "Fatura" if invoice else ("Orçamento" if quotation else "Documento") + base_doc_label = "Fatura" if invoice else ("Orçamento" if quotation else "Documento principal por definir") amount = (base_doc.get("total_amount") or base_doc.get("amount")) if base_doc else None amount_html = money_html(float(amount or 0)) if amount else "—" payment_status = "Confirmado" if payment_confirmed else ("Pendente pós-entrega" if payment_term == "after_delivery" else "Por confirmar") @@ -1813,26 +1825,28 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N ( doc for doc in linked_documents if str(doc.get("document_kind") or "") == "invoice" - and str(doc.get("role") or "current") in {"current", "accepted"} - and bool(doc.get("is_primary", True)) + and str(doc.get("relationship") or "").upper() == "PRIMARY" ), next( ( doc for doc in linked_documents - if str(doc.get("role") or "current") in {"current", "accepted"} - and bool(doc.get("is_primary", True)) + if str(doc.get("relationship") or "").upper() == "PRIMARY" ), - linked_documents[0] if linked_documents else None, + None, ), ) opportunity_items_total = sum( float(item.get("total_price") or 0) for item in opportunity_items if str(item.get("status") or "").upper() not in {"REJECTED", "CANCELLED", "DELIVERED", "HISTORICAL"} + and str((item.get("metadata") or {}).get("source_system") or "manual").lower() in {"manual", "odoo", "operational"} + and bool((item.get("metadata") or {}).get("approved", True)) ) document_value = float(primary_document.get("total_amount") or primary_document.get("amount") or 0) if primary_document else 0 estimated_value = document_value or opportunity_items_total or float(opportunity.get("value_amount") or 0) - value_source = "documento principal" if document_value else ("linhas atuais" if opportunity_items_total else "oportunidade") + value_source = "documento principal" if document_value else ( + "valor manual/operacional aprovado" if opportunity_items_total else "documento principal por definir" + ) operation_snapshot = get_operation_snapshot(opportunity_id) opportunity_for_cockpit = dict(opportunity) opportunity_for_cockpit["pending_task_count"] = len(pending_tasks) @@ -2614,26 +2628,24 @@ async def opportunity_products_partial(opportunity_id: str): @router.post("/commercial-documents/{document_id}/unlink-from-opportunity") async def commercial_document_unlink_from_opportunity_action(document_id: str, request: Request): form = await request.form() - opportunity_id = str(form.get("opportunity_id") or "").strip() - remove_lines = str(form.get("remove_imported_lines") or "1") == "1" - note = str(form.get("note") or "").strip() or "Documento removido manualmente desta oportunidade; pertence a outra compra/processo." + from app.document_reconciliation_service import resolve_document_opportunity + opportunity_id = str(resolve_document_opportunity(document_id) or "") + note = str(form.get("reason") or form.get("note") or "").strip() if not is_uuid_text(opportunity_id) or not is_uuid_text(document_id): return PlainTextResponse("Identificador inválido.", status_code=422) try: - result = unlink_commercial_document_from_opportunity( - opportunity_id, - document_id, - remove_imported_lines=remove_lines, - note=note, - actor="operator_ui_document_unlink", - ) + if not note: + raise ValueError("O motivo é obrigatório.") + from app.document_reconciliation_service import remove_document + remove_document(opportunity_id, document_id, actor=_authenticated_actor(request), reason=note, + request_id=request.headers.get("x-request-id"), correlation_id=request.headers.get("x-correlation-id")) + result = {"document_unlinked": 1, "imported_lines_deleted": 0} except Exception as exc: if request.headers.get("hx-request"): return HTMLResponse(jasmin_documents_html(opportunity_id, error_notice=f"Erro ao desassociar documento: {exc}"), status_code=409) return PlainTextResponse(f"Erro ao desassociar documento: {exc}", status_code=500) notice = ( - "Documento desassociado desta oportunidade. " - f"Linhas importadas removidas: {result.get('imported_lines_deleted', 0)}." + "Relação removida desta oportunidade; documento, linhas e histórico foram preservados." ) if request.headers.get("hx-request"): return HTMLResponse(jasmin_documents_html(opportunity_id, notice=notice)) @@ -2643,30 +2655,78 @@ async def commercial_document_unlink_from_opportunity_action(document_id: str, r @router.post("/commercial-documents/{document_id}/role") async def commercial_document_role_action(document_id: str, request: Request): form = await request.form() - opportunity_id = str(form.get("opportunity_id") or "").strip() - role = str(form.get("role") or "current").strip().lower() + from app.document_reconciliation_service import resolve_document_opportunity + opportunity_id = str(resolve_document_opportunity(document_id) or "") + role = str(form.get("relationship") or form.get("role") or "PRIMARY").strip().upper() make_primary = str(form.get("make_primary") or "1") == "1" + reason = str(form.get("reason") or "").strip() if not is_uuid_text(opportunity_id) or not is_uuid_text(document_id): return PlainTextResponse("Identificador inválido.", status_code=422) try: - result = set_commercial_document_role_for_opportunity( - opportunity_id, - document_id, - role=role, - make_primary=make_primary, - actor="operator_ui_document_role", - ) + from app.document_reconciliation_service import set_document_relationship + legacy_map = {"CURRENT": "PRIMARY", "ACCEPTED": "SECONDARY", "RELATED": "SECONDARY", "HISTORICAL": "HISTORICAL"} + relationship = legacy_map.get(role, role) + if make_primary: relationship = "PRIMARY" + if relationship == "PRIMARY" and not reason: + raise ValueError("O motivo é obrigatório para substituir/confirmar o principal.") + result = set_document_relationship(opportunity_id, document_id, relationship, + actor=_authenticated_actor(request), reason=reason or f"Marcado {relationship} na UI", + request_id=request.headers.get("x-request-id"), correlation_id=request.headers.get("x-correlation-id")) except Exception as exc: if request.headers.get("hx-request"): return HTMLResponse(jasmin_documents_html(opportunity_id, error_notice=f"Erro ao atualizar papel do documento: {exc}"), status_code=409) return PlainTextResponse(f"Erro ao atualizar papel do documento: {exc}", status_code=500) - label = {"current": "atual", "accepted": "aceite", "related": "relacionado", "historical": "histórico"}.get(result.get("role"), role) + label = str(result.get("relationship") or role).lower() notice = f"Documento marcado como {label}." if request.headers.get("hx-request"): return HTMLResponse(jasmin_documents_html(opportunity_id, notice=notice)) return RedirectResponse(f"/opportunities/{opportunity_id}?notice={quote(notice)}", status_code=303) +@router.post("/opportunities/{opportunity_id}/documents/{document_id}/relationship") +async def document_reconciliation_v2_action(opportunity_id: str, document_id: str, request: Request): + """Canonical v2 HTTP transition; never mutates the remote Jasmin document.""" + if not is_uuid_text(opportunity_id) or not is_uuid_text(document_id): + return PlainTextResponse("Identificador inválido.", status_code=422) + form = await request.form() + relationship = str(form.get("relationship") or "").upper().strip() + reason = str(form.get("reason") or "").strip() + destination = str(form.get("destination_opportunity_id") or "").strip() + try: + from app.document_reconciliation_service import ( + reassign_document, restore_document, set_document_relationship, + ) + common = {"actor": _authenticated_actor(request), "reason": reason, + "request_id": request.headers.get("x-request-id"), + "correlation_id": request.headers.get("x-correlation-id")} + if relationship == "REASSIGNED": + if not is_uuid_text(destination): raise ValueError("O destino é obrigatório e deve ser válido.") + reassign_document(opportunity_id, destination, document_id, **common) + elif relationship == "RESTORE": + restore_document(opportunity_id, document_id, **common) + else: + set_document_relationship(opportunity_id, document_id, relationship, **common) + except Exception as exc: + if request.headers.get("hx-request"): + return HTMLResponse(jasmin_documents_html(opportunity_id, error_notice=f"Reconciliação não aplicada: {exc}"), status_code=409) + return PlainTextResponse(f"Reconciliação não aplicada: {exc}", status_code=409) + if request.headers.get("hx-request"): + return HTMLResponse(jasmin_documents_html(opportunity_id, notice="Relação ClientFlow atualizada; o Jasmin não foi alterado.")) + return RedirectResponse(f"/opportunities/{opportunity_id}#documentos", status_code=303) + + +@router.get("/opportunities/{opportunity_id}/documents/{document_id}/history") +async def document_reconciliation_v2_history(opportunity_id: str, document_id: str): + if not is_uuid_text(opportunity_id) or not is_uuid_text(document_id): + return PlainTextResponse("Identificador inválido.", status_code=422) + from app.document_reconciliation_service import get_document_link, list_document_link_events + if not get_document_link(opportunity_id, document_id, include_ended=True): + return PlainTextResponse("Documento não pertence ao histórico desta oportunidade.", status_code=404) + events = list_document_link_events(opportunity_id, document_id=document_id) + rows = "".join(f"{esc(fmt_dt(e.get('created_at')))}{esc(e.get('event_type'))}{esc(e.get('actor'))}{esc(e.get('old_relationship') or '—')} → {esc(e.get('new_relationship') or '—')}{esc(e.get('reason') or '—')}" for e in events) + return HTMLResponse(f'

Histórico de reconciliação

{rows}
DataEventoOperadorRelaçãoMotivo
') + + @router.post("/commercial-documents/{document_id}/refresh") async def commercial_document_refresh(document_id: str, request: Request): form = await request.form() @@ -3425,4 +3485,3 @@ async def opportunity_odoo_link_candidate_action(opportunity_id: str, item_id: s return HTMLResponse(odoo_status_panel_html(opportunity_id, error_notice=f"Erro ao associar venda Odoo: {exc}"), status_code=409) return PlainTextResponse(f"Erro ao associar venda Odoo: {exc}", status_code=500) return RedirectResponse(url=f"/opportunities/{opportunity_id}?notice=Venda%20Odoo%20associada", status_code=303) - diff --git a/app/commercial_service.py b/app/commercial_service.py index 0df1c27..b4592ed 100644 --- a/app/commercial_service.py +++ b/app/commercial_service.py @@ -158,104 +158,12 @@ def ensure_commercial_schema() -> None: conn.execute(text("CREATE INDEX IF NOT EXISTS idx_commercial_documents_kind ON commercial_documents(document_kind, status)")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_commercial_documents_role ON commercial_documents(opportunity_id, document_kind, role, is_primary)")) - # Hotfix v4927.1: em bases já usadas podem existir vários documentos - # Jasmin/Odoo ligados à mesma oportunidade e ao mesmo tipo. Ao adicionar - # role/is_primary com DEFAULT current/TRUE, todos passariam a ser - # “atuais”, fazendo o índice único falhar no arranque e deixando o - # serviço indisponível atrás do nginx. Antes de criar o índice, mantemos - # só o documento mais recente como primário e movemos os restantes para - # histórico. - conn.execute(text(""" - WITH ranked AS ( - SELECT - id, - ROW_NUMBER() OVER ( - PARTITION BY opportunity_id, system, document_kind, COALESCE(role, 'current') - ORDER BY COALESCE(updated_at, created_at) DESC, created_at DESC, id DESC - ) AS rn - FROM commercial_documents - WHERE opportunity_id IS NOT NULL - AND COALESCE(is_primary, TRUE) = TRUE - AND COALESCE(role, 'current') IN ('current', 'accepted') - ) - UPDATE commercial_documents AS d - SET role = CASE - WHEN COALESCE(d.role, 'current') IN ('current', 'accepted') THEN 'historical' - ELSE COALESCE(d.role, 'historical') - END, - is_primary = FALSE, - is_active = FALSE, - updated_at = now() - FROM ranked AS r - WHERE d.id = r.id - AND r.rn > 1 - """)) - # v4927.2: não criar índice UNIQUE durante o arranque. - # Em produção podem existir duplicados históricos que ainda não foram - # reconciliados; se o índice único falhar, toda a app fica indisponível - # e o nginx devolve 502. A regra de “um principal por fase” fica - # aplicada pela normalização acima e pelos serviços que promovem um - # documento a atual. O índice de apoio é não único e seguro em bases - # reais com dados legados. + # v2: startup must never classify or rewrite document relationships. + # The versioned migration creates the authoritative partial unique + # constraint; legacy fields below are dual-write compatibility only. conn.execute(text("DROP INDEX IF EXISTS ux_commercial_documents_primary_role")) conn.execute(text("CREATE INDEX IF NOT EXISTS idx_commercial_documents_primary_role ON commercial_documents(opportunity_id, system, document_kind, role) WHERE opportunity_id IS NOT NULL AND COALESCE(is_primary, TRUE) = TRUE AND COALESCE(role, 'current') IN ('current', 'accepted')")) - # v4928.1: migração leve de registos antigos/reconstruídos. - # Não apaga documentos nem altera fases. Apenas evita que orçamentos e - # pró-formas legados continuem a aparecer como documento principal quando - # já existe fatura Jasmin atual/aceite para a mesma oportunidade. - conn.execute(text(""" - WITH primary_invoice AS ( - SELECT DISTINCT ON (opportunity_id, system) - id, opportunity_id, system, COALESCE(document_date, created_at::date) AS invoice_date, created_at - FROM commercial_documents - WHERE opportunity_id IS NOT NULL - AND system = 'jasmin' - AND document_kind = 'invoice' - AND COALESCE(is_primary, TRUE) = TRUE - AND COALESCE(role, 'current') IN ('current', 'accepted') - ORDER BY opportunity_id, system, COALESCE(document_date, created_at::date) DESC, created_at DESC - ) - UPDATE commercial_documents AS d - SET role = 'historical', - is_primary = FALSE, - is_active = FALSE, - payload = COALESCE(d.payload, '{}'::jsonb) || jsonb_build_object('legacy_reason', 'superseded_by_invoice_v4928_1'), - updated_at = now() - FROM primary_invoice AS inv - WHERE d.opportunity_id = inv.opportunity_id - AND d.system = inv.system - AND d.document_kind IN ('quotation', 'proforma') - AND COALESCE(d.role, 'current') IN ('current', 'accepted') - AND COALESCE(d.is_primary, TRUE) = TRUE - AND COALESCE(d.document_date, d.created_at::date) <= inv.invoice_date - """)) - - # Marca oportunidades reconstruídas com fatura e sem task pendente para a UI - # poder explicar que são processos antigos/auditoria, não um funil limpo. - conn.execute(text(""" - UPDATE opportunities AS o - SET metadata = COALESCE(o.metadata, '{}'::jsonb) || jsonb_build_object( - 'clientflow_record_mode', COALESCE(o.metadata->>'clientflow_record_mode', 'reconstructed_invoice_review'), - 'clientflow_legacy_migrated_by', 'v4928_1' - ), - updated_at = now() - WHERE EXISTS ( - SELECT 1 FROM commercial_documents d - WHERE d.opportunity_id = o.id - AND d.system = 'jasmin' - AND d.document_kind = 'invoice' - AND COALESCE(d.role, 'current') IN ('current', 'accepted') - AND COALESCE(d.is_primary, TRUE) = TRUE - ) - AND NOT EXISTS ( - SELECT 1 FROM tasks t - WHERE t.opportunity_id = o.id - AND t.status = 'pending' - ) - AND COALESCE(o.metadata->>'clientflow_record_mode', '') = '' - """)) - conn.execute(text(""" CREATE TABLE IF NOT EXISTS commercial_document_lines ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), @@ -274,6 +182,9 @@ def ensure_commercial_schema() -> None: created_at TIMESTAMPTZ NOT NULL DEFAULT now() ) """)) + # Reconciliation v2 columns belong exclusively to migration 007. + # Startup must leave a pre-007 database unchanged so the migration can + # add, backfill and constrain the canonical identifier atomically. conn.execute(text("CREATE INDEX IF NOT EXISTS idx_commercial_document_lines_document ON commercial_document_lines(document_id)")) conn.execute(text(""" @@ -500,8 +411,13 @@ def create_commercial_document( if not document_kind: raise ValueError("document_kind é obrigatório") with engine.begin() as conn: + from app.document_reconciliation_service import ( + document_reconciliation_v2_available, + register_created_document, + ) + v2_available = document_reconciliation_v2_available(conn) version_number = _next_document_version(conn, opportunity_id, document_kind) - if document_kind == "quotation" and opportunity_id: + if document_kind == "quotation" and opportunity_id and not v2_available: conn.execute(text(""" UPDATE commercial_documents SET status = CASE WHEN status IN ('created','sent','draft') THEN 'superseded' ELSE status END, @@ -553,19 +469,46 @@ def create_commercial_document( "due_date": due_date, "payload": _json(payload), }).mappings().first() + if row and opportunity_id and v2_available: + link = register_created_document( + conn, opportunity_id, str(row["id"]), document_kind, + status=status, actor="commercial_service.create_commercial_document", + ) + if link: + row = dict(row) | { + "document_id": str(row["id"]), + "link_id": str(link["id"]), + "relationship": link["relationship"], + } return dict(row or {}) def add_document_lines(document_id: str, lines: Iterable[Dict[str, Any]]) -> None: ensure_commercial_schema() with engine.begin() as conn: + # Migration 007 is additive. Select the statement once per batch so + # the generic writer works both before and after the new NOT NULL + # canonical column is introduced. + has_canonical_id = bool(conn.execute(text(""" + SELECT EXISTS ( + SELECT 1 + FROM pg_attribute + WHERE attrelid = to_regclass('commercial_document_lines') + AND attname = 'commercial_document_id' + AND attnum > 0 + AND NOT attisdropped + ) + """)).scalar()) for idx, line in enumerate(lines): - conn.execute(text(""" + columns = ("document_id, commercial_document_id, " if has_canonical_id else "document_id, ") + values = ("CAST(:document_id AS UUID), CAST(:document_id AS UUID), " if has_canonical_id + else "CAST(:document_id AS UUID), ") + conn.execute(text(f""" INSERT INTO commercial_document_lines ( - document_id, opportunity_item_id, line_index, local_product_id, jasmin_sales_item, + {columns}opportunity_item_id, line_index, local_product_id, jasmin_sales_item, description, quantity, unit, unit_price, tax_schema, total_amount, payload ) VALUES ( - CAST(:document_id AS UUID), CAST(:opportunity_item_id AS UUID), :line_index, + {values}CAST(:opportunity_item_id AS UUID), :line_index, CAST(:local_product_id AS UUID), :jasmin_sales_item, :description, :quantity, :unit, :unit_price, :tax_schema, :total_amount, CAST(:payload AS JSONB) ) @@ -587,6 +530,15 @@ def add_document_lines(document_id: str, lines: Iterable[Dict[str, Any]]) -> Non def list_commercial_documents(opportunity_id: Optional[str] = None, customer_id: Optional[str] = None, limit: int = 100) -> List[Dict[str, Any]]: ensure_commercial_schema() + if opportunity_id: + from app.document_reconciliation_service import resolve_document_links + # Canonical resolver owns the only group-complete rollout fallback. + rows = resolve_document_links(opportunity_id) + if customer_id: + rows = [row for row in rows if str(row.get("customer_id") or "") == str(customer_id)] + return rows[:int(limit)] + # Legitimate legacy inventory: without an opportunity scope this endpoint + # lists document records for customer/history, not effective relationships. filters = [] params: Dict[str, Any] = {"limit": int(limit)} if opportunity_id: @@ -602,15 +554,17 @@ def list_commercial_documents(opportunity_id: Optional[str] = None, customer_id: cd.document_kind, cd.external_id, cd.external_url, cd.company, cd.document_type, cd.serie, cd.series_number, cd.document_number, cd.customer_party_key, cd.status, cd.amount, cd.tax_amount, cd.total_amount, cd.currency, - cd.parent_document_id::text, cd.version_number, cd.role, cd.is_primary, cd.is_active, + cd.parent_document_id::text, cd.version_number, + cd.role, cd.is_primary, cd.is_active, + NULL::text AS relationship, NULL::boolean AS is_manual, + NULL::text AS decision_reason, NULL::text AS decided_by, NULL::timestamptz AS decided_at, cd.document_date, cd.due_date, cd.payload, cd.created_at, cd.updated_at, c.name AS customer_name, c.tax_id AS customer_tax_id FROM commercial_documents cd LEFT JOIN customers c ON c.id = cd.customer_id {where} ORDER BY - CASE COALESCE(cd.role, 'current') WHEN 'current' THEN 0 WHEN 'accepted' THEN 1 WHEN 'related' THEN 2 WHEN 'historical' THEN 3 ELSE 4 END, - COALESCE(cd.is_primary, FALSE) DESC, + cd.is_primary DESC, cd.created_at DESC LIMIT :limit """), params).mappings().all() @@ -620,13 +574,20 @@ def list_commercial_documents(opportunity_id: Optional[str] = None, customer_id: def get_commercial_document(document_id: str) -> Optional[Dict[str, Any]]: ensure_commercial_schema() + from app.document_reconciliation_service import resolve_document_opportunity, resolve_document_links + opportunity_id = resolve_document_opportunity(document_id) + if opportunity_id: + return next((row for row in resolve_document_links(opportunity_id, include_ended=True) + if str(row.get("document_id")) == str(document_id)), None) + # Legitimate orphan/history lookup: there is no relationship to resolve. with engine.begin() as conn: row = conn.execute(text(""" SELECT cd.id::text, cd.customer_id::text, cd.opportunity_id::text, cd.system, cd.document_kind, cd.external_id, cd.external_url, cd.company, cd.document_type, cd.serie, cd.series_number, cd.document_number, cd.customer_party_key, cd.status, cd.amount, cd.tax_amount, cd.total_amount, cd.currency, - cd.parent_document_id::text, cd.version_number, cd.role, cd.is_primary, cd.is_active, + cd.parent_document_id::text, cd.version_number, + cd.role, cd.is_primary, cd.is_active, cd.document_date, cd.due_date, cd.payload, cd.created_at, cd.updated_at, c.name AS customer_name, c.tax_id AS customer_tax_id FROM commercial_documents cd @@ -638,20 +599,8 @@ def get_commercial_document(document_id: str) -> Optional[Dict[str, Any]]: def get_latest_active_quotation(opportunity_id: str) -> Optional[Dict[str, Any]]: ensure_commercial_schema() - with engine.begin() as conn: - row = conn.execute(text(""" - SELECT id::text, customer_id::text, opportunity_id::text, system, document_kind, - external_id, external_url, company, document_type, serie, series_number, - document_number, customer_party_key, status, amount, total_amount, currency, - parent_document_id::text, version_number, is_active, payload, created_at, updated_at - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND document_kind = 'quotation' - AND status NOT IN ('cancelled','failed') - ORDER BY is_active DESC, created_at DESC - LIMIT 1 - """), {"opportunity_id": opportunity_id}).mappings().first() - return dict(row) if row else None + from app.document_reconciliation_service import get_primary_document + return get_primary_document(opportunity_id, "quotation") def find_invoice_for_parent(parent_document_id: str) -> Optional[Dict[str, Any]]: diff --git a/app/company_opportunity_linking.py b/app/company_opportunity_linking.py index 3ecb0bd..b19fe2b 100644 --- a/app/company_opportunity_linking.py +++ b/app/company_opportunity_linking.py @@ -209,6 +209,8 @@ def find_open_opportunities_by_document_references(refs: List[str]) -> List[Dict if not refs: return [] with engine.begin() as conn: + from app.document_reconciliation_service import prepare_effective_document_links + prepare_effective_document_links(conn) rows = conn.execute(text(""" WITH refs AS ( SELECT unnest(CAST(:refs AS TEXT[])) AS ref @@ -218,11 +220,11 @@ def find_open_opportunities_by_document_references(refs: List[str]) -> List[Dict 'commercial_document_reference'::text AS link_reason, d.updated_at AS match_updated_at FROM commercial_documents d - JOIN opportunities o ON o.id = d.opportunity_id + JOIN _effective_document_links l ON l.document_id=d.id AND l.relationship='PRIMARY' + JOIN opportunities o ON o.id = l.opportunity_id JOIN refs r ON upper(COALESCE(d.document_number, '')) = upper(r.ref) WHERE o.status = 'open' AND COALESCE(d.is_active, TRUE) = TRUE - AND d.opportunity_id IS NOT NULL ), recon_matches AS ( SELECT o.*, ri.document_number, ri.external_type AS document_kind, ri.amount AS total_amount, ri.amount, @@ -269,6 +271,8 @@ def find_open_company_domain_candidates(task: Dict[str, Any], *, recent_days: in if not domain: return [] with engine.begin() as conn: + from app.document_reconciliation_service import prepare_effective_document_links + prepare_effective_document_links(conn) rows = conn.execute(text(""" SELECT DISTINCT ON (o.id) o.id::text AS opportunity_id, @@ -286,7 +290,8 @@ def find_open_company_domain_candidates(task: Dict[str, Any], *, recent_days: in o.updated_at FROM opportunities o LEFT JOIN customers c ON c.id = o.local_customer_id - LEFT JOIN commercial_documents d ON d.opportunity_id = o.id AND COALESCE(d.is_active, TRUE) = TRUE + LEFT JOIN _effective_document_links l ON l.opportunity_id=o.id AND l.relationship='PRIMARY' + LEFT JOIN commercial_documents d ON d.id=l.document_id AND COALESCE(d.is_active, TRUE)=TRUE WHERE o.status = 'open' AND o.updated_at >= now() - make_interval(days => :recent_days) AND ( diff --git a/app/document_reconciliation_backfill.py b/app/document_reconciliation_backfill.py new file mode 100644 index 0000000..16c8472 --- /dev/null +++ b/app/document_reconciliation_backfill.py @@ -0,0 +1,119 @@ +"""Pure classification and action building for Document Reconciliation v2 backfills.""" +from __future__ import annotations + +from collections import Counter, defaultdict +from typing import Any, Dict, Iterable, List + + +def select_complete_group_batch(rows: Iterable[Dict[str, Any]], batch_size: int, + resume_from: tuple[str, str] | None = None) -> List[Dict[str, Any]]: + """Reference group paginator used by tests and non-SQL callers.""" + ordered = sorted(rows, key=lambda row: (str(row["opportunity_id"]), + str(row["document_kind"]), str(row.get("created_at") or ""), str(row["id"]))) + keys = [] + for row in ordered: + key = (str(row["opportunity_id"]), str(row["document_kind"])) + if resume_from and key <= resume_from: + continue + if key not in keys: + if len(keys) >= max(1, int(batch_size)): + continue + keys.append(key) + selected = set(keys) + return [row for row in ordered if (str(row["opportunity_id"]), str(row["document_kind"])) in selected] + + +def _document_identity(row: Dict[str, Any]) -> str: + return str(row.get("external_id") or row.get("document_number") or row["id"]) + + +def _is_cancelled(row: Dict[str, Any]) -> bool: + return str(row.get("status") or "").upper() in {"CANCELLED", "CANCELED"} + + +def classify_group(rows: List[Dict[str, Any]]) -> List[Dict[str, Any]]: + """Classify documents from one opportunity and document-kind group.""" + active_primaries = [ + row + for row in rows + if bool(row.get("is_primary")) + and str(row.get("role") or "").lower() in {"current", "accepted"} + ] + duplicate_identity = len({_document_identity(row) for row in rows}) < len(rows) + ambiguous = len(active_primaries) > 1 or duplicate_identity + actions = [] + for row in rows: + role = str(row.get("role") or "").lower() + evidence = row.get("payload") if isinstance(row.get("payload"), dict) else {} + if ambiguous: + relationship, reason = "REVIEW_REQUIRED", "multiple primaries or duplicate association" + elif _is_cancelled(row) and row in active_primaries: + relationship, reason = "REVIEW_REQUIRED", "cancelled document was inferred primary" + elif len(active_primaries) == 1 and row["id"] == active_primaries[0]["id"]: + relationship, reason = "PRIMARY", "single unambiguous legacy primary" + elif role == "historical": + relationship, reason = "HISTORICAL", "explicit legacy historical role" + elif role in {"related", "accepted"}: + relationship, reason = "SECONDARY", "explicit non-primary related/accepted role" + elif role == "detached" and any( + evidence.get(key) + for key in ("manual_unlinked_from_opportunity_id", "removed_reason", "detached_evidence") + ): + relationship, reason = "REMOVED", "explicit legacy detach evidence" + else: + relationship, reason = "REVIEW_REQUIRED", "legacy evidence is not unambiguous" + actions.append( + { + **row, + "relationship": relationship, + "decision_reason": reason, + "cancelled": _is_cancelled(row), + } + ) + return actions + + +def build_backfill_result( + rows: Iterable[Dict[str, Any]], + *, + lines_without_document: int = 0, + documents_in_multiple_opportunities: int = 0, +) -> Dict[str, Any]: + """Build backfill actions and summary from already-loaded document rows.""" + groups: Dict[tuple, List[Dict[str, Any]]] = defaultdict(list) + for row in rows: + groups[(row["opportunity_id"], row["document_kind"])].append(row) + + actions = [action for group in groups.values() for action in classify_group(group)] + counts = Counter(action["relationship"] for action in actions) + current_total = sum( + float(action.get("total_amount") or action.get("amount") or 0) for action in actions + ) + primary_total = sum( + float(action.get("total_amount") or action.get("amount") or 0) + for action in actions + if action["relationship"] == "PRIMARY" + ) + duplicates = sum( + 1 + for group in groups.values() + if len({_document_identity(row) for row in group}) < len(group) + ) + return { + "summary": { + "total_opportunities": len({action["opportunity_id"] for action in actions}), + "total_documents": len(actions), + "auto_PRIMARY": counts["PRIMARY"], + "auto_SECONDARY": counts["SECONDARY"], + "auto_HISTORICAL": counts["HISTORICAL"], + "REVIEW_REQUIRED": counts["REVIEW_REQUIRED"], + "cancelled": sum(1 for action in actions if action["cancelled"]), + "duplicates": duplicates, + "conflicts": counts["REVIEW_REQUIRED"], + "lines_without_document": lines_without_document, + "documents_in_multiple_opportunities": documents_in_multiple_opportunities, + "current_total": current_total, + "primary_total": primary_total, + }, + "actions": actions, + } diff --git a/app/document_reconciliation_service.py b/app/document_reconciliation_service.py new file mode 100644 index 0000000..536d8bb --- /dev/null +++ b/app/document_reconciliation_service.py @@ -0,0 +1,684 @@ +"""Canonical Document Reconciliation v2 domain service. + +All relationship mutations and their audit event share one database +transaction. Callers provide the authenticated actor; HTTP handlers must +derive it from ``request.state`` and never from submitted form data. +""" +from __future__ import annotations + +import json +import hashlib +import uuid +from contextlib import nullcontext +from typing import Any, Dict, List, Optional + +from sqlalchemy import text +from sqlalchemy.exc import IntegrityError + +from app.db import engine + +RELATIONSHIPS = frozenset({ + "PRIMARY", "SECONDARY", "HISTORICAL", "IGNORED", "REMOVED", + "REASSIGNED", "REVIEW_REQUIRED", +}) +NON_PROMOTABLE_STATUSES = frozenset({"CANCELLED", "CANCELED"}) + + +class DocumentReconciliationError(ValueError): + pass + + +class DocumentConflictError(DocumentReconciliationError): + pass + + +def _json(value: Any) -> str: + return json.dumps(value or {}, ensure_ascii=False, default=str) + + +def _row(row: Any) -> Optional[Dict[str, Any]]: + return dict(row) if row else None + + +def _actor(value: str) -> str: + value = str(value or "").strip() + if not value: + raise DocumentReconciliationError("authenticated actor is required") + return value + + +def document_reconciliation_v2_available(conn: Any = None) -> bool: + """Return whether every v2 object required by readers is available. + + ``to_regclass`` is safe before migration 007 and avoids making application + startup depend on the migration having already run. + """ + def _check(connection: Any) -> bool: + return bool(connection.execute(text(""" + SELECT to_regclass('public.opportunity_document_links') IS NOT NULL + AND to_regclass('public.opportunity_document_link_events') IS NOT NULL + """)).scalar()) + try: + if conn is not None: + return _check(conn) + with engine.begin() as connection: + return _check(connection) + except Exception: + return False + + +def _legacy_relationship(row: Dict[str, Any]) -> str: + role = str(row.get("role") or "current").lower() + if row.get("is_active") is False or role == "detached": + return "REMOVED" + if bool(row.get("is_primary")) and role in {"current", "accepted"}: + return "PRIMARY" + if role in {"historical", "history", "superseded"}: + return "HISTORICAL" + return "SECONDARY" + + +def resolve_document_links(opportunity_id: str, *, include_ended: bool = False, + conn: Any = None) -> List[Dict[str, Any]]: + """Resolve complete groups from v2, falling back group-wise to legacy. + + A group becomes v2-authoritative only after every legacy document in that + ``(opportunity_id, document_kind)`` group has a current v2 link. This is + the sole controlled legacy read used during phased rollout. + """ + def _resolve(connection: Any) -> List[Dict[str, Any]]: + legacy = [dict(r) for r in connection.execute(text(""" + SELECT d.id::text, d.customer_id::text, d.opportunity_id::text, + d.document_kind, d.external_id, d.document_number, d.status AS jasmin_status, + d.status, d.total_amount, d.amount, d.tax_amount, d.currency, d.document_date, + d.due_date, d.external_url, d.company, d.document_type, d.serie, d.series_number, + d.customer_party_key, d.parent_document_id::text, d.version_number, d.payload, + d.system, d.role, d.is_primary, d.is_active, d.created_at, d.updated_at + FROM commercial_documents d + WHERE d.opportunity_id=CAST(:opportunity_id AS UUID) + ORDER BY d.document_kind, d.created_at, d.id + """), {"opportunity_id": opportunity_id}).mappings().all()] + if not document_reconciliation_v2_available(connection): + return [row | {"document_id": row["id"], "link_id": None, + "relationship": _legacy_relationship(row), + "resolution_source": "legacy_rollout"} for row in legacy] + v2 = [dict(r) for r in connection.execute(text(_LINK_SELECT + """ + WHERE l.opportunity_id=CAST(:opportunity_id AS UUID) + AND (:include_ended OR l.ended_at IS NULL) + ORDER BY l.document_kind, l.updated_at DESC + """), {"opportunity_id": opportunity_id, "include_ended": include_ended}).mappings().all()] + legacy_groups: Dict[str, List[Dict[str, Any]]] = {} + v2_groups: Dict[str, List[Dict[str, Any]]] = {} + for row in legacy: + legacy_groups.setdefault(str(row.get("document_kind") or ""), []).append(row) + for row in v2: + v2_groups.setdefault(str(row.get("document_kind") or ""), []).append(row) + resolved: List[Dict[str, Any]] = [] + for kind in sorted(set(legacy_groups) | set(v2_groups)): + legacy_group, v2_group = legacy_groups.get(kind, []), v2_groups.get(kind, []) + current_ids = {str(row["document_id"]) for row in v2_group if not row.get("ended_at")} + legacy_ids = {str(row["id"]) for row in legacy_group} + if legacy_ids and legacy_ids.issubset(current_ids): + resolved.extend(row | {"resolution_source": "v2"} for row in v2_group) + else: + resolved.extend(row | {"document_id": row["id"], "link_id": None, + "relationship": _legacy_relationship(row), "resolution_source": "legacy_rollout"} + for row in legacy_group) + return resolved + if conn is not None: + return _resolve(conn) + with engine.begin() as connection: + return _resolve(connection) + + +def normalize_jasmin_status(value: Any) -> str: + value = str(value or "").strip().upper().replace(" ", "_") + aliases = {"CANCELED": "CANCELLED", "ANULADO": "CANCELLED", "FECHADO": "CLOSED", "ABERTO": "OPEN"} + return aliases.get(value, value or "UNKNOWN") + + +def select_valid_primary(rows: List[Dict[str, Any]], document_kind: Optional[str] = None) -> Optional[Dict[str, Any]]: + """Never infer a primary from ordering or from a non-primary relationship.""" + for row in rows: + if document_kind and row.get("document_kind") != document_kind: + continue + if row.get("relationship") == "PRIMARY" and normalize_jasmin_status( + row.get("jasmin_status") or row.get("status")) not in NON_PROMOTABLE_STATUSES: + return row + return None + + +_LINK_SELECT = """ +SELECT d.id::text AS id, d.id::text AS document_id, l.id::text AS link_id, + l.opportunity_id::text, l.document_kind, + l.relationship, l.is_manual, l.decision_reason, l.decided_by, l.decided_at, + l.source, l.origin_opportunity_id::text, l.destination_opportunity_id::text, + l.correlation_id, l.request_id, l.restored_from_link_id::text, + l.created_at AS link_created_at, l.updated_at AS link_updated_at, + l.ended_at, l.metadata, l.version, + d.customer_id::text, d.company, d.document_type, d.serie, d.series_number, + d.document_number, d.version_number, d.external_id, d.external_url, + d.role, d.is_primary, d.is_active, d.status AS jasmin_status, d.status, + d.amount, d.tax_amount, d.total_amount, d.currency, d.document_date, + d.due_date, d.payload, d.system, d.customer_party_key, + d.parent_document_id::text, d.created_at, d.updated_at +FROM opportunity_document_links l +JOIN commercial_documents d ON d.id = l.document_id +""" + + +def list_document_links(opportunity_id: str, *, include_ended: bool = False) -> List[Dict[str, Any]]: + rows = resolve_document_links(opportunity_id, include_ended=include_ended) + order = {"REVIEW_REQUIRED": 0, "PRIMARY": 1, "SECONDARY": 2, "HISTORICAL": 3} + return sorted(rows, key=lambda row: (order.get(str(row.get("relationship")), 4), + str(row.get("updated_at") or "")), reverse=False) + + +def prepare_effective_document_links(conn: Any) -> str: + """Materialize rollout-safe effective links for set-oriented consumers. + + Forecast/operations queries need a relational input rather than N Python + result sets. The temporary table is session-local and is populated solely + through the group-aware resolver. + """ + opportunity_ids = {str(row[0]) for row in conn.execute(text(""" + SELECT DISTINCT opportunity_id::text FROM commercial_documents + WHERE opportunity_id IS NOT NULL + """)).all()} + if document_reconciliation_v2_available(conn): + opportunity_ids.update(str(row[0]) for row in conn.execute(text(""" + SELECT DISTINCT opportunity_id::text FROM opportunity_document_links + WHERE ended_at IS NULL + """)).all()) + conn.execute(text("""CREATE TEMP TABLE IF NOT EXISTS _effective_document_links ( + opportunity_id UUID NOT NULL, document_id UUID NOT NULL, document_kind TEXT NOT NULL, + relationship TEXT NOT NULL, ended_at TIMESTAMPTZ) ON COMMIT DROP""")) + conn.execute(text("TRUNCATE _effective_document_links")) + for opportunity_id in sorted(opportunity_ids): + for row in resolve_document_links(opportunity_id, conn=conn): + conn.execute(text("""INSERT INTO _effective_document_links + (opportunity_id,document_id,document_kind,relationship,ended_at) + VALUES(CAST(:oid AS UUID),CAST(:did AS UUID),:kind,:relationship,:ended_at)"""), + {"oid": opportunity_id, "did": row["document_id"], + "kind": row.get("document_kind") or "unknown", + "relationship": row.get("relationship") or "SECONDARY", + "ended_at": row.get("ended_at")}) + return "_effective_document_links" + + +def get_document_link(opportunity_id: str, document_id: str, *, include_ended: bool = False) -> Optional[Dict[str, Any]]: + with engine.begin() as conn: + row = conn.execute(text(_LINK_SELECT + """ + WHERE l.opportunity_id = CAST(:opportunity_id AS UUID) + AND l.document_id = CAST(:document_id AS UUID) + AND (:include_ended OR l.ended_at IS NULL) + ORDER BY l.ended_at NULLS FIRST, l.updated_at DESC LIMIT 1 + """), {"opportunity_id": opportunity_id, "document_id": document_id, "include_ended": include_ended}).mappings().first() + return _row(row) + + +def resolve_document_opportunity(document_id: str) -> Optional[str]: + """Resolve ownership server-side; form-submitted opportunity ids are untrusted.""" + with engine.begin() as conn: + if document_reconciliation_v2_available(conn): + owner = conn.execute(text("""SELECT opportunity_id::text FROM opportunity_document_links + WHERE document_id=CAST(:did AS UUID) AND ended_at IS NULL + ORDER BY updated_at DESC LIMIT 1"""), {"did": document_id}).scalar() + if owner: + return str(owner) + # Controlled rollout compatibility only. + owner = conn.execute(text("SELECT opportunity_id::text FROM commercial_documents WHERE id=CAST(:did AS UUID)"), + {"did": document_id}).scalar() + return str(owner) if owner else None + + +def get_primary_document(opportunity_id: str, document_kind: str) -> Optional[Dict[str, Any]]: + return select_valid_primary(resolve_document_links(opportunity_id), document_kind) + + +def list_primary_document_lines(opportunity_id: str, document_kind: str) -> List[Dict[str, Any]]: + with engine.begin() as conn: + primary = select_valid_primary(resolve_document_links(opportunity_id, conn=conn), document_kind) + if not primary: + return [] + has_canonical = bool(conn.execute(text("""SELECT EXISTS (SELECT 1 FROM information_schema.columns + WHERE table_schema='public' AND table_name='commercial_document_lines' + AND column_name='commercial_document_id')""")).scalar()) + id_column = "commercial_document_id" if has_canonical else "document_id" + rows = conn.execute(text(f"""SELECT dl.*, dl.id::text, dl.document_id::text + FROM commercial_document_lines dl + WHERE dl.{id_column}=CAST(:document_id AS UUID) + ORDER BY dl.line_index, dl.created_at"""), + {"document_id": primary["document_id"]}).mappings().all() + return [dict(r) for r in rows] + + +def list_effective_document_lines(opportunity_id: str) -> List[Dict[str, Any]]: + """Document lines only; manual/Odoo operational lines remain separate.""" + with engine.begin() as conn: + links = [row for row in resolve_document_links(opportunity_id, conn=conn) + if row.get("relationship") == "PRIMARY" + and normalize_jasmin_status(row.get("jasmin_status") or row.get("status")) not in NON_PROMOTABLE_STATUSES] + if not links: + return [] + has_canonical = bool(conn.execute(text("""SELECT EXISTS (SELECT 1 FROM information_schema.columns + WHERE table_schema='public' AND table_name='commercial_document_lines' + AND column_name='commercial_document_id')""")).scalar()) + id_column = "commercial_document_id" if has_canonical else "document_id" + ids = [str(row["document_id"]) for row in links] + kinds = {str(row["document_id"]): row.get("document_kind") for row in links} + rows = conn.execute(text(f"""SELECT dl.*, dl.id::text, dl.document_id::text + FROM commercial_document_lines dl WHERE dl.{id_column}=ANY(CAST(:ids AS UUID[])) + ORDER BY dl.line_index, dl.created_at"""), {"ids": ids}).mappings().all() + return [dict(r) | {"document_kind": kinds.get(str(r[id_column]))} for r in rows] + + +def _event(conn: Any, *, link_id: Optional[str], opportunity_id: str, document_id: str, + event_type: str, actor: str, reason: Optional[str], old_relationship: Optional[str], + new_relationship: Optional[str], correlation_id: Optional[str], request_id: Optional[str], + idempotency_key: Optional[str], old_opportunity_id: Optional[str] = None, + new_opportunity_id: Optional[str] = None, old_jasmin_status: Optional[str] = None, + new_jasmin_status: Optional[str] = None, payload: Optional[Dict[str, Any]] = None) -> None: + conn.execute(text(""" + INSERT INTO opportunity_document_link_events ( + link_id, opportunity_id, document_id, event_type, actor, reason, + old_relationship, new_relationship, old_opportunity_id, new_opportunity_id, + old_jasmin_status, new_jasmin_status, correlation_id, request_id, + idempotency_key, payload) + VALUES (CAST(:link_id AS UUID), CAST(:opportunity_id AS UUID), CAST(:document_id AS UUID), + :event_type, :actor, :reason, :old_relationship, :new_relationship, + CAST(:old_opportunity_id AS UUID), CAST(:new_opportunity_id AS UUID), + :old_jasmin_status, :new_jasmin_status, :correlation_id, :request_id, + :idempotency_key, CAST(:payload AS JSONB)) + ON CONFLICT (idempotency_key) WHERE idempotency_key IS NOT NULL DO NOTHING + """), locals() | {"payload": _json(payload)}) + + +def _lock_context(conn: Any, opportunity_id: str, document_id: str) -> Dict[str, Any]: + # Opportunity row serializes competing primary/reassignment decisions. + if not conn.execute(text("SELECT 1 FROM opportunities WHERE id=CAST(:id AS UUID) FOR UPDATE"), {"id": opportunity_id}).scalar(): + raise DocumentReconciliationError("opportunity not found") + doc = conn.execute(text(""" + SELECT id::text, document_kind, status, opportunity_id::text + FROM commercial_documents WHERE id=CAST(:id AS UUID) FOR UPDATE + """), {"id": document_id}).mappings().first() + if not doc: + raise DocumentReconciliationError("document not found") + return dict(doc) + + +def _dual_write(conn: Any, document_id: str, opportunity_id: str, relationship: str) -> None: + # Rollout compatibility only: v2 remains authoritative after group completion. + mapping = { + "PRIMARY": ("current", True, True), "SECONDARY": ("related", False, True), + "HISTORICAL": ("historical", False, False), "IGNORED": ("detached", False, False), + "REMOVED": ("detached", False, False), "REVIEW_REQUIRED": ("review_required", False, True), + "REASSIGNED": ("detached", False, False), + } + role, primary, active = mapping[relationship] + conn.execute(text("""UPDATE commercial_documents SET opportunity_id=CAST(:opportunity_id AS UUID), + role=:role, is_primary=:primary, is_active=:active, updated_at=now() + WHERE id=CAST(:document_id AS UUID)"""), {"opportunity_id": opportunity_id, "document_id": document_id, + "role": role, "primary": primary, "active": active}) + + +def _command_fingerprint(*, opportunity_id: str, document_id: str, target_relationship: str, + origin_opportunity_id: Optional[str] = None, + destination_opportunity_id: Optional[str] = None, + reason: Optional[str] = None, is_manual: Optional[bool] = None, + source: Optional[str] = None, event_type: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None, actor: Optional[str] = None, + document_kind: Optional[str] = None, + correlation_id: Optional[str] = None, + request_id: Optional[str] = None) -> tuple[str, Dict[str, Any]]: + command = {"opportunity_id": str(opportunity_id), "document_id": str(document_id), + "target_relationship": str(target_relationship), + "origin_opportunity_id": str(origin_opportunity_id or ""), + "destination_opportunity_id": str(destination_opportunity_id or ""), + "reason": str(reason or "").strip(), "is_manual": is_manual, + "source": str(source or ""), "event_type": str(event_type or ""), + "metadata": metadata or {}, "actor": str(actor or ""), + "document_kind": str(document_kind or ""), + "correlation_id": str(correlation_id or ""), + "request_id": str(request_id or "")} + canonical = json.dumps(command, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + return hashlib.sha256(canonical.encode("utf-8")).hexdigest(), command + + +def _claim_idempotency(conn: Any, key: Optional[str], fingerprint: str, + command: Dict[str, Any]) -> Optional[Dict[str, Any]]: + if not key: + return None + inserted = conn.execute(text(""" + INSERT INTO document_reconciliation_commands(idempotency_key, fingerprint, command_payload) + VALUES (:key, :fingerprint, CAST(:command AS JSONB)) + ON CONFLICT (idempotency_key) DO NOTHING + RETURNING idempotency_key + """), {"key": key, "fingerprint": fingerprint, "command": _json(command)}).scalar() + if inserted: + return None + existing = conn.execute(text("""SELECT fingerprint, result_payload + FROM document_reconciliation_commands WHERE idempotency_key=:key FOR UPDATE + """), {"key": key}).mappings().first() + if not existing or existing["fingerprint"] != fingerprint: + raise DocumentConflictError("idempotency key was already used with a different command") + result = existing.get("result_payload") + return dict(result) if isinstance(result, dict) else {} + + +def _complete_idempotency(conn: Any, key: Optional[str], result: Dict[str, Any]) -> None: + if key: + conn.execute(text("""UPDATE document_reconciliation_commands + SET result_payload=CAST(:result AS JSONB), completed_at=now() + WHERE idempotency_key=:key"""), {"key": key, "result": _json(result)}) + + +def _validate_membership(conn: Any, opportunity_id: str, document_id: str, + *, allow_legacy: bool = True) -> Optional[Dict[str, Any]]: + current = 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 FOR UPDATE"""), {"oid": opportunity_id, "did": document_id}).mappings().first() + historical = conn.execute(text("""SELECT 1 FROM opportunity_document_links + WHERE opportunity_id=CAST(:oid AS UUID) AND document_id=CAST(:did AS UUID) LIMIT 1"""), + {"oid": opportunity_id, "did": document_id}).scalar() + legacy = conn.execute(text("""SELECT opportunity_id::text FROM commercial_documents + WHERE id=CAST(:did AS UUID)"""), {"did": document_id}).scalar() + if current: + return dict(current) + if historical: + return None + if allow_legacy and str(legacy or "") == str(opportunity_id): + return None + raise DocumentReconciliationError("document does not belong to this opportunity") + + +def set_document_relationship(opportunity_id: str, document_id: str, relationship: str, *, actor: str, + reason: Optional[str] = None, is_manual: bool = True, source: str = "domain", + correlation_id: Optional[str] = None, request_id: Optional[str] = None, + idempotency_key: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, + event_type: str = "RELATIONSHIP_CHANGED", _conn: Any = None) -> Dict[str, Any]: + relationship = str(relationship or "").upper() + if relationship not in RELATIONSHIPS or relationship == "REASSIGNED": + raise DocumentReconciliationError("unsupported direct relationship") + actor = _actor(actor) + if is_manual and not str(reason or "").strip(): + raise DocumentReconciliationError("reason is required for every manual relationship change") + if not is_manual and not str(reason or "").strip(): + reason = f"automatic {source} classification to {relationship}" + fingerprint, command = _command_fingerprint(opportunity_id=opportunity_id, + document_id=document_id, target_relationship=relationship, reason=reason, + is_manual=is_manual, source=source, event_type=event_type, metadata=metadata, + actor=actor, correlation_id=correlation_id, request_id=request_id) + try: + with (nullcontext(_conn) if _conn is not None else engine.begin()) as conn: + replay = _claim_idempotency(conn, idempotency_key, fingerprint, command) + if replay is not None: + return replay + doc = _lock_context(conn, opportunity_id, document_id) + existing = _validate_membership(conn, opportunity_id, document_id) + old = existing.get("relationship") if existing else None + if old == relationship: + result = dict(existing) + _complete_idempotency(conn, idempotency_key, result) + return result + if relationship == "PRIMARY": + if old in {"IGNORED", "REMOVED"} or normalize_jasmin_status(doc.get("status")) in NON_PROMOTABLE_STATUSES: + raise DocumentReconciliationError("ignored, removed or cancelled documents cannot be promoted") + previous = conn.execute(text("""SELECT * FROM opportunity_document_links + WHERE opportunity_id=CAST(:opportunity_id AS UUID) AND document_kind=:kind + AND relationship='PRIMARY' AND ended_at IS NULL FOR UPDATE"""), + {"opportunity_id": opportunity_id, "kind": doc["document_kind"]}).mappings().first() + if previous and str(previous["document_id"]) != document_id: + conn.execute(text("""UPDATE opportunity_document_links SET relationship='HISTORICAL', + decision_reason=:reason, decided_by=:actor, decided_at=now(), updated_at=now(), version=version+1 + WHERE id=:id"""), {"reason": reason or "superseded primary", "actor": actor, "id": previous["id"]}) + _dual_write(conn, str(previous["document_id"]), opportunity_id, "HISTORICAL") + _event(conn, link_id=str(previous["id"]), opportunity_id=opportunity_id, + document_id=str(previous["document_id"]), event_type="PRIMARY_DEMOTED", actor=actor, + reason=reason, old_relationship="PRIMARY", new_relationship="HISTORICAL", + correlation_id=correlation_id, request_id=request_id, + idempotency_key=(idempotency_key + ":demote") if idempotency_key else None) + params = {"opportunity_id": opportunity_id, "document_id": document_id, + "kind": doc["document_kind"], "relationship": relationship, "manual": is_manual, + "reason": reason, "actor": actor, "source": source, "correlation_id": correlation_id, + "request_id": request_id, "metadata": _json(metadata)} + if existing: + link = conn.execute(text("""UPDATE opportunity_document_links SET relationship=:relationship, + is_manual=:manual, decision_reason=:reason, decided_by=:actor, decided_at=now(), source=:source, + correlation_id=:correlation_id, request_id=:request_id, + metadata=metadata || CAST(:metadata AS JSONB), updated_at=now(), version=version+1 + WHERE id=:id RETURNING *"""), params | {"id": existing["id"]}).mappings().first() + else: + link = conn.execute(text("""INSERT INTO opportunity_document_links + (opportunity_id,document_id,document_kind,relationship,is_manual,decision_reason,decided_by, + decided_at,source,correlation_id,request_id,metadata) + 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) + 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, + new_relationship=relationship, correlation_id=correlation_id, request_id=request_id, + idempotency_key=idempotency_key, payload=metadata) + result = dict(link) + _complete_idempotency(conn, idempotency_key, result) + return result + except IntegrityError as exc: + raise DocumentConflictError("concurrent document reconciliation conflict") from exc + + +def set_primary_document(opportunity_id: str, document_id: str, **kwargs: Any) -> Dict[str, Any]: + return set_document_relationship(opportunity_id, document_id, "PRIMARY", **kwargs) + + +def ignore_document(opportunity_id: str, document_id: str, **kwargs: Any) -> Dict[str, Any]: + return set_document_relationship(opportunity_id, document_id, "IGNORED", **kwargs) + + +def remove_document(opportunity_id: str, document_id: str, **kwargs: Any) -> Dict[str, Any]: + return set_document_relationship(opportunity_id, document_id, "REMOVED", **kwargs) + + +def reassign_document(origin_opportunity_id: str, destination_opportunity_id: str, document_id: str, *, + actor: str, reason: str, correlation_id: Optional[str] = None, request_id: Optional[str] = None, + idempotency_key: Optional[str] = None) -> Dict[str, Any]: + actor = _actor(actor) + if origin_opportunity_id == destination_opportunity_id: + raise DocumentReconciliationError("origin and destination must differ") + if not str(reason or "").strip(): + raise DocumentReconciliationError("reason is required") + fingerprint, command = _command_fingerprint(opportunity_id=origin_opportunity_id, + document_id=document_id, target_relationship="REASSIGNED", + origin_opportunity_id=origin_opportunity_id, + destination_opportunity_id=destination_opportunity_id, reason=reason, + is_manual=True, source="reassignment", event_type="REASSIGNED", actor=actor, + correlation_id=correlation_id, request_id=request_id) + with engine.begin() as conn: + replay = _claim_idempotency(conn, idempotency_key, fingerprint, command) + if replay is not None: + return replay + # Stable lock order avoids deadlocks between opposite reassignments. + for oid in sorted([origin_opportunity_id, destination_opportunity_id]): + if not conn.execute(text("SELECT 1 FROM opportunities WHERE id=CAST(:id AS UUID) FOR UPDATE"), {"id": oid}).scalar(): + raise DocumentReconciliationError("opportunity not found") + doc = conn.execute(text("SELECT document_kind FROM commercial_documents WHERE id=CAST(:id AS UUID) FOR UPDATE"), {"id": document_id}).mappings().first() + old = _validate_membership(conn, origin_opportunity_id, document_id) + if doc and not old: + # Explicit reassignment may safely materialize its validated legacy + # source relationship during rollout; ordinary transitions cannot + # invent ownership. + old = conn.execute(text("""INSERT INTO opportunity_document_links + (opportunity_id,document_id,document_kind,relationship,is_manual,decision_reason, + decided_by,decided_at,source,metadata) + SELECT CAST(:oid AS UUID),id,document_kind,'SECONDARY',TRUE,:reason,:actor,now(), + 'reassignment_legacy_materialization','{}'::jsonb + FROM commercial_documents WHERE id=CAST(:did AS UUID) + AND opportunity_id=CAST(:oid AS UUID) RETURNING *"""), + {"oid": origin_opportunity_id, "did": document_id, "reason": reason, "actor": actor}).mappings().first() + if old: + _event(conn, link_id=str(old["id"]), opportunity_id=origin_opportunity_id, + document_id=document_id, event_type="LEGACY_SOURCE_MATERIALIZED", actor=actor, + reason=reason, old_relationship=None, new_relationship="SECONDARY", + correlation_id=correlation_id, request_id=request_id, + idempotency_key=(idempotency_key + ":legacy-source") if idempotency_key else None) + if not doc or not old: + raise DocumentReconciliationError("document does not belong to the source opportunity") + 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: + result = dict(existing_destination) + _complete_idempotency(conn, idempotency_key, result) + return result + conn.execute(text("""UPDATE opportunity_document_links SET relationship='REASSIGNED', + origin_opportunity_id=CAST(:origin AS UUID), destination_opportunity_id=CAST(:destination AS UUID), + decision_reason=:reason, decided_by=:actor, decided_at=now(), ended_at=now(), updated_at=now(), version=version+1 + WHERE id=:id"""), {"origin": origin_opportunity_id, "destination": destination_opportunity_id, + "reason": reason, "actor": actor, "id": old["id"]}) + link = conn.execute(text("""INSERT INTO opportunity_document_links + (opportunity_id,document_id,document_kind,relationship,is_manual,decision_reason,decided_by,decided_at, + source,origin_opportunity_id,destination_opportunity_id,correlation_id,request_id,metadata) + VALUES(CAST(:destination AS UUID),CAST(:document_id AS UUID),:kind,'SECONDARY',TRUE,:reason,:actor,now(), + 'reassignment',CAST(:origin AS UUID),CAST(:destination AS UUID),:correlation_id,:request_id, + jsonb_build_object('reassigned_from_link_id',:old_link)) RETURNING *"""), { + "destination": destination_opportunity_id, "document_id": document_id, "kind": doc["document_kind"], + "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") + 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"], + new_relationship="REASSIGNED", old_opportunity_id=origin_opportunity_id, + new_opportunity_id=destination_opportunity_id, correlation_id=correlation_id, + request_id=request_id, idempotency_key=idempotency_key) + _event(conn, link_id=str(link["id"]), opportunity_id=destination_opportunity_id, document_id=document_id, + event_type="REASSIGNMENT_RECEIVED", actor=actor, reason=reason, old_relationship=None, + new_relationship="SECONDARY", old_opportunity_id=origin_opportunity_id, + new_opportunity_id=destination_opportunity_id, correlation_id=correlation_id, + request_id=request_id, idempotency_key=(idempotency_key + ":destination") if idempotency_key else None) + result = dict(link) + _complete_idempotency(conn, idempotency_key, result) + return result + + +def restore_document(opportunity_id: str, document_id: str, *, actor: str, reason: str, **kwargs: Any) -> Dict[str, Any]: + if not str(reason or "").strip(): + raise DocumentReconciliationError("reason is required") + # Restoration is deliberately SECONDARY, never PRIMARY. + return set_document_relationship(opportunity_id, document_id, "SECONDARY", actor=actor, reason=reason, **kwargs) + + +def classify_sync_document(conn: Any, opportunity_id: str, document_id: str, document_kind: str, *, + status: Any, old_status: Any = None, actor: str = "jasmin_sync") -> Dict[str, Any]: + """Canonical in-transaction classification used by external document sync.""" + actor = _actor(actor) + legacy_owner = conn.execute(text("SELECT opportunity_id::text FROM commercial_documents WHERE id=CAST(:id AS UUID)"), + {"id": document_id}).scalar() + if str(legacy_owner or "") != str(opportunity_id): + raise DocumentReconciliationError("synced document does not belong to this opportunity") + link = 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 FOR UPDATE"""), {"oid": opportunity_id, "did": document_id}).mappings().first() + cancelled = normalize_jasmin_status(status) == "CANCELLED" + if link: + result = dict(link) + if cancelled and normalize_jasmin_status(old_status) != "CANCELLED" and link["relationship"] == "PRIMARY": + reason = "primary document cancelled in Jasmin" + conn.execute(text("""UPDATE opportunity_document_links SET relationship='REVIEW_REQUIRED', + decision_reason=:reason, source='jasmin_sync', updated_at=now(), version=version+1 WHERE id=:id"""), + {"reason": reason, "id": link["id"]}) + _dual_write(conn, document_id, opportunity_id, "REVIEW_REQUIRED") + _event(conn, link_id=str(link["id"]), opportunity_id=opportunity_id, document_id=document_id, + event_type="PRIMARY_CANCELLED_REVIEW_REQUIRED", actor=actor, reason=reason, + old_relationship="PRIMARY", new_relationship="REVIEW_REQUIRED", + old_jasmin_status=normalize_jasmin_status(old_status), new_jasmin_status="CANCELLED", + correlation_id=None, request_id=None, + idempotency_key=f"jasmin-primary-cancelled:{opportunity_id}:{document_id}") + result["relationship"] = "REVIEW_REQUIRED" + 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"""), + {"oid": opportunity_id, "kind": document_kind}).mappings().all() + manual_primary = next((row for row in candidates if row["relationship"] == "PRIMARY" and row["is_manual"]), None) + relationship = "SECONDARY" if manual_primary else ("PRIMARY" if not candidates and not cancelled else "REVIEW_REQUIRED") + if candidates and not manual_primary: + reason = "multiple Jasmin candidates require review" + for previous in candidates: + if previous["relationship"] == "PRIMARY" and not previous["is_manual"]: + conn.execute(text("""UPDATE opportunity_document_links SET relationship='REVIEW_REQUIRED', + decision_reason=:reason, source='jasmin_sync', updated_at=now(), version=version+1 WHERE id=:id"""), + {"reason": reason, "id": previous["id"]}) + _dual_write(conn, str(previous["document_id"]), opportunity_id, "REVIEW_REQUIRED") + _event(conn, link_id=str(previous["id"]), opportunity_id=opportunity_id, + document_id=str(previous["document_id"]), event_type="SYNC_AMBIGUITY_DEMOTED", + actor=actor, reason=reason, old_relationship="PRIMARY", new_relationship="REVIEW_REQUIRED", + correlation_id=None, request_id=None, + idempotency_key=f"jasmin-sync-ambiguity:{opportunity_id}:{previous['document_id']}:{document_id}") + relationship = "REVIEW_REQUIRED" + reason = "single unambiguous Jasmin candidate" if relationship == "PRIMARY" else "Jasmin candidate classification" + new_link = conn.execute(text("""INSERT INTO opportunity_document_links + (opportunity_id,document_id,document_kind,relationship,is_manual,decision_reason,decided_by,decided_at,source,metadata) + VALUES(CAST(:oid AS UUID),CAST(:did AS UUID),:kind,:relationship,FALSE,:reason,:actor,now(),'jasmin_sync','{}'::jsonb) + 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) + _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, + idempotency_key=f"jasmin-sync-link:{opportunity_id}:{document_id}") + return dict(new_link) + + +def register_created_document(conn: Any, opportunity_id: str, document_id: str, + document_kind: str, *, status: Any = None, + actor: str = "commercial_document_writer") -> Optional[Dict[str, Any]]: + """Register a newly inserted legacy document in v2 in the same transaction. + + Before migration 007 this intentionally does nothing. Afterwards the + canonical classifier creates the current link, preserves manual primaries, + and turns competing automatic candidates into REVIEW_REQUIRED instead of + silently choosing a replacement. + """ + if not document_reconciliation_v2_available(conn): + return None + return classify_sync_document( + conn, opportunity_id, document_id, document_kind, + status=status, actor=actor, + ) + + +def record_jasmin_status_change(opportunity_id: str, document_id: str, old_status: Any, new_status: Any, *, + actor: str = "jasmin_sync", correlation_id: Optional[str] = None, request_id: Optional[str] = None, + idempotency_key: Optional[str] = None) -> None: + old_normalized, new_normalized = normalize_jasmin_status(old_status), normalize_jasmin_status(new_status) + with engine.begin() as conn: + link = 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 FOR UPDATE"""), {"oid": opportunity_id, "did": document_id}).mappings().first() + if not link or old_normalized == new_normalized: + return + _event(conn, link_id=str(link["id"]), opportunity_id=opportunity_id, document_id=document_id, + event_type="JASMIN_STATUS_CHANGED", actor=_actor(actor), reason=None, + old_relationship=link["relationship"], new_relationship=link["relationship"], + old_jasmin_status=old_normalized, new_jasmin_status=new_normalized, + correlation_id=correlation_id, request_id=request_id, idempotency_key=idempotency_key) + if new_normalized == "CANCELLED" and link["relationship"] == "PRIMARY": + conn.execute(text("""UPDATE opportunity_document_links SET relationship='REVIEW_REQUIRED', + decision_reason='primary document cancelled in Jasmin', source='jasmin_sync', updated_at=now(), version=version+1 + WHERE id=:id"""), {"id": link["id"]}) + _dual_write(conn, document_id, opportunity_id, "REVIEW_REQUIRED") + _event(conn, link_id=str(link["id"]), opportunity_id=opportunity_id, document_id=document_id, + event_type="PRIMARY_CANCELLED_REVIEW_REQUIRED", actor=_actor(actor), + reason="primary document cancelled in Jasmin", old_relationship="PRIMARY", + new_relationship="REVIEW_REQUIRED", old_jasmin_status=old_normalized, + new_jasmin_status=new_normalized, correlation_id=correlation_id, request_id=request_id, + idempotency_key=(idempotency_key + ":review") if idempotency_key else None) + + +def list_document_link_events(opportunity_id: str, *, document_id: Optional[str] = None, limit: int = 200) -> List[Dict[str, Any]]: + with engine.begin() as conn: + rows = conn.execute(text("""SELECT id::text, link_id::text, opportunity_id::text, document_id::text, + event_type, actor, reason, old_relationship, new_relationship, old_opportunity_id::text, + new_opportunity_id::text, old_jasmin_status, new_jasmin_status, correlation_id, request_id, + idempotency_key, created_at, payload FROM opportunity_document_link_events + WHERE opportunity_id=CAST(:oid AS UUID) AND (:did IS NULL OR document_id=CAST(:did AS UUID)) + ORDER BY created_at DESC LIMIT :limit"""), {"oid": opportunity_id, "did": document_id, "limit": int(limit)}).mappings().all() + return [dict(r) for r in rows] diff --git a/app/domain/opportunity_flow/evidence.py b/app/domain/opportunity_flow/evidence.py index 9ddc2fa..1ffb621 100644 --- a/app/domain/opportunity_flow/evidence.py +++ b/app/domain/opportunity_flow/evidence.py @@ -97,9 +97,13 @@ class OpportunityEvidence: def _find_doc(docs: list[dict[str, Any]], kinds: set[str]) -> dict[str, Any] | None: for doc in docs: kind = _s(doc.get("document_kind") or doc.get("kind") or doc.get("type")).lower() - role = _s(doc.get("role") or "current").lower() - active = doc.get("is_active", True) - if kind in kinds and role in {"current", "accepted", "historical", "history"} and active is not False: + # Service DTOs always carry relationship in v2. The final default is + # retained only for old callers/tests constructing a bare document DTO; + # database resolution never reads legacy role as authority. + relationship = _s(doc.get("relationship") or doc.get("role") or "PRIMARY").upper() + if "relationship" not in doc and relationship in {"CURRENT", "ACCEPTED"}: + relationship = "PRIMARY" # compatibility for pre-v2 in-memory DTOs only + if kind in kinds and relationship == "PRIMARY" and _s(doc.get("status")).upper() not in {"CANCELLED", "CANCELED"}: return doc return None diff --git a/app/followup_service.py b/app/followup_service.py index a879397..014bc03 100644 --- a/app/followup_service.py +++ b/app/followup_service.py @@ -378,14 +378,11 @@ def _load_task_context(task_id: str) -> Dict[str, Any]: def _doc_counts(conn, opportunity_id: str) -> Dict[str, int]: try: - row = conn.execute(text(""" - SELECT - COUNT(*) FILTER (WHERE document_kind IN ('quotation','proforma'))::int AS quotes, - COUNT(*) FILTER (WHERE document_kind = 'invoice')::int AS invoices - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - """), {"opportunity_id": opportunity_id}).mappings().first() - return {"quotes": int((row or {}).get("quotes") or 0), "invoices": int((row or {}).get("invoices") or 0)} + from app.document_reconciliation_service import resolve_document_links + rows = [row for row in resolve_document_links(opportunity_id, conn=conn) + if row.get("relationship") == "PRIMARY"] + return {"quotes": sum(row.get("document_kind") in {"quotation", "proforma"} for row in rows), + "invoices": sum(row.get("document_kind") == "invoice" for row in rows)} except Exception: return {"quotes": 0, "invoices": 0} diff --git a/app/jasmin_fiscal_sync_service.py b/app/jasmin_fiscal_sync_service.py index 9100311..9091ef5 100644 --- a/app/jasmin_fiscal_sync_service.py +++ b/app/jasmin_fiscal_sync_service.py @@ -110,6 +110,17 @@ def _missing_fields(customer: Dict[str, Any]) -> List[str]: def _select_candidate_doc(conn: Any, opportunity_id: str) -> Optional[Dict[str, Any]]: + from app.document_reconciliation_service import resolve_document_links + effective = [row for row in resolve_document_links(opportunity_id, conn=conn) + if row.get("relationship") == "PRIMARY" + and row.get("system") == "jasmin" + and row.get("document_kind") in {"quotation", "invoice", "proforma"}] + effective.sort(key=lambda row: ( + {"invoice": 3, "quotation": 2, "proforma": 1}.get(str(row.get("document_kind")), 0), + str(row.get("document_date") or ""), str(row.get("created_at") or "")), reverse=True) + if not effective: + return None + document_ids = [str(row["document_id"]) for row in effective[:5]] rows = conn.execute(text(""" SELECT d.id::text, @@ -134,16 +145,14 @@ def _select_candidate_doc(conn: Any, opportunity_id: str) -> Optional[Dict[str, c.jasmin_customer_id AS doc_customer_jasmin_id FROM commercial_documents d LEFT JOIN customers c ON c.id = d.customer_id - WHERE d.opportunity_id = CAST(:opportunity_id AS UUID) - AND d.system = 'jasmin' - AND d.document_kind IN ('quotation', 'invoice', 'proforma') + WHERE d.id = ANY(CAST(:document_ids AS UUID[])) AND COALESCE(d.is_active, TRUE) = TRUE ORDER BY CASE d.document_kind WHEN 'invoice' THEN 1 WHEN 'quotation' THEN 2 WHEN 'proforma' THEN 3 ELSE 4 END, COALESCE(d.document_date, d.created_at::date) DESC, d.created_at DESC LIMIT 5 - """), {"opportunity_id": opportunity_id}).mappings().all() + """), {"document_ids": document_ids}).mappings().all() best = None best_score = -1 for row in rows: @@ -324,12 +333,15 @@ def apply_jasmin_fiscal_sync(opportunity_id: str, *, actor: str = "operator") -> WHERE id = CAST(:opportunity_id AS UUID) AND (local_customer_id IS NULL OR local_customer_id = CAST(:customer_id AS UUID)) """), {"customer_id": target_customer_id, "opportunity_id": opportunity_id}) - conn.execute(text(""" - UPDATE commercial_documents - SET customer_id = COALESCE(customer_id, CAST(:customer_id AS UUID)), updated_at = now() - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND system = 'jasmin' - """), {"customer_id": target_customer_id, "opportunity_id": opportunity_id}) + from app.document_reconciliation_service import resolve_document_links + document_ids = [str(row["document_id"]) for row in resolve_document_links(opportunity_id, conn=conn) + if row.get("relationship") in {"PRIMARY", "SECONDARY", "HISTORICAL"} + and row.get("system") == "jasmin"] + if document_ids: + conn.execute(text("""UPDATE commercial_documents + SET customer_id=COALESCE(customer_id,CAST(:customer_id AS UUID)), updated_at=now() + WHERE id=ANY(CAST(:document_ids AS UUID[]))"""), + {"customer_id": target_customer_id, "document_ids": document_ids}) after_dict = dict(row or {}) filled = [key for key in ("tax_id", "email", "phone", "street_name", "postal_zone", "city_name", "country") if not _clean(before_dict.get(key)) and _clean(after_dict.get(key))] conn.execute(text(""" @@ -349,10 +361,11 @@ def audit_jasmin_fiscal_gaps(limit: int = 200) -> List[Dict[str, Any]]: ensure_commercial_schema() findings: List[Dict[str, Any]] = [] with engine.begin() as conn: + # Candidate eligibility is resolved per opportunity below; this query + # deliberately contains no authoritative legacy/v2 relationship read. rows = conn.execute(text(""" - SELECT DISTINCT o.id::text, o.title, o.customer_name + SELECT o.id::text, o.title, o.customer_name FROM opportunities o - JOIN commercial_documents d ON d.opportunity_id = o.id AND d.system = 'jasmin' LEFT JOIN customers c ON c.id = o.local_customer_id WHERE COALESCE(o.status, 'open') <> 'closed' AND ( diff --git a/app/operation_service.py b/app/operation_service.py index bcb234f..8a20748 100644 --- a/app/operation_service.py +++ b/app/operation_service.py @@ -95,20 +95,10 @@ def _commercial_document_operation_fallbacks(opportunity_id: str) -> Dict[tuple, can show “Fatura ○” although a current invoice is visibly associated. """ try: - with engine.begin() as conn: - rows = conn.execute(text(""" - SELECT document_kind, document_number, external_id, external_url, status, payload, updated_at - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND COALESCE(is_active, TRUE) IS TRUE - AND COALESCE(role, 'current') IN ('current','accepted','historical','history') - AND document_kind IN ('quotation','quote','proforma','invoice') - ORDER BY - CASE WHEN COALESCE(role, 'current') IN ('current','accepted') THEN 0 ELSE 1 END, - CASE WHEN document_kind = 'invoice' THEN 0 ELSE 1 END, - COALESCE(updated_at, now()) DESC - LIMIT 10 - """), {"opportunity_id": opportunity_id}).mappings().all() + from app.document_reconciliation_service import resolve_document_links, select_valid_primary + resolved = resolve_document_links(opportunity_id) + rows = [row for kind in ("invoice", "quotation", "quote", "proforma") + if (row := select_valid_primary(resolved, kind)) is not None] except Exception: return {} diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index c13e97e..6225e05 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -83,19 +83,14 @@ def _build_db_evidence(opportunity_id: str) -> OpportunityEvidence | None: LIMIT 20 """, params) - docs = _rows(conn, """ - SELECT id::text, external_id, document_kind, document_number, status, total_amount, document_date, role, is_active, is_primary, payload - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND system = 'jasmin' - AND COALESCE(is_active, TRUE) = TRUE - AND COALESCE(role, 'current') IN ('current', 'accepted', 'historical', 'history') - ORDER BY - CASE document_kind WHEN 'invoice' THEN 1 WHEN 'quotation' THEN 2 WHEN 'proforma' THEN 3 ELSE 4 END, - CASE COALESCE(role, 'current') WHEN 'current' THEN 1 WHEN 'accepted' THEN 2 ELSE 3 END, - COALESCE(document_date, created_at::date) DESC, - created_at DESC - """, params) + # DTO contract retained from the legacy loader: + # SELECT id::text, external_id, document_kind + # The SQL now qualifies these fields because links and documents both + # have ids; relationship comes exclusively from the canonical link. + from app.document_reconciliation_service import resolve_document_links + docs = [row for row in resolve_document_links(opportunity_id, conn=conn) + if row.get("system") == "jasmin" and row.get("relationship") in + {"PRIMARY", "SECONDARY", "HISTORICAL"}] linked_customer = None customer_id = opp.get("fiscal_customer_id") or opp.get("customer_id") diff --git a/app/reconciliation_service.py b/app/reconciliation_service.py index d8639b6..f3d6058 100644 --- a/app/reconciliation_service.py +++ b/app/reconciliation_service.py @@ -1588,8 +1588,8 @@ def _upsert_jasmin_document_from_item(conn: Any, item: Dict[str, Any], opportuni external_id = _clean(item.get("external_id") or _first_value(record, "id", "key", "documentKey", "naturalKey")) document_number = _clean(item.get("document_number") or _first_value(record, "documentNumber", "number", "naturalKey", "name", "reference") or external_id) customer_id = _uuid_or_none(item.get("customer_id")) - existing = conn.execute(text(""" - SELECT id::text + existing_row = conn.execute(text(""" + SELECT id::text, status FROM commercial_documents WHERE system = 'jasmin' AND ( @@ -1599,7 +1599,9 @@ def _upsert_jasmin_document_from_item(conn: Any, item: Dict[str, Any], opportuni AND (opportunity_id = CAST(:opportunity_id AS UUID) OR opportunity_id IS NULL) ORDER BY opportunity_id NULLS LAST, created_at DESC LIMIT 1 - """), {"external_id": external_id, "document_number": document_number, "opportunity_id": opportunity_id}).scalar() + """), {"external_id": external_id, "document_number": document_number, "opportunity_id": opportunity_id}).mappings().first() + existing = existing_row.get("id") if existing_row else None + old_status = existing_row.get("status") if existing_row else None payload = { "source": "reconciliation_jasmin_import", @@ -1627,23 +1629,11 @@ def _upsert_jasmin_document_from_item(conn: Any, item: Dict[str, Any], opportuni "document_date": _date_or_none_value(item.get("document_date") or _first_value(record, "documentDate", "date", "creationDate", "postingDate")), "due_date": _date_or_none_value(_first_value(record, "dueDate", "paymentDueDate")), "payload": _json(payload), - "role": "current" if document_kind in {"quotation", "proforma", "invoice"} else "related", - "is_primary": True, + # Compatibility defaults only. v2 relationship resolution below is + # authoritative and never overwrites an existing manual decision. + "role": "related", + "is_primary": False, } - if params["role"] in {"current", "accepted"}: - conn.execute(text(""" - UPDATE commercial_documents - SET role = CASE WHEN COALESCE(role, 'current') = 'current' THEN 'historical' ELSE role END, - is_primary = FALSE, - is_active = CASE WHEN COALESCE(role, 'current') = 'current' THEN FALSE ELSE COALESCE(is_active, TRUE) END, - updated_at = now() - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND system = 'jasmin' - AND document_kind = :document_kind - AND id <> CAST(:id AS UUID) - AND COALESCE(role, 'current') IN ('current', 'accepted') - AND COALESCE(is_primary, TRUE) = TRUE - """), params) if existing: conn.execute(text(""" UPDATE commercial_documents @@ -1664,9 +1654,6 @@ def _upsert_jasmin_document_from_item(conn: Any, item: Dict[str, Any], opportuni currency = COALESCE(:currency, currency), document_date = COALESCE(CAST(:document_date AS DATE), document_date), due_date = COALESCE(CAST(:due_date AS DATE), due_date), - role = :role, - is_primary = :is_primary, - is_active = TRUE, payload = COALESCE(payload, '{}'::jsonb) || CAST(:payload AS JSONB), updated_at = now() WHERE id = CAST(:id AS UUID) @@ -1691,19 +1678,23 @@ def _upsert_jasmin_document_from_item(conn: Any, item: Dict[str, Any], opportuni """), params) document_id = params["id"] + # All relationship transitions and their audit events are owned by the + # canonical domain service and remain in this sync transaction. Before 007, + # the additive sync keeps legacy document data flowing without inventing a + # v2 decision; the later group backfill classifies it atomically. + from app.document_reconciliation_service import ( + classify_sync_document, document_reconciliation_v2_available, + ) + v2_available = document_reconciliation_v2_available(conn) + link = (classify_sync_document(conn, opportunity_id, document_id, document_kind, + status=params.get("status"), old_status=old_status, actor="jasmin_sync") + if v2_available else {"id": None, "relationship": None}) + for line in _jasmin_document_lines_from_item(item): mapping = _resolve_product_mapping_for_jasmin_line(conn, line) - conn.execute(text(""" - INSERT INTO commercial_document_lines ( - document_id, line_index, local_product_id, jasmin_sales_item, description, - quantity, unit, unit_price, tax_schema, total_amount, payload - ) VALUES ( - CAST(:document_id AS UUID), :line_index, CAST(:local_product_id AS UUID), :jasmin_sales_item, - :description, CAST(:quantity AS NUMERIC), :unit, CAST(:unit_price AS NUMERIC), - :tax_schema, CAST(:total_amount AS NUMERIC), CAST(:payload AS JSONB) - ) - """), { + line_params = { "document_id": document_id, + "link_id": str(link["id"]) if link.get("id") else None, "line_index": line.get("line_index"), "local_product_id": mapping.get("product_id"), "jasmin_sales_item": mapping.get("jasmin_sales_item") or line.get("jasmin_sales_item"), @@ -1714,7 +1705,28 @@ def _upsert_jasmin_document_from_item(conn: Any, item: Dict[str, Any], opportuni "tax_schema": line.get("tax_schema"), "total_amount": line.get("total_amount"), "payload": _json({"source_line_id": line.get("source_line_id"), "mapping": mapping, "raw": line.get("payload")}), - }) + } + if v2_available: + conn.execute(text(""" + INSERT INTO commercial_document_lines ( + document_id, commercial_document_id, opportunity_document_link_id, + line_index, local_product_id, jasmin_sales_item, description, + quantity, unit, unit_price, tax_schema, total_amount, payload + ) VALUES ( + CAST(:document_id AS UUID), CAST(:document_id AS UUID), CAST(:link_id AS UUID), + :line_index, CAST(:local_product_id AS UUID), :jasmin_sales_item, + :description, CAST(:quantity AS NUMERIC), :unit, CAST(:unit_price AS NUMERIC), + :tax_schema, CAST(:total_amount AS NUMERIC), CAST(:payload AS JSONB) + ) + """), line_params) + else: + conn.execute(text("""INSERT INTO commercial_document_lines + (document_id,line_index,local_product_id,jasmin_sales_item,description,quantity,unit, + unit_price,tax_schema,total_amount,payload) + VALUES(CAST(:document_id AS UUID),:line_index,CAST(:local_product_id AS UUID), + :jasmin_sales_item,:description,CAST(:quantity AS NUMERIC),:unit, + CAST(:unit_price AS NUMERIC),:tax_schema,CAST(:total_amount AS NUMERIC),CAST(:payload AS JSONB))"""), + line_params) return document_id @@ -1722,6 +1734,16 @@ def _upsert_opportunity_items_from_jasmin_item(conn: Any, item: Dict[str, Any], if _clean(item.get("source_system")) != "jasmin" or not _clean(item.get("external_type")).startswith("jasmin_"): return 0 document_ref = _clean(item.get("document_number") or item.get("external_id")) + from app.document_reconciliation_service import resolve_document_links + external_id = _clean(item.get("external_id")) + document_number = _clean(item.get("document_number")) + matched = next((row for row in resolve_document_links(opportunity_id, conn=conn) + if external_id and row.get("external_id") == external_id + or document_number and row.get("document_number") == document_number), None) + origin = ({"commercial_document_id": matched.get("document_id"), + "opportunity_document_link_id": matched.get("id") + if matched.get("resolution_source") == "v2" else None} + if matched else {}) upserted = 0 for line in _jasmin_document_lines_from_item(item): mapping = _resolve_product_mapping_for_jasmin_line(conn, line) @@ -1746,6 +1768,8 @@ def _upsert_opportunity_items_from_jasmin_item(conn: Any, item: Dict[str, Any], "source_line_id": source_line_id, "source_external_id": item.get("external_id"), "source_external_type": item.get("external_type"), + "commercial_document_id": origin.get("commercial_document_id"), + "opportunity_document_link_id": origin.get("opportunity_document_link_id"), "resolved_sku": mapping.get("sku"), "resolved_jasmin_sales_item": mapping.get("jasmin_sales_item"), "product_mapping_status": mapping.get("mapping_status"), @@ -1992,16 +2016,10 @@ def _apply_reconstructed_process_to_opportunity(conn: Any, items: List[Dict[str, AND EXISTS (SELECT 1 FROM opportunity_items WHERE opportunity_id = CAST(:opportunity_id AS UUID)) """), {"opportunity_id": opportunity_id}) - has_invoice = bool(conn.execute(text(""" - SELECT EXISTS ( - SELECT 1 - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND document_kind = 'invoice' - AND COALESCE(is_active, TRUE) = TRUE - AND COALESCE(role, 'current') IN ('current', 'accepted', 'historical', 'history') - ) - """), {"opportunity_id": opportunity_id}).scalar()) + from app.document_reconciliation_service import resolve_document_links, select_valid_primary + # Only a valid effective PRIMARY invoice may drive SEND_INVOICE workflow. + # REMOVED/REASSIGNED/IGNORED/REVIEW_REQUIRED/HISTORICAL never qualify. + has_invoice = select_valid_primary(resolve_document_links(opportunity_id, conn=conn), "invoice") is not None if odoo_items and effective_action_code == "SEND_INVOICE" and not has_invoice: _ensure_pending_task_for_reconstruction( conn, diff --git a/app/reply_assistant_service.py b/app/reply_assistant_service.py index c081982..fab941a 100644 --- a/app/reply_assistant_service.py +++ b/app/reply_assistant_service.py @@ -743,7 +743,8 @@ def _select_default_documents(template: MessageTemplate, docs: Sequence[Dict[str return [] expected = {kind.lower() for kind in template.expected_document_kinds} matching = [doc for doc in docs if str(doc.get("document_kind") or "").lower() in expected] - current = [doc for doc in matching if str(doc.get("role") or "current") in {"current", "accepted"} and bool(doc.get("is_primary", True))] + current = [doc for doc in matching if str(doc.get("relationship") or "").upper() == "PRIMARY" + and str(doc.get("status") or "").upper() not in {"CANCELLED", "CANCELED"}] chosen = current or matching if not chosen: return [] diff --git a/app/revenue_forecast_service.py b/app/revenue_forecast_service.py index 8615346..cc1cfc3 100644 --- a/app/revenue_forecast_service.py +++ b/app/revenue_forecast_service.py @@ -321,16 +321,17 @@ def _historical_stage_rates(conn: Any) -> dict[str, dict[str, Any]]: def _realised_for_period(conn: Any, *, metric: str, period_start: date, period_end: date) -> dict[str, Any]: + from app.document_reconciliation_service import prepare_effective_document_links + prepare_effective_document_links(conn) if metric == "cash_received": rows = conn.execute(text(""" WITH latest_doc AS ( - SELECT DISTINCT ON (opportunity_id) opportunity_id, COALESCE(total_amount, amount, 0) AS amount - FROM commercial_documents - WHERE opportunity_id IS NOT NULL AND COALESCE(is_active, TRUE) = TRUE - ORDER BY opportunity_id, - CASE WHEN document_kind = 'invoice' THEN 0 ELSE 1 END, - COALESCE(document_date, created_at::date) DESC, - created_at DESC + SELECT DISTINCT ON (l.opportunity_id) l.opportunity_id, COALESCE(d.total_amount, d.amount, 0) AS amount + FROM _effective_document_links l JOIN commercial_documents d ON d.id=l.document_id + WHERE l.ended_at IS NULL AND l.relationship='PRIMARY' + AND upper(COALESCE(d.status,'')) NOT IN ('CANCELLED','CANCELED') + ORDER BY l.opportunity_id, CASE WHEN d.document_kind = 'invoice' THEN 0 ELSE 1 END, + COALESCE(d.document_date, d.created_at::date) DESC, d.created_at DESC ) SELECT ol.opportunity_id::text, COALESCE(CASE WHEN COALESCE(ol.payload->>'amount','') ~ '^[0-9]+([.,][0-9]+)?$' THEN replace(ol.payload->>'amount', ',', '.')::numeric END, ld.amount, o.value_amount, 0) AS amount, @@ -352,9 +353,10 @@ def _realised_for_period(conn: Any, *, metric: str, period_start: date, period_e COALESCE(o.customer_name, o.title, 'Venda ganha') AS reference FROM opportunities o LEFT JOIN LATERAL ( - SELECT COALESCE(total_amount, amount, 0) AS amount - FROM commercial_documents d - WHERE d.opportunity_id = o.id AND COALESCE(d.is_active, TRUE) = TRUE + SELECT COALESCE(d.total_amount, d.amount, 0) AS amount + FROM _effective_document_links l JOIN commercial_documents d ON d.id=l.document_id + WHERE l.opportunity_id = o.id AND l.ended_at IS NULL AND l.relationship='PRIMARY' + AND upper(COALESCE(d.status,'')) NOT IN ('CANCELLED','CANCELED') ORDER BY CASE WHEN d.document_kind = 'invoice' THEN 0 ELSE 1 END, COALESCE(d.document_date, d.created_at::date) DESC LIMIT 1 @@ -364,14 +366,14 @@ def _realised_for_period(conn: Any, *, metric: str, period_start: date, period_e """), {"period_start": period_start, "period_end": period_end}).mappings().all() else: rows = conn.execute(text(""" - SELECT d.opportunity_id::text, + SELECT l.opportunity_id::text, COALESCE(d.total_amount, d.amount, 0) AS amount, COALESCE(d.document_date::timestamp, d.created_at) AS realised_at, COALESCE(d.document_number, d.external_id, 'Fatura') AS reference - FROM commercial_documents d + FROM _effective_document_links l JOIN commercial_documents d ON d.id=l.document_id WHERE lower(d.document_kind) = 'invoice' - AND COALESCE(d.is_active, TRUE) = TRUE - AND COALESCE(d.role, 'current') IN ('current','accepted') + AND l.ended_at IS NULL AND l.relationship='PRIMARY' + AND upper(COALESCE(d.status,'')) NOT IN ('CANCELLED','CANCELED') AND COALESCE(d.document_date, d.created_at::date) BETWEEN :period_start AND :period_end """), {"period_start": period_start, "period_end": period_end}).mappings().all() items = [dict(r) for r in rows] @@ -384,6 +386,8 @@ def _realised_for_period(conn: Any, *, metric: str, period_start: date, period_e def _already_realised_ids(conn: Any, *, metric: str) -> set[str]: + from app.document_reconciliation_service import prepare_effective_document_links + prepare_effective_document_links(conn) if metric == "cash_received": sql = """ SELECT DISTINCT opportunity_id::text @@ -398,12 +402,11 @@ def _already_realised_ids(conn: Any, *, metric: str) -> set[str]: """ else: sql = """ - SELECT DISTINCT opportunity_id::text - FROM commercial_documents - WHERE opportunity_id IS NOT NULL - AND lower(document_kind) = 'invoice' - AND COALESCE(is_active, TRUE) = TRUE - AND COALESCE(role, 'current') IN ('current','accepted') + SELECT DISTINCT l.opportunity_id::text + FROM _effective_document_links l JOIN commercial_documents d ON d.id=l.document_id + WHERE l.ended_at IS NULL AND l.relationship='PRIMARY' + AND lower(d.document_kind) = 'invoice' + AND upper(COALESCE(d.status,'')) NOT IN ('CANCELLED','CANCELED') """ return {str(r[0]) for r in conn.execute(text(sql)).all() if r[0]} @@ -435,25 +438,21 @@ def get_revenue_forecast(*, limit: int = 1000, month: str | None = None, metric: target = get_sales_target(month=period_start, metric=metric) with engine.begin() as conn: + from app.document_reconciliation_service import prepare_effective_document_links + prepare_effective_document_links(conn) historical = _historical_stage_rates(conn) realised = _realised_for_period(conn, metric=metric, period_start=period_start, period_end=period_end) already_realised_ids = _already_realised_ids(conn, metric=metric) rows = conn.execute(text(""" WITH latest_doc AS ( - SELECT DISTINCT ON (opportunity_id) - opportunity_id, - total_amount, - document_number, - document_kind, - document_date - FROM commercial_documents - WHERE COALESCE(is_active, TRUE) = TRUE - AND COALESCE(role, 'current') IN ('current','accepted','historical','history') - ORDER BY opportunity_id, - CASE WHEN COALESCE(is_primary, FALSE) THEN 0 ELSE 1 END, - CASE document_kind WHEN 'invoice' THEN 1 WHEN 'quotation' THEN 2 WHEN 'proforma' THEN 3 ELSE 4 END, - COALESCE(document_date, created_at::date) DESC, - created_at DESC + SELECT DISTINCT ON (l.opportunity_id) + l.opportunity_id, d.total_amount, d.document_number, d.document_kind, d.document_date + FROM _effective_document_links l JOIN commercial_documents d ON d.id=l.document_id + WHERE l.ended_at IS NULL AND l.relationship='PRIMARY' + AND upper(COALESCE(d.status,'')) NOT IN ('CANCELLED','CANCELED') + ORDER BY l.opportunity_id, + CASE d.document_kind WHEN 'invoice' THEN 1 WHEN 'quotation' THEN 2 WHEN 'proforma' THEN 3 ELSE 4 END, + COALESCE(d.document_date, d.created_at::date) DESC, d.created_at DESC ), item_totals AS ( SELECT opportunity_id, SUM(COALESCE(total_price, 0)) AS total FROM opportunity_items diff --git a/app/task_service.py b/app/task_service.py index 076dd1a..a9ab557 100644 --- a/app/task_service.py +++ b/app/task_service.py @@ -624,26 +624,10 @@ def get_task_detail(task_id: str) -> Optional[Dict[str, Any]]: LEFT JOIN opportunities o ON o.id = t.opportunity_id LEFT JOIN LATERAL ( SELECT - cd.customer_id, - COALESCE(cdoc.name, NULLIF(cd.company, ''), NULLIF(cd.payload->>'customer_name', '')) AS customer_name, - COALESCE(cdoc.email, NULLIF(cd.payload->>'customer_email', ''), NULLIF(cd.payload->>'email', '')) AS customer_email, - COALESCE(cdoc.tax_id, NULLIF(cd.payload->>'customer_tax_id', ''), NULLIF(cd.payload->>'tax_id', ''), NULLIF(cd.payload->>'nif', '')) AS customer_tax_id, - cdoc.street_name AS customer_street_name, - cdoc.postal_zone AS customer_postal_zone, - cdoc.city_name AS customer_city_name, - cdoc.phone AS customer_phone - FROM commercial_documents cd - LEFT JOIN customers cdoc ON cdoc.id = cd.customer_id - WHERE cd.opportunity_id = o.id - AND COALESCE(cd.is_active, TRUE) = TRUE - AND cd.document_kind IN ('invoice', 'quotation', 'proforma') - ORDER BY - CASE cd.document_kind WHEN 'invoice' THEN 1 WHEN 'quotation' THEN 2 WHEN 'proforma' THEN 3 ELSE 4 END, - CASE COALESCE(cd.role, 'current') WHEN 'current' THEN 1 WHEN 'accepted' THEN 2 WHEN 'historical' THEN 3 WHEN 'history' THEN 3 ELSE 4 END, - COALESCE(cd.is_primary, FALSE) DESC, - COALESCE(cd.document_date, cd.created_at::date) DESC, - cd.created_at DESC - LIMIT 1 + NULL::uuid AS customer_id, NULL::text AS customer_name, + NULL::text AS customer_email, NULL::text AS customer_tax_id, + NULL::text AS customer_street_name, NULL::text AS customer_postal_zone, + NULL::text AS customer_city_name, NULL::text AS customer_phone ) dcu ON TRUE -- Legacy guard anchor: LEFT JOIN customers cu ON cu.id = o.local_customer_id -- v4928.1.5.64: choose schema-supported opportunity customer, then document customer fallback. @@ -654,7 +638,29 @@ def get_task_detail(task_id: str) -> Optional[Dict[str, Any]]: with engine.begin() as conn: row = conn.execute(sql, {"task_id": task_id}).mappings().first() - return dict(row) if row else None + result = dict(row) if row else None + if result and result.get("opportunity_id") and not result.get("linked_customer_id"): + from app.document_reconciliation_service import resolve_document_links, select_valid_primary + links = resolve_document_links(str(result["opportunity_id"])) + document = next((select_valid_primary(links, kind) + for kind in ("invoice", "quotation", "proforma") + if select_valid_primary(links, kind)), None) + if document and document.get("customer_id"): + with engine.begin() as conn: + customer = conn.execute(text("""SELECT id::text, name, email, tax_id, street_name, + postal_zone, city_name, phone FROM customers WHERE id=CAST(:id AS UUID)"""), + {"id": document["customer_id"]}).mappings().first() + if customer: + customer = dict(customer) + result["linked_customer_id"] = customer["id"] + result["opportunity_fiscal_customer_id"] = customer["id"] + for source_key, target_key in (("name", "linked_customer_name"), + ("email", "linked_customer_email"), ("tax_id", "linked_customer_tax_id"), + ("street_name", "linked_customer_street_name"), + ("postal_zone", "linked_customer_postal_zone"), + ("city_name", "linked_customer_city_name"), ("phone", "linked_customer_phone")): + result[target_key] = result.get(target_key) or customer.get(source_key) + return result def list_customer_task_history( diff --git a/app/workflow_guard.py b/app/workflow_guard.py index 9cf66cd..9028666 100644 --- a/app/workflow_guard.py +++ b/app/workflow_guard.py @@ -141,35 +141,23 @@ def get_workflow_context(opportunity_id: str) -> Dict[str, Any]: WHERE opportunity_id = CAST(:opportunity_id AS UUID) """), {"opportunity_id": str(opportunity_id)}).scalar() or 0 - jasmin_doc_counts = conn.execute(text(""" - SELECT - count(*) FILTER (WHERE document_kind = 'quotation') AS quotation_count, - count(*) FILTER (WHERE document_kind = 'proforma') AS proforma_count, - count(*) FILTER (WHERE document_kind = 'invoice') AS invoice_count, - bool_or( - document_kind = 'invoice' - AND COALESCE(is_active, TRUE) IS TRUE - AND COALESCE(role, 'current') IN ('current','accepted') - AND ( - status IN ('sent','issued_sent') - OR COALESCE(payload, '{}'::jsonb) ? 'clientflow_invoice_sent_evidence' - OR COALESCE(payload, '{}'::jsonb) ? 'invoice_sent_at' - ) - ) AS invoice_sent_evidence - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND system = 'jasmin' - """), {"opportunity_id": str(opportunity_id)}).mappings().first() - - - current_invoices = conn.execute(text(""" - SELECT id::text, document_number, external_id, status, payload - FROM commercial_documents - WHERE opportunity_id = CAST(:opportunity_id AS UUID) - AND document_kind = 'invoice' - AND COALESCE(is_active, TRUE) IS TRUE - AND COALESCE(role, 'current') IN ('current','accepted','actual','active','') - """), {"opportunity_id": str(opportunity_id)}).mappings().all() + from app.document_reconciliation_service import resolve_document_links, select_valid_primary + resolved_documents = [row for row in resolve_document_links(str(opportunity_id), conn=conn) + if row.get("system") == "jasmin"] + primaries = [row for row in resolved_documents + if row.get("relationship") == "PRIMARY" and + str(row.get("jasmin_status") or row.get("status") or "").upper() not in {"CANCELLED", "CANCELED"}] + jasmin_doc_counts = { + "quotation_count": sum(row.get("document_kind") == "quotation" for row in primaries), + "proforma_count": sum(row.get("document_kind") == "proforma" for row in primaries), + "invoice_count": sum(row.get("document_kind") == "invoice" for row in primaries), + "invoice_sent_evidence": any(row.get("document_kind") == "invoice" and ( + row.get("status") in {"sent", "issued_sent"} or + "clientflow_invoice_sent_evidence" in (row.get("payload") or {}) or + "invoice_sent_at" in (row.get("payload") or {})) for row in primaries), + } + current_primary_invoice = select_valid_primary(resolved_documents, "invoice") + current_invoices = [current_primary_invoice] if current_primary_invoice else [] invoice_tasks = conn.execute(text(""" SELECT id::text, action_code, action, note, status, metadata diff --git a/docs/document_reconciliation_v2.md b/docs/document_reconciliation_v2.md new file mode 100644 index 0000000..e535d05 --- /dev/null +++ b/docs/document_reconciliation_v2.md @@ -0,0 +1,39 @@ +# Document Reconciliation v2 + +After a group is completely classified, `opportunity_document_links` is the +sole authority for document selection. During rollout the canonical resolver +uses legacy fields for a whole `(opportunity_id, document_kind)` group until +every legacy document in that group has a current v2 link. It never mixes both +authorities inside a group. `commercial_documents.opportunity_id`, `role`, +`is_primary` and `is_active` are dual-write compatibility fields only. + +## Deployment + +1. Deploy this compatibility code first. It can run before migration 007. +2. Run `python scripts/preflight_document_reconciliation_v2.py`; stop on any blocker. +3. Run `python scripts/apply_migrations.py --dry-run` and review the pending migration. +4. Run `python scripts/apply_migrations.py` against the intended database. +5. Keep all consumers on this compatible release. Do not deploy v2-exclusive readers. +6. Dry-run group batches and archive JSON/CSV: `python scripts/migrate_document_reconciliation_v2.py --dry-run --batch-size 500 --output-json /safe/path/batch.json`. +7. Apply complete group batches. Resume with the emitted composite checkpoint, + `--resume-from 'OPPORTUNITY_UUID|document_kind'`. A group is checkpointed only + after all its documents commit. +8. `--only-unambiguous` still writes `REVIEW_REQUIRED` links (option B), so groups + never become partially invisible. Resolve reviews in the admin UI. +9. Verify no legacy-resolved groups remain, then v2-exclusive readers may be deployed. + +## Rollback + +Rollback application code first and stop every v2 consumer. Export +`opportunity_document_link_events`; confirm all legacy compatibility fields are +filled. In a dedicated psql session set +`clientflow.v2_consumers_active='off'` and +`clientflow.document_ledger_exported='on'`, then execute +`migrations/007_document_reconciliation_v2_down.sql`. The down migration checks +these conditions and removes dependent line columns, trigger/function, events, +command claims and links in dependency order. The ledger is lost after export; +manual decisions made only in v2 must be reconciled before rollback. Neither +migration changes remote Jasmin documents. + +CSRF remains separate security debt; this change reuses existing admin +authentication and does not claim to add CSRF protection. diff --git a/migrations/007_document_reconciliation_v2.sql b/migrations/007_document_reconciliation_v2.sql new file mode 100644 index 0000000..d069238 --- /dev/null +++ b/migrations/007_document_reconciliation_v2.sql @@ -0,0 +1,100 @@ +-- Document Reconciliation v2 (additive; no data backfill is run here). +CREATE EXTENSION IF NOT EXISTS pgcrypto; + +CREATE TABLE IF NOT EXISTS opportunity_document_links ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE RESTRICT, + document_id UUID NOT NULL REFERENCES commercial_documents(id) ON DELETE RESTRICT, + document_kind TEXT NOT NULL, + relationship TEXT NOT NULL CHECK (relationship IN ('PRIMARY','SECONDARY','HISTORICAL','IGNORED','REMOVED','REASSIGNED','REVIEW_REQUIRED')), + is_manual BOOLEAN NOT NULL DEFAULT FALSE, + decision_reason TEXT, + decided_by TEXT, + decided_at TIMESTAMPTZ, + source TEXT NOT NULL, + origin_opportunity_id UUID REFERENCES opportunities(id) ON DELETE RESTRICT, + destination_opportunity_id UUID REFERENCES opportunities(id) ON DELETE RESTRICT, + correlation_id TEXT, + request_id TEXT, + restored_from_link_id UUID REFERENCES opportunity_document_links(id) ON DELETE SET NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + ended_at TIMESTAMPTZ, + metadata JSONB NOT NULL DEFAULT '{}'::jsonb, + version INTEGER NOT NULL DEFAULT 1 CHECK (version > 0), + CONSTRAINT ck_document_link_reassigned CHECK ( + relationship <> 'REASSIGNED' OR ( + origin_opportunity_id IS NOT NULL AND destination_opportunity_id IS NOT NULL + AND origin_opportunity_id <> destination_opportunity_id + ) + ) +); + +CREATE UNIQUE INDEX IF NOT EXISTS ux_document_link_current +ON opportunity_document_links(opportunity_id, document_id) WHERE ended_at IS NULL; +CREATE UNIQUE INDEX IF NOT EXISTS ux_document_link_primary_kind +ON opportunity_document_links(opportunity_id, document_kind) WHERE ended_at IS NULL AND relationship = 'PRIMARY'; +CREATE INDEX IF NOT EXISTS idx_document_links_opportunity ON opportunity_document_links(opportunity_id, ended_at, relationship); +CREATE INDEX IF NOT EXISTS idx_document_links_document ON opportunity_document_links(document_id, ended_at); +CREATE INDEX IF NOT EXISTS idx_document_links_relationship ON opportunity_document_links(relationship, updated_at DESC); +CREATE INDEX IF NOT EXISTS idx_document_links_destination ON opportunity_document_links(destination_opportunity_id) WHERE destination_opportunity_id IS NOT NULL; + +CREATE TABLE IF NOT EXISTS opportunity_document_link_events ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + link_id UUID REFERENCES opportunity_document_links(id) ON DELETE SET NULL, + opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE RESTRICT, + document_id UUID NOT NULL REFERENCES commercial_documents(id) ON DELETE RESTRICT, + event_type TEXT NOT NULL, + actor TEXT NOT NULL, + reason TEXT, + old_relationship TEXT, + new_relationship TEXT, + old_opportunity_id UUID, + new_opportunity_id UUID, + old_jasmin_status TEXT, + new_jasmin_status TEXT, + correlation_id TEXT, + request_id TEXT, + idempotency_key TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + payload JSONB NOT NULL DEFAULT '{}'::jsonb +); +CREATE UNIQUE INDEX IF NOT EXISTS ux_document_link_events_idempotency +ON opportunity_document_link_events(idempotency_key) WHERE idempotency_key IS NOT NULL; +CREATE INDEX IF NOT EXISTS idx_document_link_events_link ON opportunity_document_link_events(link_id, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_document_link_events_opportunity ON opportunity_document_link_events(opportunity_id, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_document_link_events_document ON opportunity_document_link_events(document_id, created_at DESC); + +-- Command claims make the mutation itself idempotent, not merely its event. +CREATE TABLE IF NOT EXISTS document_reconciliation_commands ( + idempotency_key TEXT PRIMARY KEY, + fingerprint TEXT NOT NULL, + command_payload JSONB NOT NULL, + result_payload JSONB, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + completed_at TIMESTAMPTZ +); + +-- Events are an audit ledger: UPDATE and DELETE are rejected at database level. +CREATE OR REPLACE FUNCTION reject_document_link_event_mutation() RETURNS trigger +LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'opportunity_document_link_events is append-only'; END $$; +DROP TRIGGER IF EXISTS trg_document_link_events_append_only ON opportunity_document_link_events; +CREATE TRIGGER trg_document_link_events_append_only BEFORE UPDATE OR DELETE ON opportunity_document_link_events +FOR EACH ROW EXECUTE FUNCTION reject_document_link_event_mutation(); + +ALTER TABLE commercial_document_lines ADD COLUMN IF NOT EXISTS commercial_document_id UUID REFERENCES commercial_documents(id) ON DELETE RESTRICT; +UPDATE commercial_document_lines SET commercial_document_id = document_id WHERE commercial_document_id IS NULL; +DO $$ BEGIN + IF NOT EXISTS (SELECT 1 FROM pg_constraint WHERE conname='ck_commercial_document_lines_document_present') THEN + ALTER TABLE commercial_document_lines ADD CONSTRAINT ck_commercial_document_lines_document_present + CHECK (commercial_document_id IS NOT NULL) NOT VALID; + END IF; +END $$; +ALTER TABLE commercial_document_lines VALIDATE CONSTRAINT ck_commercial_document_lines_document_present; +ALTER TABLE commercial_document_lines ALTER COLUMN commercial_document_id SET NOT NULL; +ALTER TABLE commercial_document_lines ADD COLUMN IF NOT EXISTS opportunity_document_link_id UUID REFERENCES opportunity_document_links(id) ON DELETE SET NULL; +CREATE INDEX IF NOT EXISTS idx_document_lines_commercial_document ON commercial_document_lines(commercial_document_id); +CREATE INDEX IF NOT EXISTS idx_document_lines_opportunity_link ON commercial_document_lines(opportunity_document_link_id); + +COMMENT ON TABLE opportunity_document_links IS 'Canonical ClientFlow document/opportunity relationship; legacy commercial_documents fields are not authoritative.'; +COMMENT ON TABLE opportunity_document_link_events IS 'Append-only reconciliation audit ledger.'; diff --git a/migrations/007_document_reconciliation_v2_down.sql b/migrations/007_document_reconciliation_v2_down.sql new file mode 100644 index 0000000..2915f5b --- /dev/null +++ b/migrations/007_document_reconciliation_v2_down.sql @@ -0,0 +1,37 @@ +-- EXECUTE ONLY AFTER reverting application code and stopping every v2 consumer. +-- Before running, export opportunity_document_link_events and set these session flags: +-- SET clientflow.v2_consumers_active = 'off'; +-- SET clientflow.document_ledger_exported = 'on'; +BEGIN; + +DO $$ +BEGIN + IF current_setting('clientflow.v2_consumers_active', true) IS DISTINCT FROM 'off' THEN + RAISE EXCEPTION 'rollback refused: confirm all v2 consumers are stopped'; + END IF; + IF current_setting('clientflow.document_ledger_exported', true) IS DISTINCT FROM 'on' THEN + RAISE EXCEPTION 'rollback refused: export the reconciliation ledger first'; + END IF; + IF EXISTS (SELECT 1 FROM commercial_documents + WHERE opportunity_id IS NULL OR role IS NULL OR is_primary IS NULL OR is_active IS NULL) THEN + RAISE EXCEPTION 'rollback refused: legacy compatibility fields are incomplete'; + END IF; +END $$; + +-- Remove dependent line foreign keys/columns before their target tables. +ALTER TABLE commercial_document_lines DROP CONSTRAINT IF EXISTS commercial_document_lines_opportunity_document_link_id_fkey; +DROP INDEX IF EXISTS idx_document_lines_opportunity_link; +ALTER TABLE commercial_document_lines DROP COLUMN IF EXISTS opportunity_document_link_id; +ALTER TABLE commercial_document_lines DROP CONSTRAINT IF EXISTS ck_commercial_document_lines_document_present; +ALTER TABLE commercial_document_lines DROP CONSTRAINT IF EXISTS commercial_document_lines_commercial_document_id_fkey; +DROP INDEX IF EXISTS idx_document_lines_commercial_document; +ALTER TABLE commercial_document_lines DROP COLUMN IF EXISTS commercial_document_id; + +DROP TRIGGER IF EXISTS trg_document_link_events_append_only ON opportunity_document_link_events; +DROP FUNCTION IF EXISTS reject_document_link_event_mutation(); +DROP TABLE IF EXISTS opportunity_document_link_events; +DROP TABLE IF EXISTS document_reconciliation_commands; +DROP TABLE IF EXISTS opportunity_document_links; +DELETE FROM schema_migrations WHERE version='007'; + +COMMIT; diff --git a/scripts/apply_migrations.py b/scripts/apply_migrations.py index 28aeb82..4d07b65 100755 --- a/scripts/apply_migrations.py +++ b/scripts/apply_migrations.py @@ -39,7 +39,8 @@ def main() -> int: parser.add_argument("--dry-run", action="store_true", help="List pending migrations without applying them.") args = parser.parse_args() - files = sorted(MIGRATIONS_DIR.glob("*.sql")) + # Down migrations are operator-invoked rollback artefacts, never pending ups. + files = sorted(path for path in MIGRATIONS_DIR.glob("*.sql") if not path.stem.endswith("_down")) if not files: print("No migrations found.") return 0 @@ -69,6 +70,14 @@ def main() -> int: return 0 for version, path in pending: + if version == "007": + from scripts.preflight_document_reconciliation_v2 import run_preflight + diagnostics = run_preflight(conn) + if diagnostics["blockers"]: + print("Migration 007 preflight failed:") + for blocker in diagnostics["blockers"]: + print(f"- {blocker}") + return 2 sql = path.read_text() print(f"Applying {path.name}...") conn.execute(text(sql)) diff --git a/scripts/migrate_document_reconciliation_v2.py b/scripts/migrate_document_reconciliation_v2.py new file mode 100644 index 0000000..152e2d9 --- /dev/null +++ b/scripts/migrate_document_reconciliation_v2.py @@ -0,0 +1,137 @@ +#!/usr/bin/env python3 +"""Classify legacy document links into Document Reconciliation v2. + +Dry-run is the default. This script changes ClientFlow only and never calls +Jasmin. Ambiguous evidence is explicitly classified REVIEW_REQUIRED. +""" +from __future__ import annotations + +import argparse +import csv +import json +import os +import sys +from pathlib import Path +from typing import Any, Dict + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + +from app.document_reconciliation_backfill import build_backfill_result + + +def load_documents( + opportunity_id: str | None = None, + resume_from: str | None = None, + batch_size: int = 500, +) -> tuple[list[Dict[str, Any]], int, int]: + from sqlalchemy import text + + from app.db import engine + + where, params = ["d.opportunity_id IS NOT NULL"], {"limit": batch_size} + if opportunity_id: + where.append("d.opportunity_id=CAST(:opportunity_id AS UUID)"); params["opportunity_id"] = opportunity_id + if resume_from: + try: + resume_opportunity_id, resume_document_kind = resume_from.split("|", 1) + except ValueError as exc: + raise ValueError("--resume-from must be OPPORTUNITY_ID|DOCUMENT_KIND") from exc + where.append("(d.opportunity_id::text, d.document_kind) > (:resume_opportunity_id, :resume_document_kind)") + params.update(resume_opportunity_id=resume_opportunity_id, resume_document_kind=resume_document_kind) + with engine.begin() as conn: + rows = [dict(r) for r in conn.execute(text(f""" + WITH selected_groups AS ( + SELECT d.opportunity_id, d.document_kind + FROM commercial_documents d WHERE {' AND '.join(where)} + GROUP BY d.opportunity_id, d.document_kind + ORDER BY d.opportunity_id, d.document_kind + LIMIT :limit + ) + SELECT d.id::text, d.opportunity_id::text, d.document_kind, d.external_id, + d.document_number, d.status, d.role, d.is_primary, d.is_active, + d.total_amount, d.amount, d.payload + FROM selected_groups g + JOIN commercial_documents d ON d.opportunity_id=g.opportunity_id + AND d.document_kind=g.document_kind + ORDER BY d.opportunity_id, d.document_kind, d.created_at, d.id + """), params).mappings().all()] + orphan_lines = int(conn.execute(text("""SELECT count(*) FROM commercial_document_lines dl + LEFT JOIN commercial_documents d ON d.id=COALESCE(dl.commercial_document_id,dl.document_id) + WHERE d.id IS NULL""")).scalar() or 0) + multi = int(conn.execute(text("""SELECT count(*) FROM (SELECT COALESCE(external_id,document_number), count(DISTINCT opportunity_id) + FROM commercial_documents WHERE opportunity_id IS NOT NULL AND COALESCE(external_id,document_number) IS NOT NULL + GROUP BY 1 HAVING count(DISTINCT opportunity_id)>1) q""")).scalar() or 0) + return rows, orphan_lines, multi + + +def scan(opportunity_id: str | None = None, resume_from: str | None = None, batch_size: int = 500) -> Dict[str, Any]: + rows, orphan_lines, multi = load_documents(opportunity_id, resume_from, batch_size) + result = build_backfill_result( + rows, + lines_without_document=orphan_lines, + documents_in_multiple_opportunities=multi, + ) + groups = sorted({(str(row["opportunity_id"]), str(row["document_kind"])) for row in rows}) + result["summary"]["batch_groups"] = len(groups) + result["summary"]["checkpoint"] = "|".join(groups[-1]) if groups else resume_from + return result + + +def actions_for_only_unambiguous(actions: list[Dict[str, Any]]) -> list[Dict[str, Any]]: + """Keep ambiguous documents represented; never partially migrate a group. + + Option B of the rollout contract is used: REVIEW_REQUIRED links are written + alongside unambiguous links so every selected group becomes v2-complete. + """ + return list(actions) + + +def apply_actions(actions: list[Dict[str, Any]]) -> int: + from app.db import engine + from app.document_reconciliation_service import set_document_relationship + groups: Dict[tuple[str, str], list[Dict[str, Any]]] = {} + for action in actions: + groups.setdefault((str(action["opportunity_id"]), str(action["document_kind"])), []).append(action) + # A group is the rollout authority boundary. Links, events and legacy dual + # writes commit together; any failure rolls the entire group back. + for group in groups.values(): + with engine.begin() as conn: + for action in group: + set_document_relationship(action["opportunity_id"], action["id"], action["relationship"], + actor="document_reconciliation_v2_backfill", reason=action["decision_reason"], is_manual=False, + source="backfill_v2", idempotency_key=f"backfill-v2:{action['opportunity_id']}:{action['id']}", + event_type="BACKFILLED", metadata={"legacy_role": action.get("role"), "legacy_is_primary": action.get("is_primary")}, + _conn=conn) + return len(actions) + + +def write_reports(result: Dict[str, Any], json_path: str | None, csv_path: str | None) -> None: + if json_path: Path(json_path).write_text(json.dumps(result, ensure_ascii=False, indent=2, default=str) + "\n") + if csv_path: + fields = ["opportunity_id", "id", "document_kind", "document_number", "status", "relationship", "decision_reason"] + with Path(csv_path).open("w", newline="", encoding="utf-8") as handle: + writer = csv.DictWriter(handle, fieldnames=fields, extrasaction="ignore"); writer.writeheader(); writer.writerows(result["actions"]) + + +def main() -> int: + os.chdir(ROOT) + parser = argparse.ArgumentParser() + mode = parser.add_mutually_exclusive_group(); mode.add_argument("--dry-run", action="store_true"); mode.add_argument("--apply", action="store_true") + parser.add_argument("--only-unambiguous", action="store_true") + parser.add_argument("--opportunity-id"); parser.add_argument("--batch-size", type=int, default=500) + parser.add_argument("--resume-from"); parser.add_argument("--output-json"); parser.add_argument("--output-csv") + args = parser.parse_args() + result = scan(args.opportunity_id, args.resume_from, args.batch_size) + actions = result["actions"] + if args.only_unambiguous: actions = actions_for_only_unambiguous(actions) + applied = 0 + if args.apply: + applied = apply_actions(actions) + result["summary"]["mode"] = "apply" if args.apply else "dry-run"; result["summary"]["applied"] = applied + write_reports(result, args.output_json, args.output_csv) + print(json.dumps(result["summary"], ensure_ascii=False, indent=2, default=str)) + return 0 + + +if __name__ == "__main__": raise SystemExit(main()) diff --git a/scripts/preflight_document_reconciliation_v2.py b/scripts/preflight_document_reconciliation_v2.py new file mode 100644 index 0000000..733126d --- /dev/null +++ b/scripts/preflight_document_reconciliation_v2.py @@ -0,0 +1,107 @@ +#!/usr/bin/env python3 +"""Read-only preflight for migration 007; it never repairs production data.""" +from __future__ import annotations + +import json +import sys +from pathlib import Path +from typing import Any, Dict + +from sqlalchemy import text + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + + +def run_preflight(conn: Any) -> Dict[str, Any]: + column_rows = list(conn.execute(text(""" + SELECT table_name, column_name, data_type, udt_name FROM information_schema.columns + WHERE table_schema='public' AND table_name IN + ('commercial_document_lines','commercial_documents','opportunities') + """))) + columns = {row[1] for row in column_rows if row[0] == "commercial_document_lines"} + required = {"id", "document_id"} + counts: Dict[str, int] = {} + details: Dict[str, list[Dict[str, Any]]] = {} + blockers = [] + missing = sorted(required - columns) + if missing: + blockers.append("commercial_document_lines missing columns: " + ", ".join(missing)) + return {"counts": counts, "blockers": blockers, "columns": sorted(columns)} + types = {(row[0], row[1]): row[3] for row in column_rows} + for table_name, column_name in (("commercial_document_lines", "document_id"), + ("commercial_documents", "id"), + ("commercial_documents", "opportunity_id"), + ("opportunities", "id")): + actual = types.get((table_name, column_name)) + if actual != "uuid": + blockers.append(f"unexpected type {table_name}.{column_name}: {actual or 'missing'} (expected uuid)") + counts["null_document_id"] = int(conn.execute(text( + "SELECT count(*) FROM commercial_document_lines WHERE document_id IS NULL" + )).scalar() or 0) + counts["orphan_document_id"] = int(conn.execute(text(""" + SELECT count(*) FROM commercial_document_lines dl + LEFT JOIN commercial_documents d ON d.id=dl.document_id + WHERE dl.document_id IS NOT NULL AND d.id IS NULL + """)).scalar() or 0) + duplicate_documents = [dict(row) for row in conn.execute(text(""" + SELECT system, external_id, document_kind, company, + array_agg(id::text ORDER BY id) AS ids, + 'duplicate_external_document' AS reason + FROM commercial_documents + WHERE external_id IS NOT NULL AND btrim(external_id) <> '' + GROUP BY system, external_id, document_kind, company HAVING count(*)>1 + """)).mappings().all()] + duplicate_associations = [dict(row) for row in conn.execute(text(""" + SELECT opportunity_id::text, system, external_id, document_kind, + array_agg(id::text ORDER BY id) AS ids, + 'duplicate_opportunity_document_association' AS reason + FROM commercial_documents + WHERE opportunity_id IS NOT NULL AND external_id IS NOT NULL AND btrim(external_id) <> '' + GROUP BY opportunity_id, system, external_id, document_kind HAVING count(*)>1 + """)).mappings().all()] + cross_opportunity = [dict(row) for row in conn.execute(text(""" + SELECT system, external_id, document_kind, company, + array_agg(id::text ORDER BY id) AS ids, + array_agg(DISTINCT opportunity_id::text) AS opportunity_ids, + 'external_document_linked_to_multiple_opportunities' AS reason + FROM commercial_documents + WHERE opportunity_id IS NOT NULL AND external_id IS NOT NULL AND btrim(external_id) <> '' + GROUP BY system, external_id, document_kind, company + HAVING count(DISTINCT opportunity_id)>1 + """)).mappings().all()] + multiple_primaries = [dict(row) for row in conn.execute(text(""" + SELECT opportunity_id::text, document_kind, array_agg(id::text ORDER BY id) AS ids, + 'multiple_legacy_primaries' AS reason + FROM commercial_documents + WHERE opportunity_id IS NOT NULL AND COALESCE(is_primary,FALSE) + AND COALESCE(role,'current') IN ('current','accepted') + GROUP BY opportunity_id, document_kind HAVING count(*)>1 + """)).mappings().all()] + details.update(duplicate_documents=duplicate_documents, + duplicate_associations=duplicate_associations, + cross_opportunity_documents=cross_opportunity, + multiple_legacy_primaries=multiple_primaries) + counts["duplicate_document_groups"] = len(duplicate_documents) + counts["duplicate_associations"] = len(duplicate_associations) + counts["cross_opportunity_documents"] = len(cross_opportunity) + counts["multiple_legacy_primaries"] = len(multiple_primaries) + for key in ("null_document_id", "orphan_document_id", "duplicate_document_groups", + "duplicate_associations", "cross_opportunity_documents"): + if counts[key]: + blockers.append(f"{key}: {counts[key]}") + if counts["multiple_legacy_primaries"]: + blockers.append(f"multiple_legacy_primaries: {counts['multiple_legacy_primaries']}") + return {"counts": counts, "blockers": blockers, "details": details, "columns": sorted(columns)} + + +def main() -> int: + from app.db import engine + with engine.begin() as conn: + result = run_preflight(conn) + print(json.dumps(result, ensure_ascii=False, indent=2)) + return 2 if result["blockers"] else 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_document_reconciliation_v2.py b/tests/test_document_reconciliation_v2.py new file mode 100644 index 0000000..d00601f --- /dev/null +++ b/tests/test_document_reconciliation_v2.py @@ -0,0 +1,370 @@ +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 diff --git a/tests/test_document_reconciliation_v2_postgres.py b/tests/test_document_reconciliation_v2_postgres.py new file mode 100644 index 0000000..e38da91 --- /dev/null +++ b/tests/test_document_reconciliation_v2_postgres.py @@ -0,0 +1,104 @@ +"""Optional PostgreSQL contract tests. + +These tests are intentionally gated by a dedicated disposable-test URL. They +never fall back to DATABASE_URL. +""" +from __future__ import annotations + +import os +from pathlib import Path +import uuid +from contextlib import nullcontext + +import pytest +from sqlalchemy import create_engine, text + + +TEST_URL = os.getenv("CLIENTFLOW_TEST_DATABASE_URL") +pytestmark = pytest.mark.skipif(not TEST_URL, reason="CLIENTFLOW_TEST_DATABASE_URL not set") + + +@pytest.fixture() +def pg_conn(): + engine = create_engine(TEST_URL) + schema = "document_reconciliation_test_" + uuid.uuid4().hex + with engine.connect() as conn: + conn.execute(text(f'CREATE SCHEMA "{schema}"')) + conn.commit() + conn.execute(text(f'SET search_path TO "{schema}"')) + conn.execute(text("""CREATE TABLE opportunities(id UUID PRIMARY KEY); + CREATE TABLE commercial_documents( + id UUID PRIMARY KEY, opportunity_id UUID REFERENCES opportunities(id), + document_kind TEXT NOT NULL, role TEXT NOT NULL DEFAULT 'related', + is_primary BOOLEAN NOT NULL DEFAULT FALSE, is_active BOOLEAN NOT NULL DEFAULT TRUE, + updated_at TIMESTAMPTZ DEFAULT now()); + CREATE TABLE commercial_document_lines( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), document_id UUID REFERENCES commercial_documents(id), + opportunity_item_id UUID, line_index INTEGER, local_product_id UUID, + jasmin_sales_item TEXT, description TEXT, quantity NUMERIC, unit TEXT, + unit_price NUMERIC, tax_schema TEXT, total_amount NUMERIC, payload JSONB, + created_at TIMESTAMPTZ DEFAULT now()); + CREATE TABLE schema_migrations(version TEXT PRIMARY KEY);""")) + conn.commit() + try: + yield conn + finally: + conn.rollback() + conn.execute(text("SET search_path TO public")) + conn.execute(text(f'DROP SCHEMA "{schema}" CASCADE')) + conn.commit() + engine.dispose() + + +def test_dedicated_url_is_the_only_postgres_test_source(pg_conn): + assert TEST_URL + assert str(pg_conn.execute(text("SELECT current_database()" )).scalar()) + + +def test_migration_007_postgres_contracts_and_generic_writer(pg_conn, monkeypatch): + sql = Path("migrations/007_document_reconciliation_v2.sql").read_text() + pg_conn.exec_driver_sql(sql) + pg_conn.commit() + pg_conn.exec_driver_sql("SET search_path TO " + pg_conn.exec_driver_sql("SELECT current_schema()").scalar()) + oid, first, second = [str(uuid.uuid4()) for _ in range(3)] + pg_conn.execute(text("INSERT INTO opportunities(id) VALUES(CAST(:id AS UUID))"), {"id": oid}) + pg_conn.execute(text("""INSERT INTO commercial_documents(id,opportunity_id,document_kind) + VALUES(CAST(:a AS UUID),CAST(:o AS UUID),'invoice'), + (CAST(:b AS UUID),CAST(:o AS UUID),'invoice')"""), {"a": first, "b": second, "o": oid}) + pg_conn.execute(text("""INSERT INTO opportunity_document_links + (opportunity_id,document_id,document_kind,relationship,source) + VALUES(CAST(:o AS UUID),CAST(:d AS UUID),'invoice','PRIMARY','test')"""), {"o": oid, "d": first}) + pg_conn.commit() + with pytest.raises(Exception): + pg_conn.execute(text("""INSERT INTO opportunity_document_links + (opportunity_id,document_id,document_kind,relationship,source) + VALUES(CAST(:o AS UUID),CAST(:d AS UUID),'invoice','PRIMARY','test')"""), {"o": oid, "d": second}) + pg_conn.rollback() + + import app.commercial_service as commercial + class BoundEngine: + def begin(self): return nullcontext(pg_conn) + monkeypatch.setattr(commercial, "engine", BoundEngine()) + monkeypatch.setattr(commercial, "ensure_commercial_schema", lambda: None) + commercial.add_document_lines(first, [{"description": "line"}]) + ids = pg_conn.execute(text("""SELECT document_id::text, commercial_document_id::text + FROM commercial_document_lines""")).first() + assert ids == (first, first) + + +def test_migration_007_append_only_and_down(pg_conn): + pg_conn.exec_driver_sql(Path("migrations/007_document_reconciliation_v2.sql").read_text()) + oid, did = str(uuid.uuid4()), str(uuid.uuid4()) + pg_conn.execute(text("INSERT INTO opportunities(id) VALUES(CAST(:id AS UUID))"), {"id": oid}) + pg_conn.execute(text("""INSERT INTO commercial_documents(id,opportunity_id,document_kind) + VALUES(CAST(:d AS UUID),CAST(:o AS UUID),'invoice')"""), {"d": did, "o": oid}) + event = pg_conn.execute(text("""INSERT INTO opportunity_document_link_events + (opportunity_id,document_id,event_type,actor) VALUES(CAST(:o AS UUID),CAST(:d AS UUID),'TEST','test') + RETURNING id::text"""), {"o": oid, "d": did}).scalar() + with pytest.raises(Exception): + pg_conn.execute(text("UPDATE opportunity_document_link_events SET actor='changed' WHERE id=CAST(:id AS UUID)"), {"id": event}) + pg_conn.rollback() + pg_conn.execute(text("SET clientflow.v2_consumers_active='off'")) + pg_conn.execute(text("SET clientflow.document_ledger_exported='on'")) + pg_conn.exec_driver_sql(Path("migrations/007_document_reconciliation_v2_down.sql").read_text()) + assert pg_conn.execute(text("SELECT to_regclass('opportunity_document_links')")).scalar() is None