Files
clientflow_backend/scripts/simulate_blif_flow_v2.py
2026-08-15 23:03:55 +00:00

843 lines
49 KiB
Python

#!/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() -> dict[str, Any]:
with engine.connect() as conn:
conn = conn.execution_options(isolation_level="AUTOCOMMIT")
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if tuple(identity) != ("clientflow_codex_shadow", "clientflow_codex", "on"):
raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
conn.execute(text("BEGIN READ ONLY"))
try:
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:
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() -> dict[str, Any]:
data = _load()
opportunities = data["opportunities"]
if len(opportunities) != 328:
raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}")
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 "<NONE>",
row["raw_v2"]["effective_operational_action"] or "<WAIT/NONE>",
row["safe_v2"]["effective_operational_action"] or "<WAIT/NONE>",
) for row in universe)
v2_actions = Counter(row["safe_v2"]["effective_operational_action"] or "<WAIT/NONE>" 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 "<NONE>" 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()