#!/usr/bin/env python3 """Read-only BLIF Flow v2 projection against the isolated shadow snapshot.""" from __future__ import annotations import json import re from collections import Counter, defaultdict from dataclasses import replace from datetime import datetime, timezone from pathlib import Path from typing import Any, Iterable from sqlalchemy import text from app.db import engine from app.domain.opportunity_flow.v2 import ( EffectiveOperationalDecision, derive_business_facts, derive_effective_operational_action, derive_safe_operational_action, derive_v2_operational_queue, suppress_duplicate_representation, ) from app.operations_service import get_operations_summary from app.opportunity_next_action_service import get_opportunity_next_actions from app.opportunity_service import list_opportunities PROJECTION = Path("/tmp/blif_flow_v2_projection.json") COMPARISON = Path("/tmp/blif_flow_v2_operations_comparison.txt") AMBIGUOUS = Path("/tmp/blif_flow_v2_ambiguous_cases.json") PROMOTIONS = Path("/tmp/blif_flow_v2_promotions_audit.json") INVOICE_WITHOUT_PAYMENT = Path("/tmp/blif_flow_v2_invoice_without_payment.json") REVIEW_AUDIT = Path("/tmp/blif_flow_v2_review_audit.json") BACKLOG_DELTA = Path("/tmp/blif_flow_v2_backlog_delta.json") CURRENT_DELTA = Path("/tmp/blif_flow_v2_current_delta.json") MATERIAL_IDENTITY = Path("/tmp/blif_flow_v2_material_identity.json") TERMINAL = {"WON", "LOST", "NO_INTEREST", "ARCHIVED", "COMPLETED", "CLOSED"} ORDER_INTENT = re.compile(r"\b(quero|queremos|pretendo|pretendemos|aceito|aceitamos|adjudic|encomendar|encomenda|avançar|avancar|proceder)\b", re.I) ORDER_CHANGE_VERB = re.compile(r"\b(alterar|alteração|alteracao|mudar|mudança|mudanca|trocar|substituir|corrigir|retificar)\b", re.I) ORDER_CHANGE_SUBJECT = re.compile(r"\b(produto|modelo|quantidade|morada|entrega|nif|fiscal|faturação|faturacao|condições|condicoes)\b", re.I) QUOTE_REQUEST = re.compile(r"\b(preço|preco|orçamento|orcamento|cotação|cotacao|proposta|quote)\b", re.I) PAYMENT_PROOF = re.compile(r"\b(comprovativo|transferência|transferencia|pagamento efetuado|pago|liquidado)\b", re.I) SIMULATION_AT = datetime(2026, 8, 15, tzinfo=timezone.utc) def _s(value: Any) -> str: return str(value or "").strip() def _dt(value: Any) -> datetime | None: if isinstance(value, datetime): return value if value.tzinfo else value.replace(tzinfo=timezone.utc) if not value: return None try: parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00")) return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc) except ValueError: return None def _jsonable(value: Any) -> Any: if isinstance(value, datetime): return value.isoformat() if isinstance(value, dict): return {key: _jsonable(item) for key, item in value.items()} if isinstance(value, (list, tuple)): return [_jsonable(item) for item in value] return value def _compact(value: Any, limit: int = 260) -> str: result = re.sub(r"\s+", " ", _s(value)) return result if len(result) <= limit else result[: limit - 1].rstrip() + "…" def _payload(value: Any) -> dict[str, Any]: if isinstance(value, dict): return value if isinstance(value, str) and value.strip(): try: parsed = json.loads(value) return parsed if isinstance(parsed, dict) else {} except ValueError: pass return {} def _group(rows: Iterable[dict[str, Any]], key: str = "opportunity_id") -> dict[str, list[dict[str, Any]]]: result: dict[str, list[dict[str, Any]]] = defaultdict(list) for row in rows: result[_s(row.get(key))].append(dict(row)) return result def _load( *, expected_database: str = "clientflow_codex_shadow", expected_user: str | None = "clientflow_codex", require_read_only: bool = True, ) -> dict[str, Any]: with engine.connect() as conn: conn = conn.execution_options(isolation_level="AUTOCOMMIT") transaction_started = False try: # A fresh connection commonly reports transaction_read_only=off. # Establish the protected transaction on the same connection used # for every factual read before validating a strict audit. if require_read_only: conn.execute(text("BEGIN READ ONLY")) transaction_started = True identity = conn.execute(text( "SELECT current_database(), current_user, current_setting('transaction_read_only')" )).one() if identity[0] != expected_database or (expected_user and identity[1] != expected_user): raise RuntimeError(f"refusing unexpected database identity: {identity!r}") if require_read_only and identity[2] != "on": raise RuntimeError(f"read-only simulation requires transaction_read_only=on: {identity!r}") # Preserve the development/test path's established ordering: check # its identity first, then protect the factual reads themselves. if not require_read_only: conn.execute(text("BEGIN READ ONLY")) transaction_started = True opportunities = [dict(row) for row in conn.execute(text(""" SELECT o.*, o.id::text AS id, o.local_customer_id::text, c.name AS linked_customer_name, c.tax_id, c.email AS fiscal_email, c.street_name, c.postal_zone, c.city_name, c.phone AS fiscal_phone FROM opportunities o LEFT JOIN customers c ON c.id=o.local_customer_id ORDER BY o.created_at, o.id """)).mappings()] tasks = [dict(row) for row in conn.execute(text(""" SELECT id::text, opportunity_id::text, action_code, action, note, status, due_at, created_at, done_at, metadata FROM tasks WHERE opportunity_id IS NOT NULL ORDER BY created_at """)).mappings()] messages = [dict(row) for row in conn.execute(text(""" SELECT o.id::text AS opportunity_id, m.id::text, m.direction, COALESCE(m.clean_body,m.raw_body,'') AS body, m.created_at, m.source_system, m.metadata FROM opportunities o JOIN messages m ON m.conversation_id=o.conversation_id WHERE m.source_system IN ('chatwoot','chatwoot_backfill') ORDER BY m.created_at """)).mappings()] communications = [dict(row) for row in conn.execute(text(""" SELECT id::text, opportunity_id::text, direction, classification, subject, body, status, created_at, metadata FROM communications WHERE opportunity_id IS NOT NULL ORDER BY created_at """)).mappings()] documents = [dict(row) for row in conn.execute(text(""" SELECT DISTINCT ON (COALESCE(l.opportunity_id,d.opportunity_id),d.id) COALESCE(l.opportunity_id,d.opportunity_id)::text AS opportunity_id, d.id::text, d.document_kind, d.document_type, d.document_number, d.external_id, d.status, d.payload, d.created_at, d.updated_at, COALESCE(l.relationship, CASE WHEN d.is_primary THEN 'PRIMARY' ELSE d.role END, 'PRIMARY') AS relationship, l.ended_at FROM commercial_documents d LEFT JOIN opportunity_document_links l ON l.document_id=d.id AND l.ended_at IS NULL WHERE COALESCE(l.opportunity_id,d.opportunity_id) IS NOT NULL ORDER BY COALESCE(l.opportunity_id,d.opportunity_id),d.id,l.updated_at DESC NULLS LAST """)).mappings()] links = [dict(row) for row in conn.execute(text(""" SELECT id::text, opportunity_id::text, system, external_type, external_id, external_name, status, payload, created_at, updated_at, last_synced_at FROM operation_links ORDER BY created_at """)).mappings()] reconciliation = [dict(row) for row in conn.execute(text(""" SELECT id::text, opportunity_id::text, title, description, document_number, status, suggested_action, confidence, payload, created_at FROM reconciliation_items WHERE status IN ('open','needs_review','conflict') ORDER BY created_at """)).mappings()] finally: if transaction_started: conn.execute(text("ROLLBACK")) return { "identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]}, "opportunities": opportunities, "tasks": _group(tasks), "messages": _group(messages), "communications": _group(communications), "documents": _group(documents), "links": _group(links), "reconciliation": _group(reconciliation), "all_reconciliation": reconciliation, } def _event(row: dict[str, Any]) -> dict[str, Any]: return { "id": row.get("id"), "at": _jsonable(row.get("created_at")), "direction": row.get("direction"), "classification": row.get("classification"), "text": _compact(row.get("body") or row.get("subject")), } def _derive_record(opp: dict[str, Any], data: dict[str, Any], v1: dict[str, Any], v1_item: dict[str, Any] | None) -> dict[str, Any]: oid = _s(opp["id"]) messages = data["messages"].get(oid, []) comms = data["communications"].get(oid, []) tasks = data["tasks"].get(oid, []) docs = data["documents"].get(oid, []) links = data["links"].get(oid, []) recons = data["reconciliation"].get(oid, []) events = sorted(messages + comms, key=lambda row: _dt(row.get("created_at")) or datetime.min.replace(tzinfo=timezone.utc)) inbound = [row for row in events if _s(row.get("direction")).lower() == "inbound"] outbound = [row for row in events if _s(row.get("direction")).lower() == "outbound"] latest_in, latest_out = (inbound[-1] if inbound else None), (outbound[-1] if outbound else None) inbound_text = "\n".join(_s(row.get("body") or row.get("subject")) for row in inbound) inbound_classes = {_s(row.get("classification")).upper() for row in inbound} request_kind = "quote" if inbound_classes & {"SEND_QUOTE", "SEND_PROFORMA"} or QUOTE_REQUEST.search(inbound_text) else "info" order_rows = [row for row in inbound if _s(row.get("classification")).upper() in {"SEND_PROFORMA", "CONFIRM_PAYMENT", "SEND_INVOICE"} or ORDER_INTENT.search(_s(row.get("body") or row.get("subject")))] order_intent_at = _dt(order_rows[-1].get("created_at")) if order_rows else None current_docs = [row for row in docs if _s(row.get("relationship")).upper() == "PRIMARY" and _s(row.get("status")).lower() not in {"cancelled", "canceled", "failed"}] proformas = [row for row in current_docs if _s(row.get("document_kind")).lower() in {"quotation", "quote", "proforma"} and _s(row.get("status")).lower() != "converted"] invoices = [row for row in current_docs if _s(row.get("document_kind")).lower() == "invoice"] proforma = proformas[-1] if proformas else None invoice = invoices[-1] if invoices else None doc_number = _s((proforma or {}).get("document_number") or (proforma or {}).get("external_id")) proforma_payload = _payload((proforma or {}).get("payload")) sent_outbound = next((row for row in reversed(outbound) if doc_number and doc_number.casefold() in _s(row.get("body") or row.get("subject")).casefold()), None) proforma_sent = bool(proforma and ( (proforma or {}).get("sent_at") or proforma_payload.get("sent_at") or proforma_payload.get("clientflow_sent_evidence") or _s((proforma or {}).get("status")).lower() in {"sent", "issued_sent"} or sent_outbound )) payment_links = [row for row in links if row.get("system") == "clientflow" and row.get("external_type") == "payment"] payment = next((row for row in reversed(payment_links) if _s(row.get("status")).lower() == "confirmed"), None) payment_proof_rows = [row for row in inbound if _s(row.get("classification")).upper() == "CONFIRM_PAYMENT" or PAYMENT_PROOF.search(_s(row.get("body") or row.get("subject")))] odoo_sales = [row for row in links if row.get("system") == "odoo" and row.get("external_type") == "sale_order" and _s(row.get("status")).lower() not in {"not_found", "no_order", "cancelled"}] validation = [row for row in links if row.get("system") == "odoo" and row.get("external_type") in {"physical_validation", "physical_status"} and _s(row.get("status")).lower() in {"validated", "ready_to_ship", "shipped", "done", "delivered"}] fulfilled = any(row.get("external_type") in {"physical_status", "delivery"} and _s(row.get("status")).lower() in {"shipped", "done", "delivered"} for row in links) change_rows = [ row for row in inbound if ORDER_CHANGE_VERB.search(_s(row.get("body") or row.get("subject"))) and ORDER_CHANGE_SUBJECT.search(_s(row.get("body") or row.get("subject"))) ] change_at = _dt(change_rows[-1].get("created_at")) if change_rows else None proforma_at = _dt((proforma or {}).get("created_at")) material_change = bool(change_at and proforma_at and change_at > proforma_at) fiscal_complete = bool(opp.get("local_customer_id") and opp.get("tax_id") and opp.get("fiscal_email") and opp.get("street_name") and opp.get("postal_zone") and opp.get("city_name")) fiscal_conflict = bool((_payload(opp.get("metadata")).get("fiscal_conflict") or _payload(opp.get("metadata")).get("has_nif_conflict"))) blockers = [] if fiscal_conflict: blockers.append("Conflicting fiscal/NIF evidence.") conflict_recons = [row for row in recons if _s(row.get("status")).lower() in {"needs_review", "conflict"}] if conflict_recons: blockers.append("Unresolved document reconciliation conflict.") reconstructed = _s(_payload(opp.get("metadata")).get("clientflow_record_mode")) in { "reconstructed_invoice_review", "historical_reconstructed", "legacy_review" } or any(marker in _s(opp.get("title")).casefold() for marker in ("processo reconstruído", "sem oportunidade")) if reconstructed and not (invoice or payment or odoo_sales): blockers.append("Reconstructed process lacks corroborating structured evidence.") if invoice and not payment: blockers.append("Structured invoice exists without confirmed payment evidence; correction/reconstruction flow is unspecified.") if odoo_sales and (not payment or not invoice): blockers.append("Odoo execution evidence exists without the mandatory linked payment and invoice evidence.") status = _s(opp.get("status")).upper() stage = _s(opp.get("stage")).upper() lost = status in {"LOST", "NO_INTEREST"} or stage in {"LOST", "NO_INTEREST", "ARCHIVED"} info_sent = bool(latest_out and (not latest_in or _dt(latest_out.get("created_at")) >= _dt(latest_in.get("created_at")))) followup_satisfied = any( _s(task.get("status")).lower() == "pending" and _s(task.get("action_code")).upper().startswith("FOLLOW_UP_") and latest_in and _dt(latest_in.get("created_at")) > (_dt(task.get("created_at")) or datetime.max.replace(tzinfo=timezone.utc)) for task in tasks ) sparse = not events and not current_docs and not links review_required = ( bool(material_change and (payment or invoice)) or (reconstructed and sparse) or bool(invoice and not payment) or bool(odoo_sales and (not payment or not invoice)) ) facts = derive_business_facts( opportunity_id=oid, terminal=status in TERMINAL or stage in TERMINAL, explicitly_lost=lost, review_required=review_required, fiscal_blocked=fiscal_conflict, document_reconciliation_required=False, customer_request=bool(inbound), request_kind=request_kind, latest_relevant_inbound_at=_dt((latest_in or {}).get("created_at")), latest_relevant_outbound_at=_dt((latest_out or {}).get("created_at")), info_or_offer_sent=info_sent, order_intent=bool(order_rows), order_intent_at=order_intent_at, fiscal_identity_evidence=fiscal_complete, proforma_exists=bool(proforma), proforma_sent=proforma_sent, proforma_created_at=proforma_at, proforma_sent_at=_dt((sent_outbound or {}).get("created_at")), potential_payment_evidence=bool(payment_proof_rows and not payment), payment_confirmed=bool(payment), payment_confirmed_at=_dt((payment or {}).get("created_at")), invoice_exists=bool(invoice), invoice_created_at=_dt((invoice or {}).get("created_at")), odoo_order_exists=bool(odoo_sales), odoo_order_validated=bool(validation), fulfillment_complete=fulfilled, material_order_change=material_change, material_order_change_at=change_at, later_customer_inbound_satisfies_followup=followup_satisfied, blockers=blockers, audit_task_codes=[f"{task.get('action_code')}:{task.get('status')}" for task in tasks], ) decision = derive_v2_operational_queue(facts) confidence = decision.confidence ambiguity = [] if sparse and not lost: confidence = "low" ambiguity.append("No message, structured document, payment, or Odoo evidence is linked.") if stage in {"QUOTE_SENT", "PROFORMA_SENT", "WAITING_PAYMENT"} and not proforma: confidence = "low" ambiguity.append("V1 stage suggests a formal offer, but no current structured proforma is linked.") if stage == "PAYMENT_CONFIRMED" and not payment: confidence = "low" ambiguity.append("V1 stage says payment confirmed, but no confirmed payment operation link exists.") if stage in {"WON", "SHIPPED", "ODOO_ORDER_CREATED", "IN_PRODUCTION"} and not odoo_sales: confidence = "low" ambiguity.append("V1 stage implies execution, but no Odoo sale-order link exists.") decision = replace(decision, confidence=confidence) v1_action = _s((v1_item or {}).get("current_action_code") or v1.get("action_code")) or None v1_queue = _s((v1_item or {}).get("operational_queue")) or "not_current" pending = [task for task in tasks if _s(task.get("status")).lower() == "pending"] call_task = next((task for task in pending if _s(task.get("action_code")).upper() == "CALL_CUSTOMER"), None) call_due = bool(call_task and (not _dt(call_task.get("due_at")) or _dt(call_task.get("due_at")) <= SIMULATION_AT)) followup_codes = ( {"FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA"} if decision.business_state == "AWAITING_CUSTOMER" else {"FOLLOW_UP_PAYMENT"} if decision.business_state == "AWAITING_PAYMENT" else set() ) followup_task = next(( task for task in pending if _s(task.get("action_code")).upper() in followup_codes ), None) followup_satisfied_now = bool( followup_task and latest_in and _dt(latest_in.get("created_at")) > (_dt(followup_task.get("created_at")) or SIMULATION_AT) ) or bool(followup_task and payment and _dt(payment.get("created_at")) > (_dt(followup_task.get("created_at")) or SIMULATION_AT)) followup_due = bool( followup_task and not followup_satisfied_now and _dt(followup_task.get("due_at")) and _dt(followup_task.get("due_at")) <= SIMULATION_AT ) followup_future = bool( followup_task and not followup_satisfied_now and _dt(followup_task.get("due_at")) and _dt(followup_task.get("due_at")) > SIMULATION_AT ) formal_doc_for_reconciliation = bool(current_docs) reconciliation_blocking = bool( formal_doc_for_reconciliation and (conflict_recons or v1_action == "RECONCILE_DOCUMENTS") ) fiscal_required = decision.next_action in {"CREATE_PROFORMA", "CREATE_INVOICE"} if blockers and any("conflict" in blocker.casefold() or "without confirmed payment" in blocker.casefold() for blocker in blockers): diagnostic_status = "conflicting_evidence" elif sparse or (reconstructed and not (invoice and payment)) or (odoo_sales and (not invoice or not payment)): diagnostic_status = "incomplete_history" elif ambiguity: diagnostic_status = "ambiguous" else: diagnostic_status = "clear" raw = derive_effective_operational_action( decision, integration_exception=v1_queue == "exception", scheduled_call_current=call_due, due_followup_action=_s((followup_task or {}).get("action_code")).upper() if followup_due else None, future_followup_action=_s((followup_task or {}).get("action_code")).upper() if followup_future else None, fiscal_complete=fiscal_complete, fiscal_required=fiscal_required, reconciliation_blocking=reconciliation_blocking, diagnostic_status=diagnostic_status, ) strong_current_evidence = bool( raw.precedence in {"integration_exception", "scheduled_call", "document_prerequisite", "fiscal_prerequisite"} or (decision.business_state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "EXCEPTION"} and diagnostic_status == "conflicting_evidence") or followup_due or (decision.next_action in {"SEND_PROFORMA"} and proforma) or (decision.next_action == "CONFIRM_PAYMENT" and payment_proof_rows) or (decision.next_action == "CREATE_INVOICE" and payment) or (decision.next_action in {"PREPARE_ORDER", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY"} and invoice and payment) or (decision.next_action == "CREATE_PROFORMA" and order_rows and not proforma) or (decision.next_action in {"SEND_INFO", "SEND_QUOTE"} and latest_in and (not latest_out or _dt(latest_in.get("created_at")) > _dt(latest_out.get("created_at")))) or decision.operational_queue in {"waiting", "not_current"} ) safe = derive_safe_operational_action( raw, v1_action=v1_action, v1_queue=v1_queue, strong_current_evidence=strong_current_evidence, ) raw_v2 = raw.to_dict() safe_v2 = safe.to_dict() v2_action = raw.effective_operational_action if decision.business_state == "REVIEW_REQUIRED": classification = "REVIEW_REQUIRED" elif ambiguity: classification = "AMBIGUOUS" elif v1_action == v2_action and v1_queue == raw.effective_operational_queue: classification = "UNCHANGED" elif v1_queue in {"do_now", "review", "exception"} and raw.effective_operational_queue in {"waiting", "not_current"}: classification = "DEMOTED_TO_WAITING" elif v1_queue in {"waiting", "backlog", "not_current"} and raw.effective_operational_queue in {"do_now", "review", "exception"}: classification = "PROMOTED_TO_CURRENT" else: classification = "ACTION_CHANGED" return _jsonable({ "opportunity_id": oid, "title": opp.get("title"), "customer": opp.get("linked_customer_name") or opp.get("customer_name") or opp.get("customer_email"), "v1": {"commercial_stage": opp.get("stage"), "lifecycle_state": opp.get("lifecycle_state"), "current_action": v1_action, "operational_queue": v1_queue, "reason": (v1_item or {}).get("eligibility_reason_code") or v1.get("reason")}, "raw_v2": raw_v2, "safe_v2": safe_v2, # Compatibility alias for first-iteration report consumers. "v2": raw_v2, "classification": classification, "evidence": { "latest_relevant_inbound": _event(latest_in) if latest_in else None, "latest_relevant_outbound": _event(latest_out) if latest_out else None, "order_intent": [_event(row) for row in order_rows[-3:]], "fiscal_customer_identity": {"complete": fiscal_complete, "customer_id": _s(opp.get("local_customer_id")), "tax_id_present": bool(opp.get("tax_id"))}, "proforma": [{key: _jsonable(row.get(key)) for key in ("id", "external_id", "document_kind", "document_number", "status", "relationship", "created_at")} for row in proformas], "payment": [{key: _jsonable(row.get(key)) for key in ("id", "status", "external_name", "created_at")} for row in payment_links], "invoice": [{key: _jsonable(row.get(key)) for key in ("id", "external_id", "document_number", "status", "relationship", "created_at")} for row in invoices], "odoo": [{key: _jsonable(row.get(key)) for key in ("id", "external_type", "external_id", "external_name", "status", "created_at")} for row in links if row.get("system") == "odoo"], "blockers": blockers, "pending_tasks_for_audit_only": [ {key: _jsonable(task.get(key)) for key in ("id", "action_code", "status", "due_at", "created_at")} for task in tasks if task.get("status") == "pending" ], "ambiguity": ambiguity, "strong_current_evidence": strong_current_evidence, "reconciliation_blocking": reconciliation_blocking, "fiscal_required_for_transition": fiscal_required, "scheduled_call_current": call_due, "followup_due": followup_due, "followup_future": followup_future, "followup_satisfied": followup_satisfied_now, "diagnostic_status": diagnostic_status, "is_reconstructed": reconstructed, "reconciliation": [{ "id": row.get("id"), "document_number": row.get("document_number"), "status": row.get("status"), "created_at": _jsonable(row.get("created_at")), } for row in recons], "material_order_change_evidence": [_event(row) for row in change_rows[-3:]], }, }) def _material_keys(row: dict[str, Any]) -> set[str]: evidence = row["evidence"] keys = set() for link in evidence.get("odoo", []): if link.get("external_type") != "sale_order": continue if _s(link.get("external_id")): keys.add(f"odoo_sale_id:{_s(link['external_id']).casefold()}") sale_name = _s(link.get("external_name")) if sale_name and re.fullmatch(r"[A-Z]{1,4}[-/]?[0-9]{2,}", sale_name, re.I): keys.add(f"odoo_sale_name:{sale_name.casefold()}") for kind in ("invoice", "proforma"): for doc in evidence.get(kind, []): if _s(doc.get("external_id")): keys.add(f"jasmin_{kind}_id:{_s(doc['external_id']).casefold()}") if _s(doc.get("document_number")): keys.add(f"jasmin_{kind}_number:{_s(doc['document_number']).casefold()}") return keys def _apply_material_identity(records: list[dict[str, Any]]) -> list[dict[str, Any]]: parent = {row["opportunity_id"]: row["opportunity_id"] for row in records} def find(value: str) -> str: while parent[value] != value: parent[value] = parent[parent[value]] value = parent[value] return value def union(left: str, right: str) -> None: a, b = find(left), find(right) if a != b: parent[b] = a by_key: dict[str, list[str]] = defaultdict(list) for row in records: keys = sorted(_material_keys(row)) row["material_identity_keys"] = keys for key in keys: by_key[key].append(row["opportunity_id"]) for ids in by_key.values(): for oid in ids[1:]: union(ids[0], oid) groups: dict[str, list[dict[str, Any]]] = defaultdict(list) for row in records: groups[find(row["opportunity_id"])].append(row) report = [] for group in groups.values(): if len(group) < 2: row = group[0] row["material_process_key"] = next(iter(row["material_identity_keys"]), f"opportunity:{row['opportunity_id']}") row["canonical_process_id"] = row["opportunity_id"] row["duplicate_process_ids"] = [] continue def score(row: dict[str, Any]) -> tuple[int, str]: evidence = row["evidence"] value = 0 value += 50 if not evidence.get("is_reconstructed") else 0 value += 20 if evidence.get("invoice") else 0 value += 20 if evidence.get("payment") else 0 value += 15 if evidence.get("proforma") else 0 value += 15 if any(link.get("external_type") == "sale_order" for link in evidence.get("odoo", [])) else 0 value += 10 if evidence.get("latest_relevant_inbound") or evidence.get("latest_relevant_outbound") else 0 value += 8 if "processo reconstruído" in _s(row.get("title")).casefold() else 0 value -= 8 if "sem oportunidade" in _s(row.get("title")).casefold() else 0 return value, row["opportunity_id"] canonical = max(group, key=score) common = set(canonical["material_identity_keys"]) for row in group: common &= set(row["material_identity_keys"]) preferred = sorted(common, key=lambda key: (0 if key.startswith("odoo_sale_id:") else 1, key)) process_key = preferred[0] if preferred else sorted(canonical["material_identity_keys"])[0] duplicates = [row["opportunity_id"] for row in group if row is not canonical] canonical["material_process_key"] = process_key canonical["canonical_process_id"] = canonical["opportunity_id"] canonical["duplicate_process_ids"] = duplicates for duplicate in group: if duplicate is canonical: continue duplicate["material_process_key"] = process_key duplicate["canonical_process_id"] = canonical["opportunity_id"] duplicate["duplicate_process_ids"] = [] duplicate["raw_v2"] = suppress_duplicate_representation( EffectiveOperationalDecision(**duplicate["raw_v2"]), canonical_process_id=canonical["opportunity_id"], ).to_dict() duplicate["safe_v2"] = suppress_duplicate_representation( EffectiveOperationalDecision(**duplicate["safe_v2"]), canonical_process_id=canonical["opportunity_id"], ).to_dict() duplicate["classification"] = "DUPLICATE_REPRESENTATION" report.append({ "material_process_key": process_key, "canonical_process_id": canonical["opportunity_id"], "duplicate_process_ids": duplicates, "identity_keys": sorted(set.intersection(*(set(row["material_identity_keys"]) for row in group))), "canonical_reason": "Highest factual completeness; prefers non-reconstructed process and explicit reconstructed process over an unassociated synthetic record.", }) return report def _totals(items: Iterable[dict[str, Any]], queue_key: str) -> dict[str, int]: counts = Counter(_s(item.get(queue_key)) or "not_current" for item in items) return { "current_work": sum(counts[name] for name in ("do_now", "review", "exception")), "do_now": counts["do_now"], "review": counts["review"], "waiting": counts["waiting"], "backlog": counts["backlog"], "exception": counts["exception"], "not_current": counts["not_current"], } def collect( *, expected_database: str = "clientflow_codex_shadow", expected_user: str | None = "clientflow_codex", require_read_only: bool = True, expected_opportunity_count: int | None = None, require_opportunities: bool = False, ) -> dict[str, Any]: data = _load( expected_database=expected_database, expected_user=expected_user, require_read_only=require_read_only, ) opportunities = data["opportunities"] if expected_opportunity_count is not None and len(opportunities) != expected_opportunity_count: raise RuntimeError( f"expected {expected_opportunity_count} opportunities, found {len(opportunities)}" ) if require_opportunities and not opportunities: raise RuntimeError("Flow v2 projection requires at least one opportunity") ids = [_s(opp["id"]) for opp in opportunities] v1_decisions = get_opportunity_next_actions(ids) operations = get_operations_summary(limit=200) all_v1_items = [] for key in ("work_items", "waiting_items", "backlog_items", "not_current_items"): all_v1_items.extend(operations.get(key, [])) by_opp = {_s(item.get("opportunity_id")): item for item in all_v1_items if item.get("opportunity_id")} records = [_derive_record(opp, data, v1_decisions.get(_s(opp["id"]), {}), by_opp.get(_s(opp["id"]))) for opp in opportunities] material_identity = _apply_material_identity(records) classifications = Counter(row["classification"] for row in records) standalone = [] for item in all_v1_items: if item.get("opportunity_id"): continue action = _s(item.get("current_action_code")) or None queue = _s(item.get("operational_queue")) or "backlog" projection = { "business_state": None, "business_next_action": None, "effective_operational_action": action, "effective_operational_queue": queue, "reason": "Standalone canonical Operations work is outside the standard commercial flow and is preserved.", "confidence": "high", "precedence": "preserved_non_opportunity", } standalone.append({ "candidate_key": item.get("work_item_key") or f"standalone:{item.get('source')}:{item.get('id')}", "opportunity_id": None, "title": item.get("title"), "customer": item.get("customer_name"), "v1": {"current_action": action, "operational_queue": queue, "reason": item.get("eligibility_reason_code")}, "raw_v2": dict(projection), "safe_v2": dict(projection), "classification": "UNCHANGED", "source": item.get("source"), }) universe = records + standalone transitions = Counter(( row["v1"]["current_action"] or "", row["raw_v2"]["effective_operational_action"] or "", row["safe_v2"]["effective_operational_action"] or "", ) for row in universe) v2_actions = Counter(row["safe_v2"]["effective_operational_action"] or "" for row in universe) for action in ( "SEND_INFO", "SEND_QUOTE", "CREATE_PROFORMA", "SEND_PROFORMA", "CONFIRM_PAYMENT", "CREATE_INVOICE", "PREPARE_ORDER", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY", ): v2_actions.setdefault(action, 0) current_v1 = [row for row in universe if row["v1"]["operational_queue"] in {"do_now", "review", "exception"}] false_negatives = [] current_obligation_mapping = [] for row in current_v1: oid = row.get("opportunity_id") mapping = { "opportunity_id": oid, "title": row.get("title"), "customer": row.get("customer"), "v1_action": row["v1"]["current_action"], "v1_queue": row["v1"]["operational_queue"], "raw_business_state": row["raw_v2"].get("business_state"), "raw_action": row["raw_v2"]["effective_operational_action"], "raw_queue": row["raw_v2"]["effective_operational_queue"], "safe_action": row["safe_v2"]["effective_operational_action"], "safe_queue": row["safe_v2"]["effective_operational_queue"], "disposition": "UNCHANGED" if ( row["v1"]["current_action"] == row["safe_v2"]["effective_operational_action"] and row["v1"]["operational_queue"] == row["safe_v2"]["effective_operational_queue"] ) else "REPLACED", } current_obligation_mapping.append(mapping) if row["safe_v2"]["effective_operational_queue"] not in {"do_now", "review", "exception"} or row["v1"]["current_action"] != row["safe_v2"]["effective_operational_action"]: false_negatives.append({ "opportunity_id": oid, "title": row.get("title"), "customer": row.get("customer"), "v1_action": row["v1"]["current_action"], "v1_queue": row["v1"]["operational_queue"], "raw_state": row["raw_v2"].get("business_state"), "raw_action": row["raw_v2"]["effective_operational_action"], "safe_action": row["safe_v2"]["effective_operational_action"], "safe_queue": row["safe_v2"]["effective_operational_queue"], "reason": row["safe_v2"]["reason"], "pending_tasks": row.get("evidence", {}).get("pending_tasks_for_audit_only", []), }) promotion_audit = [] for row in records: if row["v1"]["operational_queue"] in {"do_now", "review", "exception"}: continue if row["raw_v2"]["effective_operational_queue"] not in {"do_now", "review", "exception"}: continue # Reproduce the 45-item first-simulation promotion cohort: cases without # a stage/evidence ambiguity, plus the nine Odoo-only cases that the # first model had incorrectly promoted toward completion. Invoice-only # REVIEW_REQUIRED cases were already classified as review, not promotion. odoo_only_completion_error = any( "Odoo execution evidence exists" in blocker for blocker in row["evidence"].get("blockers", []) ) and not row["evidence"].get("invoice") if row["evidence"].get("ambiguity"): continue if row["raw_v2"]["business_state"] == "REVIEW_REQUIRED" and not odoo_only_completion_error: continue if row["raw_v2"]["precedence"] in {"fiscal_prerequisite", "document_prerequisite"}: audit_class = "BLOCKER_PRECEDENCE_ERROR" elif row["safe_v2"]["precedence"] == "safe_ambiguous_review": audit_class = "AMBIGUOUS_REVIEW" elif row["safe_v2"]["effective_operational_action"] == row["raw_v2"]["effective_operational_action"]: audit_class = "REAL_PROMOTION" else: audit_class = "HISTORICAL_EVIDENCE_FALSE_POSITIVE" promotion_audit.append({ "opportunity_id": row["opportunity_id"], "customer": row["customer"], "title": row["title"], "v1_action": row["v1"]["current_action"], "v1_queue": row["v1"]["operational_queue"], "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"], "evidence": row["evidence"], "confidence": row["raw_v2"]["confidence"], "classification": audit_class, }) invoice_without_payment = [] for row in records: if not row["evidence"]["invoice"] or row["evidence"]["payment"]: continue metadata = _payload(next(opp for opp in opportunities if _s(opp["id"]) == row["opportunity_id"]).get("metadata")) text_blob = json.dumps(row["evidence"], ensure_ascii=False).casefold() if _s(metadata.get("payment_terms")).lower() in {"after_delivery", "payment_after_delivery", "pos_entrega"}: category = "PAYMENT_AFTER_INVOICE_ALLOWED" elif "comprovativo" in text_blob or "pagamento" in text_blob: category = "PAYMENT_EVIDENCE_MISSING" elif any(term in text_blob for term in ("nota de crédito", "nota de credito", "corrigir", "anular")): category = "FINANCIAL_CORRECTION_REQUIRED" elif not row["evidence"]["order_intent"] and not row["evidence"]["proforma"]: category = "PREMATURE_INVOICE" else: category = "UNKNOWN_REVIEW" invoice_without_payment.append({ "opportunity_id": row["opportunity_id"], "title": row["title"], "customer": row["customer"], "classification": category, "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"], "evidence": row["evidence"], }) review_audit = [] for row in universe: if row["safe_v2"]["effective_operational_queue"] != "review": continue diagnostic = row.get("evidence", {}).get("diagnostic_status", "clear") if row["v1"]["operational_queue"] == "review" and row["v1"]["current_action"] == row["safe_v2"]["effective_operational_action"]: category = "EXISTING_VALID_REVIEW" elif row["v1"]["operational_queue"] not in {"do_now", "review", "exception"}: category = "REAL_CURRENT_PROMOTION" elif diagnostic == "incomplete_history": category = "HISTORICAL_INCOMPLETE" elif diagnostic == "ambiguous": category = "DIAGNOSTIC_ONLY" else: category = "ACTIONABLE_REVIEW" review_audit.append({ "opportunity_id": row.get("opportunity_id"), "title": row.get("title"), "v1": row["v1"], "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"], "diagnostic_status": diagnostic, "classification": category, }) backlog_delta = [] for row in universe: if row["v1"]["operational_queue"] != "backlog" or row["safe_v2"]["effective_operational_queue"] == "backlog": continue safe_queue = row["safe_v2"]["effective_operational_queue"] if row["safe_v2"]["precedence"] == "duplicate_representation": category = "DEDUPLICATED" elif safe_queue in {"do_now", "review", "exception"}: category = "PROMOTED_TO_CURRENT" elif safe_queue == "waiting": category = "MOVED_TO_WAITING" elif row.get("evidence", {}).get("strong_current_evidence"): category = "LEGITIMATE_BACKLOG_REMOVAL" else: category = "SHOULD_REMAIN_BACKLOG" backlog_delta.append({ "opportunity_id": row.get("opportunity_id"), "title": row.get("title"), "v1_action": row["v1"]["current_action"], "safe_v2": row["safe_v2"], "evidence": row.get("evidence", {}), "classification": category, }) current_delta = [] for row in universe: if row["v1"]["operational_queue"] in {"do_now", "review", "exception"}: continue if row["safe_v2"]["effective_operational_queue"] not in {"do_now", "review", "exception"}: continue precedence = row["safe_v2"]["precedence"] diagnostic = row.get("evidence", {}).get("diagnostic_status", "clear") if precedence == "due_followup": category = "DUE_FOLLOW_UP" elif precedence in {"fiscal_prerequisite", "document_prerequisite", "integration_exception"}: category = "BLOCKER" elif precedence == "duplicate_representation": category = "DUPLICATE" elif row.get("evidence", {}).get("strong_current_evidence") and row["safe_v2"]["effective_operational_queue"] != "review": category = "REAL_NEW_OBLIGATION" elif diagnostic in {"ambiguous", "incomplete_history"}: category = "DIAGNOSTIC_ONLY" elif row["safe_v2"]["effective_operational_queue"] == "review": category = "ACTIONABLE_REVIEW" else: category = "FALSE_PROMOTION" current_delta.append({ "opportunity_id": row.get("opportunity_id"), "title": row.get("title"), "v1": row["v1"], "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"], "evidence": row.get("evidence", {}), "classification": category, }) def action_counts(rows: list[dict[str, Any]]) -> dict[str, int]: return dict(Counter(row["safe_v2"]["effective_operational_action"] or "" for row in rows)) action_count_scopes = { "all_candidates": action_counts(universe), "current_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] in {"do_now", "review", "exception"}]), "do_now_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "do_now"]), "review_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "review"]), "waiting_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "waiting"]), "backlog_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "backlog"]), } v1_totals = _totals([{"queue": row["v1"]["operational_queue"]} for row in universe], "queue") raw_totals = _totals([{"queue": row["raw_v2"]["effective_operational_queue"]} for row in universe], "queue") safe_totals = _totals([{"queue": row["safe_v2"]["effective_operational_queue"]} for row in universe], "queue") result = { "generated_at": datetime.now(timezone.utc), "database": data["identity"], "opportunity_count": len(records), "candidate_universe_count": len(universe), "v1_totals": v1_totals, "raw_v2_totals": raw_totals, "safe_v2_totals": safe_totals, "classifications": dict(classifications), "v2_actions": dict(v2_actions), "transitions": [{"v1_action": old, "raw_v2_action": raw, "safe_v2_action": safe, "count": count} for (old, raw, safe), count in transitions.most_common()], "possible_false_negatives": false_negatives, "v1_current_obligation_mapping": current_obligation_mapping, "promotions_audit": promotion_audit, "invoice_without_payment": invoice_without_payment, "summary": { "raw_ambiguous_count": sum(row["raw_v2"]["confidence"] != "high" for row in records), "safe_overrides_count": sum( (row["raw_v2"]["effective_operational_action"], row["raw_v2"]["effective_operational_queue"]) != (row["safe_v2"]["effective_operational_action"], row["safe_v2"]["effective_operational_queue"]) for row in universe ), "preserved_v1_obligations": sum(row["disposition"] == "UNCHANGED" for row in current_obligation_mapping), "real_promotions": sum(row["classification"] == "REAL_PROMOTION" for row in promotion_audit), "rejected_promotions": sum(row["classification"] == "HISTORICAL_EVIDENCE_FALSE_POSITIVE" for row in promotion_audit), "ambiguous_promotions": sum(row["classification"] == "AMBIGUOUS_REVIEW" for row in promotion_audit), "fiscal_blockers_preserved": sum(row["raw_v2"]["precedence"] == "fiscal_prerequisite" for row in records), "reconciliation_blockers_preserved": sum(row["raw_v2"]["precedence"] == "document_prerequisite" for row in records), "non_opportunity_canonical_work_preserved": len(standalone), "diagnostic_ambiguous_not_current": sum( row.get("evidence", {}).get("diagnostic_status") in {"ambiguous", "incomplete_history"} and row["safe_v2"]["effective_operational_queue"] == "not_current" for row in records ), "actionable_review": sum(row["classification"] in {"ACTIONABLE_REVIEW", "EXISTING_VALID_REVIEW", "REAL_CURRENT_PROMOTION"} for row in review_audit), "due_followups": sum(row["safe_v2"]["precedence"] == "due_followup" for row in records), "safe_v2_current_minus_v1": len(current_delta), "backlog_delta_explained": len(backlog_delta), "duplicate_material_processes": len(material_identity), "duplicate_current_cards_suppressed": sum(len(row["duplicate_process_ids"]) for row in material_identity), }, "action_counts": action_count_scopes, "review_audit": review_audit, "backlog_delta": backlog_delta, "current_delta": current_delta, "material_identity": material_identity, "standalone_canonical_items": standalone, "opportunities": records, } return _jsonable(result) def _named(records: list[dict[str, Any]], name: str) -> list[dict[str, Any]]: folded = name.casefold() return [row for row in records if folded in f"{_s(row.get('title'))} {_s(row.get('customer'))}".casefold()] def _comparison(result: dict[str, Any]) -> str: lines = [ "BLIF FLOW V2 SHADOW SIMULATION", "", f"Database: {result['database']}", f"Opportunities: {result['opportunity_count']}", f"Comparable candidate universe: {result['candidate_universe_count']}", "", "CENTRO DE TRABALHO", "metric V1 RAW V2 SAFE V2", ] for key in ("current_work", "do_now", "review", "waiting", "backlog", "exception", "not_current"): lines.append( f"{key:<30} {result['v1_totals'].get(key, 0):>5}" f" {result['raw_v2_totals'].get(key, 0):>7} {result['safe_v2_totals'].get(key, 0):>7}" ) lines += ["", "SUMMARY"] for key, value in result["summary"].items(): lines.append(f"{key}: {value}") for scope in ("all_candidates", "current_only", "do_now_only", "review_only", "waiting_only", "backlog_only"): lines += ["", f"SAFE V2 ACTIONS — {scope}"] for action, count in sorted(result["action_counts"][scope].items(), key=lambda item: (-item[1], item[0])): lines.append(f"{action:<36} {count:>5}") lines += ["", "CLASSIFICATIONS"] for name in ("UNCHANGED", "ACTION_CHANGED", "DEMOTED_TO_WAITING", "PROMOTED_TO_CURRENT", "REVIEW_REQUIRED", "AMBIGUOUS"): lines.append(f"{name:<30} {result['classifications'].get(name, 0):>5}") lines += ["", "V1 ACTION -> RAW V2 ACTION -> SAFE V2 ACTION"] for row in result["transitions"]: lines.append(f"{row['v1_action']} -> {row['raw_v2_action']} -> {row['safe_v2_action']}: {row['count']}") lines += ["", "POSSIBLE FALSE NEGATIVES"] if not result["possible_false_negatives"]: lines.append("None.") for row in result["possible_false_negatives"]: lines.append(json.dumps(row, ensure_ascii=False, sort_keys=True)) lines += ["", "EVERY CURRENT V1 OBLIGATION -> V2"] for row in result["v1_current_obligation_mapping"]: lines.append(json.dumps(row, ensure_ascii=False, sort_keys=True)) for name in ("INSTALBEIRA", "ENGEXICON", "RZSOLAR", "X MAT", "CONSTRURECUP", "PANORAMIC SUCCESS"): lines += ["", name] matches = _named(result["opportunities"], name) lines.extend(json.dumps(row, ensure_ascii=False, sort_keys=True) for row in matches) if not matches: lines.append("No opportunity title/customer match.") return "\n".join(lines) + "\n" def main() -> None: result = collect() PROJECTION.write_text(json.dumps(result, ensure_ascii=False, indent=2), encoding="utf-8") ambiguous = [row for row in result["opportunities"] if row["classification"] in {"AMBIGUOUS", "REVIEW_REQUIRED"}] AMBIGUOUS.write_text(json.dumps(ambiguous, ensure_ascii=False, indent=2), encoding="utf-8") PROMOTIONS.write_text(json.dumps(result["promotions_audit"], ensure_ascii=False, indent=2), encoding="utf-8") INVOICE_WITHOUT_PAYMENT.write_text(json.dumps(result["invoice_without_payment"], ensure_ascii=False, indent=2), encoding="utf-8") REVIEW_AUDIT.write_text(json.dumps(result["review_audit"], ensure_ascii=False, indent=2), encoding="utf-8") BACKLOG_DELTA.write_text(json.dumps(result["backlog_delta"], ensure_ascii=False, indent=2), encoding="utf-8") CURRENT_DELTA.write_text(json.dumps(result["current_delta"], ensure_ascii=False, indent=2), encoding="utf-8") MATERIAL_IDENTITY.write_text(json.dumps(result["material_identity"], ensure_ascii=False, indent=2), encoding="utf-8") COMPARISON.write_text(_comparison(result), encoding="utf-8") print(_comparison(result), end="") print(f"Output: {PROJECTION}\nOutput: {COMPARISON}\nOutput: {AMBIGUOUS}\nOutput: {PROMOTIONS}\nOutput: {INVOICE_WITHOUT_PAYMENT}\nOutput: {REVIEW_AUDIT}\nOutput: {BACKLOG_DELTA}\nOutput: {CURRENT_DELTA}\nOutput: {MATERIAL_IDENTITY}") if __name__ == "__main__": main()