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
| Data | Evento | Operador | Relação | Motivo |
{rows}
')
+
+
@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