5 Commits

Author SHA1 Message Date
plx
32e7773957 feat: add BLIF Flow v2 shadow compare cutover layer 2026-08-16 00:12:20 +00:00
plx
3b60c3fb70 feat: apply validated BLIF Flow v2 task repairs 2026-08-16 00:02:01 +00:00
plx
acbbd1a84b feat: add BLIF Flow v2 data repair planner 2026-08-15 23:50:54 +00:00
plx
e99da64d8b feat: add BLIF Flow v2 persistence scaffolding 2026-08-15 23:32:49 +00:00
plx
20bc91dec5 feat: add shadow BLIF Flow v2 projection 2026-08-15 23:03:55 +00:00
27 changed files with 3593 additions and 390 deletions

View File

@@ -40,6 +40,16 @@ INTERNAL_ACTION_CODES = {
ACTION_CODES = TRIAGE_ACTION_CODES | INTERNAL_ACTION_CODES ACTION_CODES = TRIAGE_ACTION_CODES | INTERNAL_ACTION_CODES
# Persistence vocabulary for the shadow-only Flow v2 projection. These are not
# added to TRIAGE_ACTION_CODES or ACTION_CODES, so no existing task/LLM/runtime
# behavior changes.
FLOW_V2_BUSINESS_ACTION_CODES = {
"CREATE_PROFORMA",
"CREATE_INVOICE",
"VALIDATE_ODOO_ORDER",
"COMPLETE_OPPORTUNITY",
}
ACTION_MAP = { ACTION_MAP = {
"CALL_CUSTOMER": { "CALL_CUSTOMER": {
"route": "vendas", "route": "vendas",

View File

@@ -1544,20 +1544,14 @@ def _parse_opportunity_dt(value: object):
def _opportunity_lifecycle_state(opp: dict) -> str: def _opportunity_lifecycle_state(opp: dict) -> str:
state = str(opp.get("lifecycle_state") or "active").strip().lower() or "active" state = str(opp.get("lifecycle_state") or "active").strip().lower() or "active"
now = datetime.now(timezone.utc) now = datetime.now(timezone.utc)
pending_call = str(opp.get("pending_follow_up_action_code") or "").strip().upper() == "CALL_CUSTOMER" pending_followup = str(opp.get("pending_follow_up_action_code") or "").strip().upper()
pending_call_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at")) pending_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at"))
if pending_call: if pending_followup:
return "follow_up_due" if pending_call_due and pending_call_due <= now else "scheduled_follow_up" return "follow_up_due" if pending_due and pending_due <= now else "scheduled_follow_up"
# Compatibility only: a stale denormalized lifecycle marker cannot invent # Compatibility-only timestamps/lifecycle markers cannot create work. An
# a scheduled call when no active CALL_CUSTOMER task exists. # active pending follow-up task and its due_at are the operational source.
if state == "scheduled_follow_up": if state in {"scheduled_follow_up", "follow_up_due"}:
return "active" return "active"
nurture_until = _parse_opportunity_dt(opp.get("nurture_until"))
next_follow_up = _parse_opportunity_dt(opp.get("next_follow_up_at"))
if state == "nurture" and nurture_until and nurture_until <= now:
return "follow_up_due"
if state in {"awaiting_customer", "active", "scheduled_follow_up"} and next_follow_up and next_follow_up <= now:
return "follow_up_due"
return state return state

View File

@@ -0,0 +1,216 @@
"""Persistence scaffolding for the rebuildable BLIF Flow v2 projection.
The factual sources remain authoritative. This module writes only the additive
projection/audit tables introduced by migration 011 and never updates stages or
tasks.
"""
from __future__ import annotations
import hashlib
import json
import re
from datetime import datetime, timezone
from typing import Any, Iterable
from sqlalchemy import text
from app.config import settings
FLOW_VERSION = "blif-flow-v2-shadow-3"
WRITE_DATABASE_ALLOWLIST = frozenset({"clientflow_codex_test"})
def _jsonable(value: Any) -> Any:
if isinstance(value, datetime):
return value.isoformat()
if isinstance(value, dict):
return {str(key): _jsonable(item) for key, item in value.items()}
if isinstance(value, (list, tuple)):
return [_jsonable(item) for item in value]
return value
def _evidence_refs(row: dict[str, Any]) -> list[dict[str, Any]]:
evidence = row.get("evidence") or {}
refs: list[dict[str, Any]] = []
for name in ("latest_relevant_inbound", "latest_relevant_outbound"):
item = evidence.get(name)
if item and item.get("id"):
refs.append({"source": "message", "role": name, "id": item["id"], "at": item.get("at")})
for name, source in (("proforma", "jasmin_proforma"), ("invoice", "jasmin_invoice"),
("payment", "payment"), ("odoo", "odoo"),
("reconciliation", "reconciliation")):
for item in evidence.get(name) or []:
if item.get("id"):
refs.append({
"source": source, "id": item["id"],
"external_id": item.get("external_id"),
"document_number": item.get("document_number"),
"external_type": item.get("external_type"),
"status": item.get("status"),
})
return refs
def _projection_value(row: dict[str, Any], derived_at: datetime) -> dict[str, Any]:
raw = row["raw_v2"]
duplicate = row.get("canonical_process_id") != row.get("opportunity_id")
evidence_refs = _evidence_refs(row)
reason_code = str(raw.get("precedence") or "business_transition").upper()
stable = {
"opportunity_id": row["opportunity_id"],
"material_process_key": row["material_process_key"],
"canonical_opportunity_id": row["canonical_process_id"],
"is_duplicate_representation": duplicate,
"business_state": raw["business_state"],
"business_next_action": raw.get("business_next_action"),
"diagnostic_status": raw.get("diagnostic_status") or row.get("evidence", {}).get("diagnostic_status") or "clear",
"confidence": raw.get("confidence") or "low",
"reason_code": reason_code,
"reason_text": raw.get("reason") or "",
"evidence_refs": evidence_refs,
"flow_version": FLOW_VERSION,
}
fingerprint = hashlib.sha256(
json.dumps(_jsonable(stable), ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8")
).hexdigest()
return stable | {"source_fingerprint": fingerprint, "derived_at": derived_at}
def _derive_all() -> list[dict[str, Any]]:
# Reuse the validated shadow evidence adapter without making it authoritative.
from scripts.simulate_blif_flow_v2 import collect
report = collect(
expected_database="clientflow_codex_test",
expected_user="clientflow_codex_test",
# collect() still opens its factual read phase with BEGIN READ ONLY;
# the session default may be read-write in the isolated test database.
require_read_only=False,
)
return list(report["opportunities"])
def rebuild_blif_flow_v2_projection(
*,
mode: str | None = None,
derived_rows: Iterable[dict[str, Any]] | None = None,
allowed_databases: frozenset[str] = WRITE_DATABASE_ALLOWLIST,
target_schema: str = "public",
connection: Any | None = None,
) -> dict[str, Any]:
"""Idempotently rebuild projection rows and state-change transitions.
``off`` is a no-op. ``shadow`` and ``compare`` write only the same additive
projection tables. Compare affects observation at the read boundary, never
the persisted business rows. Authoritative remains deliberately disabled.
"""
selected_mode = str(mode or settings.blif_flow_v2_mode or "off").strip().lower()
if selected_mode == "off":
return {"mode": "off", "projection_count": 0, "transitions_written": 0, "disabled": True}
if selected_mode not in {"shadow", "compare"}:
raise RuntimeError(f"BLIF Flow v2 mode {selected_mode!r} is disabled; only off/shadow/compare are safe")
if not re.fullmatch(r"[a-z_][a-z0-9_]*", target_schema):
raise ValueError("invalid target_schema")
rows = list(derived_rows) if derived_rows is not None else _derive_all()
derived_at = datetime.now(timezone.utc)
values = [_projection_value(row, derived_at) for row in rows]
if len({value["opportunity_id"] for value in values}) != len(values):
raise RuntimeError("Flow v2 projection requires exactly one derived row per opportunity")
owns_connection = connection is None
if connection is None:
from app.db import engine
conn = engine.connect().execution_options(isolation_level="AUTOCOMMIT")
else:
conn = connection
try:
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if identity[0] not in allowed_databases:
raise RuntimeError(f"refusing Flow v2 projection write to database {identity[0]!r}")
if owns_connection:
conn.exec_driver_sql("BEGIN READ WRITE")
try:
if conn.execute(text("SELECT current_setting('transaction_read_only')")).scalar_one() != "off":
raise RuntimeError("Flow v2 projection rebuild requires an explicit READ WRITE transaction")
conn.exec_driver_sql(f'SET LOCAL search_path TO "{target_schema}"')
existing = {
str(row["opportunity_id"]): dict(row)
for row in conn.execute(text("""
SELECT opportunity_id::text, business_state, source_fingerprint
FROM opportunity_flow_state_v2
""")).mappings()
}
transitions_written = 0
for value in values:
previous = existing.get(value["opportunity_id"])
if previous is None or previous["business_state"] != value["business_state"]:
result = conn.execute(text("""
INSERT INTO opportunity_flow_transitions (
opportunity_id, from_state, to_state, reason_code, reason_text,
evidence_refs, flow_version, source_fingerprint
) VALUES (
CAST(:opportunity_id AS UUID), :from_state, :to_state, :reason_code,
:reason_text, CAST(:evidence_refs AS JSONB), :flow_version, :source_fingerprint
)
ON CONFLICT (opportunity_id, from_state, to_state, flow_version, source_fingerprint)
DO NOTHING
"""), {
**value, "from_state": previous["business_state"] if previous else None,
"to_state": value["business_state"],
"evidence_refs": json.dumps(_jsonable(value["evidence_refs"]), ensure_ascii=False),
})
transitions_written += result.rowcount
conn.execute(text("""
INSERT INTO opportunity_flow_state_v2 (
opportunity_id, material_process_key, canonical_opportunity_id,
is_duplicate_representation, business_state, business_next_action,
diagnostic_status, confidence, reason_code, reason_text, evidence_refs,
flow_version, source_fingerprint, derived_at, updated_at
) VALUES (
CAST(:opportunity_id AS UUID), :material_process_key,
CAST(:canonical_opportunity_id AS UUID), :is_duplicate_representation,
:business_state, :business_next_action, :diagnostic_status, :confidence,
:reason_code, :reason_text, CAST(:evidence_refs_json AS JSONB), :flow_version,
:source_fingerprint, :derived_at, now()
)
ON CONFLICT (opportunity_id) DO UPDATE SET
material_process_key=EXCLUDED.material_process_key,
canonical_opportunity_id=EXCLUDED.canonical_opportunity_id,
is_duplicate_representation=EXCLUDED.is_duplicate_representation,
business_state=EXCLUDED.business_state,
business_next_action=EXCLUDED.business_next_action,
diagnostic_status=EXCLUDED.diagnostic_status,
confidence=EXCLUDED.confidence,
reason_code=EXCLUDED.reason_code,
reason_text=EXCLUDED.reason_text,
evidence_refs=EXCLUDED.evidence_refs,
flow_version=EXCLUDED.flow_version,
source_fingerprint=EXCLUDED.source_fingerprint,
derived_at=EXCLUDED.derived_at,
updated_at=CASE
WHEN opportunity_flow_state_v2.source_fingerprint IS DISTINCT FROM EXCLUDED.source_fingerprint
THEN now() ELSE opportunity_flow_state_v2.updated_at END
"""), value | {
"evidence_refs_json": json.dumps(_jsonable(value["evidence_refs"]), ensure_ascii=False)
})
if owns_connection:
conn.exec_driver_sql("COMMIT")
except Exception:
if owns_connection:
conn.exec_driver_sql("ROLLBACK")
raise
finally:
if owns_connection:
conn.close()
return {
"mode": selected_mode, "projection_count": len(values),
"canonical_count": sum(not value["is_duplicate_representation"] for value in values),
"duplicate_count": sum(value["is_duplicate_representation"] for value in values),
"transitions_written": transitions_written, "disabled": False,
}

View File

@@ -21,6 +21,10 @@ class Settings(BaseSettings):
# TIMESTAMPTZ nem os índices usados pelo schema core. # TIMESTAMPTZ nem os índices usados pelo schema core.
database_url: str database_url: str
clientflow_persist: bool = True clientflow_persist: bool = True
# Flow v2 remains non-authoritative. Compare observes V1/V2 differences
# while returning V1; authoritative is accepted by config only to fail
# closed at the switching boundary.
blif_flow_v2_mode: Literal["off", "shadow", "compare", "authoritative"] = "off"
clientflow_webhook_secret: str = "" clientflow_webhook_secret: str = ""
clientflow_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy" clientflow_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy"

View File

@@ -249,15 +249,7 @@ def build_opportunity_evidence(
quote = _find_doc(docs, QUOTE_KINDS) quote = _find_doc(docs, QUOTE_KINDS)
invoice = _find_doc(docs, INVOICE_KINDS) invoice = _find_doc(docs, INVOICE_KINDS)
# A existência do documento comercial prova apenas que o orçamento foi quote_sent = bool(quote) or _completed_send_quote_task_evidence(tasks)
# criado/associado. O envio ao cliente exige evidência própria.
#
# Compatibilidade histórica: QUOTE_SENT é também uma afirmação canónica
# explícita de que o orçamento já foi enviado.
quote_sent = (
stage == "QUOTE_SENT"
or _completed_send_quote_task_evidence(tasks)
)
pending_task = None pending_task = None
invalid_payment_task = None invalid_payment_task = None

View File

@@ -0,0 +1,159 @@
"""Pure BLIF Flow v2 historical repair classification and simulation helpers.
The functions in this module never access the database and never mutate business
state. They deliberately treat tasks as obligations/history, not factual proof.
"""
from __future__ import annotations
from dataclasses import asdict, dataclass, field
from datetime import datetime, timezone
from typing import Any
TASK_CLASSIFICATIONS = {
"VALID_CURRENT", "SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE",
"PREMATURE", "AMBIGUOUS",
}
FOLLOWUP_ACTIONS = {
"FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT", "CALL_CUSTOMER",
}
STAGE_RANK = {
"INQUIRY": 0, "AWAITING_CUSTOMER": 1, "PROFORMA_REQUIRED": 2,
"PROFORMA_CREATED": 3, "AWAITING_PAYMENT": 4, "INVOICE_REQUIRED": 5,
"INVOICE_CREATED": 6, "ODOO_ORDER_REQUIRED": 6,
"ODOO_ORDER_CREATED": 7, "ODOO_ORDER_VALIDATED": 8, "COMPLETED": 9,
}
ACTION_REQUIRED_RANK = {
"SEND_INFO": 0, "SEND_QUOTE": 0, "CREATE_PROFORMA": 2,
"SEND_PROFORMA": 3, "FOLLOW_UP_PROFORMA": 4, "CONFIRM_PAYMENT": 4,
"FOLLOW_UP_PAYMENT": 4, "CREATE_INVOICE": 5, "SEND_INVOICE": 6,
"PREPARE_ORDER": 6, "VALIDATE_ODOO_ORDER": 7,
"COMPLETE_OPPORTUNITY": 8,
}
@dataclass(frozen=True)
class TaskRepairContext:
task_id: str
opportunity_id: str | None
action_code: str
created_at: datetime | None = None
due_at: datetime | None = None
business_state: str | None = None
business_next_action: str | None = None
material_process_key: str | None = None
is_duplicate_representation: bool = False
canonical_opportunity_id: str | None = None
later_inbound_event: dict[str, Any] | None = None
later_outbound_event: dict[str, Any] | None = None
proforma_exists: bool = False
proforma_sent: bool = False
payment_confirmed: bool = False
invoice_exists: bool = False
invoice_sent: bool = False
odoo_order_exists: bool = False
odoo_order_validated: bool = False
terminal: bool = False
same_obligation_task_id: str | None = None
evidence_refs: tuple[dict[str, Any], ...] = field(default_factory=tuple)
@dataclass(frozen=True)
class TaskRepairDecision:
classification: str
resolution_code: str | None
reason: str
confidence: str
safety_tier: str
auto_repair_safe: bool
human_review_required: bool
resolved_by_event_id: str | None = None
superseded_by_task_id: str | None = None
def to_dict(self) -> dict[str, Any]:
return asdict(self)
def _decision(classification: str, reason: str, *, event: dict[str, Any] | None = None,
superseded_by: str | None = None, tier: str = "HIGH") -> TaskRepairDecision:
if classification not in TASK_CLASSIFICATIONS:
raise ValueError(classification)
auto = tier == "HIGH" and classification != "VALID_CURRENT"
return TaskRepairDecision(
classification, None if classification == "VALID_CURRENT" else classification,
reason, "high" if tier == "HIGH" else "medium" if tier == "MEDIUM" else "low",
tier, auto, tier != "HIGH",
(event or {}).get("opportunity_event_id"), superseded_by,
)
def classify_pending_task(context: TaskRepairContext) -> TaskRepairDecision:
"""Classify one pending task using facts available at the audit instant."""
action = context.action_code.upper()
state = (context.business_state or "").upper()
if context.is_duplicate_representation:
return _decision(
"DUPLICATE",
f"Obligation belongs only to a duplicate representation of canonical process {context.canonical_opportunity_id}.",
)
if context.same_obligation_task_id:
return _decision(
"DUPLICATE", "The same material obligation has another canonical pending task.",
superseded_by=context.same_obligation_task_id,
)
if action in FOLLOWUP_ACTIONS:
event = context.later_inbound_event
if action == "FOLLOW_UP_PAYMENT" and context.payment_confirmed:
return _decision("SATISFIED_BY_EVENT", "Confirmed payment fact satisfies the payment follow-up.")
if event:
return _decision("SATISFIED_BY_EVENT", "A later inbound customer event satisfies the follow-up.", event=event)
return _decision("VALID_CURRENT", "No later satisfying event exists; age or overdue status alone never closes a follow-up.")
satisfied = {
"SEND_INFO": context.later_outbound_event,
"SEND_QUOTE": context.later_outbound_event,
"CREATE_PROFORMA": context.proforma_exists,
"SEND_PROFORMA": context.proforma_sent,
"CONFIRM_PAYMENT": context.payment_confirmed,
"CREATE_INVOICE": context.invoice_exists,
"SEND_INVOICE": context.invoice_sent,
"PREPARE_ORDER": context.odoo_order_exists,
"VALIDATE_ODOO_ORDER": context.odoo_order_validated,
"COMPLETE_OPPORTUNITY": context.terminal,
}.get(action, False)
if satisfied:
event = satisfied if isinstance(satisfied, dict) else None
return _decision("SATISFIED_BY_EVENT", f"Later factual evidence satisfies {action}.", event=event)
required = ACTION_REQUIRED_RANK.get(action)
rank = STAGE_RANK.get(state)
if required is not None and rank is not None:
if rank > required:
return _decision("SUPERSEDED", f"Factual process advanced to {state}, beyond the {action} obligation.")
if rank < required:
return _decision("PREMATURE", f"{action} requires prerequisites not present in factual state {state}.")
if action == "SEND_PROFORMA" and not context.proforma_exists:
return _decision("PREMATURE", "No structured current proforma exists; a send task is not document evidence.")
if action == "SEND_INVOICE" and not context.invoice_exists:
return _decision("PREMATURE", "No structured invoice exists and protected proforma/payment prerequisites are absent.")
if action == context.business_next_action:
return _decision("VALID_CURRENT", "Task matches the current factual Flow v2 obligation.")
if action in {"SUPPORT", "MARK_NO_INTEREST", "REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS"}:
return _decision("VALID_CURRENT", "Historically valid operator obligation has no factual evidence of satisfaction or supersession.")
if not context.opportunity_id:
return _decision("VALID_CURRENT", "Standalone obligation is outside opportunity business transitions and is preserved.")
return _decision("AMBIGUOUS", "Available factual evidence does not deterministically establish validity or safe removal.", tier="LOW")
def simulate_high_repairs(task_rows: list[dict[str, Any]]) -> dict[str, Any]:
removed = [row for row in task_rows if row["safety_tier"] == "HIGH" and row["auto_repair_safe"]]
counts = {name: sum(row["classification"] == name for row in removed) for name in TASK_CLASSIFICATIONS}
return {"pending_before": len(task_rows), "pending_after": len(task_rows) - len(removed),
"removed": removed, "removed_by_classification": counts}

View File

@@ -8,7 +8,6 @@ from .types import (
ACTION_CLOSE_OPPORTUNITY, ACTION_CLOSE_OPPORTUNITY,
ACTION_CONFIRM_PAYMENT, ACTION_CONFIRM_PAYMENT,
ACTION_CREATE_QUOTE, ACTION_CREATE_QUOTE,
ACTION_SEND_QUOTE,
ACTION_FOLLOW_UP, ACTION_FOLLOW_UP,
ACTION_FOLLOW_UP_PAYMENT, ACTION_FOLLOW_UP_PAYMENT,
ACTION_NO_ACTION, ACTION_NO_ACTION,
@@ -124,26 +123,16 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
return OpportunityDecision(next_action, "Conflito fiscal/NIF bloqueia ações financeiras.", blocked_actions=blocked_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version) return OpportunityDecision(next_action, "Conflito fiscal/NIF bloqueia ações financeiras.", blocked_actions=blocked_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
if not e.has_fiscal_customer: if not e.has_fiscal_customer:
# A ausência de cliente fiscal é prontidão operacional, não intenção blocked_actions.extend(_blocked(profile, code, "cliente fiscal por associar") for code in SENSITIVE_DOCUMENT_ACTIONS)
# comercial. Não deve substituir a próxima ação da oportunidade. next_action = _action(profile, ACTION_VALIDATE_FISCAL_CUSTOMER, "Associar/validar cliente fiscal antes de documentos oficiais.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#cliente" if e.opportunity_id else None)
# O bloqueio fiscal é aplicado apenas mais abaixo quando uma transição return OpportunityDecision(next_action, "Cliente fiscal ainda não associado.", blocked_actions=blocked_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
# concreta (por exemplo faturação após pagamento confirmado) exige
# efetivamente os dados fiscais.
warnings.append(
"Cliente fiscal ainda não associado; validar apenas quando uma "
"operação documental atual exigir dados fiscais."
)
if e.has_reconciliation_candidate: if e.has_reconciliation_candidate:
next_action = _action(profile, ACTION_RECONCILE_DOCUMENTS, f"Confirmar evidência encontrada: {e.reconciliation_label or 'documento/candidato'}.", priority="alta", target_url="/reconciliation") next_action = _action(profile, ACTION_RECONCILE_DOCUMENTS, f"Confirmar evidência encontrada: {e.reconciliation_label or 'documento/candidato'}.", priority="alta", target_url="/reconciliation")
return OpportunityDecision(next_action, "Há evidência de reconciliação por validar.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version) return OpportunityDecision(next_action, "Há evidência de reconciliação por validar.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
if not e.has_quote and not e.has_invoice: if not e.has_quote and not e.has_invoice:
# Não promover automaticamente qualquer oportunidade para orçamento. if e.quote_sent:
# A criação/reconciliação de orçamento só é trabalho atual quando o
# estágio comercial demonstra que o cliente pediu ou já recebeu um
# orçamento.
if e.stage == "QUOTE_SENT" and e.quote_sent:
next_action = _action( next_action = _action(
profile, profile,
ACTION_RECONCILE_DOCUMENTS, ACTION_RECONCILE_DOCUMENTS,
@@ -162,26 +151,9 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
profile_name=profile.name, profile_name=profile.name,
decision_version=profile.version, decision_version=profile.version,
) )
next_action = _action(profile, ACTION_CREATE_QUOTE, "Criar/enviar orçamento antes de pedir pagamento ou emitir fatura.", target_url=f"/opportunities/{e.opportunity_id}#documentos" if e.opportunity_id else None)
if e.stage == "QUOTE_REQUESTED":
next_action = _action(
profile,
ACTION_CREATE_QUOTE,
"Criar/enviar orçamento solicitado pelo cliente.",
target_url=f"/opportunities/{e.opportunity_id}#documentos" if e.opportunity_id else None,
)
available_actions.append(next_action) available_actions.append(next_action)
return OpportunityDecision( return OpportunityDecision(next_action, "Ainda não há orçamento/fatura associado.", available_actions=available_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state="no_document", physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
next_action,
"Existe pedido de orçamento e ainda não há documento comercial associado.",
available_actions=available_actions,
warnings=warnings,
commercial_stage=COMMERCIAL_STAGE_REVIEW,
financial_state="no_document",
physical_state=_physical_state(e),
profile_name=profile.name,
decision_version=profile.version,
)
if e.payment_terms == PAYMENT_AFTER_DELIVERY: if e.payment_terms == PAYMENT_AFTER_DELIVERY:
if e.stage == "SHIPMENT_CREATED" and not e.payment_confirmed: if e.stage == "SHIPMENT_CREATED" and not e.payment_confirmed:
@@ -254,9 +226,7 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
next_action = _action(profile, ACTION_PREPARE_ORDER, "Pagamento após entrega: criar/associar venda Odoo e avançar preparação sem exigir pagamento confirmado.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None) next_action = _action(profile, ACTION_PREPARE_ORDER, "Pagamento após entrega: criar/associar venda Odoo e avançar preparação sem exigir pagamento confirmado.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None)
return OpportunityDecision(next_action, "Condição pós-entrega permite avançar Odoo/preparação sem pagamento prévio.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version) return OpportunityDecision(next_action, "Condição pós-entrega permite avançar Odoo/preparação sem pagamento prévio.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
if e.odoo_ready and not e.order_shipped: if e.odoo_ready and not e.order_shipped:
if not e.has_invoice and ( if not e.has_invoice and e.has_fiscal_customer and not e.fiscal_data_complete:
not e.has_fiscal_customer or not e.fiscal_data_complete
):
blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos")) blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos"))
next_action = _action( next_action = _action(
profile, profile,
@@ -295,76 +265,21 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
next_action = _action(profile, ACTION_WAIT_PRODUCTION, "Pagamento após entrega: venda Odoo criada; aguardar WH/OUT ficar pronto/concluído antes de emitir fatura.", priority="normal", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None) next_action = _action(profile, ACTION_WAIT_PRODUCTION, "Pagamento após entrega: venda Odoo criada; aguardar WH/OUT ficar pronto/concluído antes de emitir fatura.", priority="normal", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None)
return OpportunityDecision(next_action, "Aguardar estado da encomenda/WH-OUT no Odoo; ordens de fabrico são apenas detalhe técnico.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version) return OpportunityDecision(next_action, "Aguardar estado da encomenda/WH-OUT no Odoo; ordens de fabrico são apenas detalhe técnico.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
# Um orçamento criado/associado ainda não significa orçamento enviado.
# O envio ao cliente é uma obrigação humana/documental própria e deve
# acontecer antes de qualquer etapa de pagamento.
if (
e.has_quote
and not e.quote_sent
and not e.has_invoice
and not e.payment_confirmed
):
next_action = _action(
profile,
ACTION_SEND_QUOTE,
f"Orçamento {e.quote_number or ''} criado/associado. Enviar o documento ao cliente.",
priority="alta",
target_url=(
f"/tasks/{e.pending_task_id}"
if e.pending_task_id and e.pending_task_action_code == ACTION_SEND_QUOTE
else f"/opportunities/{e.opportunity_id}#documentos"
if e.opportunity_id
else None
),
document_id=e.quote_id,
document_number=e.quote_number,
)
available_actions.append(next_action)
return OpportunityDecision(
next_action,
"O orçamento existe, mas ainda não há evidência de que tenha sido enviado ao cliente.",
available_actions=available_actions,
warnings=warnings,
commercial_stage=COMMERCIAL_STAGE_REVIEW,
financial_state=_financial_state(e),
physical_state=_physical_state(e),
profile_name=profile.name,
decision_version=profile.version,
)
# Default/BLIF normal sequence: budget document, payment, invoice, then preparation/shipping. # Default/BLIF normal sequence: budget document, payment, invoice, then preparation/shipping.
if ( if e.payment_terms in {PAYMENT_BEFORE_SHIPPING, "", "undefined", "agreement"} and e.has_quote and not e.payment_confirmed:
e.payment_terms in {PAYMENT_BEFORE_SHIPPING, "", "undefined", "agreement"}
and e.has_quote
and e.quote_sent
and not e.has_invoice
and not e.payment_confirmed
):
next_action = _action( next_action = _action(
profile, profile,
ACTION_NO_ACTION, ACTION_CONFIRM_PAYMENT,
f"Orçamento {e.quote_number or ''} enviado. Aguardar decisão do cliente ou evidência de pagamento.", f"Orçamento {e.quote_number or ''} associado. Confirmar pagamento antes de emitir fatura.",
force_label="Aguardar cliente / pagamento", priority="alta",
target_url=f"/opportunities/{e.opportunity_id}#operacao" if e.opportunity_id else None, target_url=f"/opportunities/{e.opportunity_id}#operacao" if e.opportunity_id else None,
document_id=e.quote_id, document_id=e.quote_id,
document_number=e.quote_number, document_number=e.quote_number,
) )
return OpportunityDecision( available_actions.append(next_action)
next_action, return OpportunityDecision(next_action, "Fluxo normal BLIF exige pagamento confirmado depois do orçamento e antes da fatura.", available_actions=available_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_WAITING_PAYMENT, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
"O orçamento foi enviado e ainda não existe evidência de pagamento que exija validação. "
"Aguardar o cliente; o follow-up comercial assume quando ficar devido.",
warnings=warnings,
commercial_stage=COMMERCIAL_STAGE_WAITING_PAYMENT,
financial_state=_financial_state(e),
physical_state=_physical_state(e),
profile_name=profile.name,
decision_version=profile.version,
)
if e.payment_confirmed and not e.has_invoice and e.has_fiscal_customer and not e.fiscal_data_complete:
if e.payment_confirmed and not e.has_invoice and (
not e.has_fiscal_customer or not e.fiscal_data_complete
):
blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos")) blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos"))
next_action = _action( next_action = _action(
profile, profile,
@@ -464,38 +379,5 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
next_action = _action(profile, ACTION_PREPARE_ORDER, f"Fatura {e.invoice_number or ''} e pagamento confirmados. Criar/validar venda Odoo e preparação.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None) next_action = _action(profile, ACTION_PREPARE_ORDER, f"Fatura {e.invoice_number or ''} e pagamento confirmados. Criar/validar venda Odoo e preparação.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None)
return OpportunityDecision(next_action, "Fatura e pagamento OK; falta validar execução/Odoo.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version) return OpportunityDecision(next_action, "Fatura e pagamento OK; falta validar execução/Odoo.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
if e.stage in {"NEW_LEAD", "INFO_REQUESTED", "INFO_SENT"}: next_action = _action(profile, ACTION_FOLLOW_UP, "Rever tarefas, documentos e próximos contactos.", priority="baixa", target_url=f"/opportunities/{e.opportunity_id}" if e.opportunity_id else None)
next_action = _action( return OpportunityDecision(next_action, "Sem regra específica aplicável; manter em acompanhamento.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_QUOTE_SENT, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
profile,
ACTION_NO_ACTION,
"Sem transição documental atual. Manter o estágio comercial e aguardar a próxima obrigação operacional real.",
target_url=f"/opportunities/{e.opportunity_id}" if e.opportunity_id else None,
)
return OpportunityDecision(
next_action,
"O estágio comercial atual não exige orçamento, faturação ou follow-up imediato gerado pelo motor central.",
warnings=warnings,
commercial_stage=e.stage,
financial_state=_financial_state(e),
physical_state=_physical_state(e),
profile_name=profile.name,
decision_version=profile.version,
)
next_action = _action(
profile,
ACTION_FOLLOW_UP,
"Rever tarefas, documentos e próximos contactos.",
priority="baixa",
target_url=f"/opportunities/{e.opportunity_id}" if e.opportunity_id else None,
)
return OpportunityDecision(
next_action,
"Sem regra específica aplicável; manter em acompanhamento.",
warnings=warnings,
commercial_stage=e.stage or COMMERCIAL_STAGE_QUOTE_SENT,
financial_state=_financial_state(e),
physical_state=_physical_state(e),
profile_name=profile.name,
decision_version=profile.version,
)

View File

@@ -21,8 +21,7 @@ ACTION_NO_ACTION = "NO_ACTION"
ACTION_OPEN_TASK = "OPEN_TASK" ACTION_OPEN_TASK = "OPEN_TASK"
ACTION_VALIDATE_FISCAL_CUSTOMER = "VALIDATE_FISCAL_CUSTOMER" ACTION_VALIDATE_FISCAL_CUSTOMER = "VALIDATE_FISCAL_CUSTOMER"
ACTION_RECONCILE_DOCUMENTS = "RECONCILE_DOCUMENTS" ACTION_RECONCILE_DOCUMENTS = "RECONCILE_DOCUMENTS"
ACTION_CREATE_QUOTE = "CREATE_QUOTE" ACTION_CREATE_QUOTE = "CREATE_JASMIN_QUOTE"
ACTION_SEND_QUOTE = "SEND_QUOTE"
ACTION_CONFIRM_PAYMENT = "CONFIRM_PAYMENT" ACTION_CONFIRM_PAYMENT = "CONFIRM_PAYMENT"
ACTION_SEND_INVOICE = "SEND_INVOICE" ACTION_SEND_INVOICE = "SEND_INVOICE"
ACTION_CONFIRM_ORDER = "CONFIRM_ORDER" ACTION_CONFIRM_ORDER = "CONFIRM_ORDER"

View File

@@ -0,0 +1,273 @@
"""Pure, shadow-only BLIF Flow v2 business-state projection.
This module has no database or V1 dependencies. In particular, task fields are
kept only for audit output and never establish document or payment facts.
"""
from __future__ import annotations
from dataclasses import asdict, dataclass, field
from datetime import datetime
from typing import Any
@dataclass(frozen=True)
class BusinessFacts:
opportunity_id: str = ""
terminal: bool = False
explicitly_lost: bool = False
exception: bool = False
review_required: bool = False
fiscal_blocked: bool = False
document_reconciliation_required: bool = False
customer_request: bool = False
request_kind: str = "info" # info | quote
latest_relevant_inbound_at: datetime | None = None
latest_relevant_outbound_at: datetime | None = None
info_or_offer_sent: bool = False
order_intent: bool = False
order_intent_at: datetime | None = None
fiscal_identity_evidence: bool = False
proforma_exists: bool = False
proforma_sent: bool = False
proforma_created_at: datetime | None = None
proforma_sent_at: datetime | None = None
potential_payment_evidence: bool = False
payment_confirmed: bool = False
payment_confirmed_at: datetime | None = None
invoice_exists: bool = False
invoice_created_at: datetime | None = None
odoo_order_exists: bool = False
odoo_order_validated: bool = False
fulfillment_complete: bool = False
material_order_change: bool = False
material_order_change_at: datetime | None = None
later_customer_inbound_satisfies_followup: bool = False
blockers: tuple[str, ...] = ()
audit_task_codes: tuple[str, ...] = field(default=(), compare=False)
@dataclass(frozen=True)
class FlowV2Decision:
business_state: str
next_action: str | None
operational_queue: str
reason: str
confidence: str = "high"
def to_dict(self) -> dict[str, Any]:
return asdict(self)
@dataclass(frozen=True)
class EffectiveOperationalDecision:
business_state: str
business_next_action: str | None
effective_operational_action: str | None
effective_operational_queue: str
reason: str
confidence: str
precedence: str
diagnostic_status: str = "clear"
legacy_preserved_action: bool = False
def to_dict(self) -> dict[str, Any]:
return asdict(self)
def derive_business_facts(**evidence: Any) -> BusinessFacts:
"""Normalize factual adapter output without inferring facts from tasks."""
allowed = BusinessFacts.__dataclass_fields__
values = {key: value for key, value in evidence.items() if key in allowed}
for key in ("blockers", "audit_task_codes"):
if key in values and not isinstance(values[key], tuple):
values[key] = tuple(values[key] or ())
return BusinessFacts(**values)
def _after(left: datetime | None, right: datetime | None) -> bool:
return bool(left and right and left > right)
def derive_business_state(facts: BusinessFacts) -> str:
"""Derive the current state from strongest present-tense facts."""
if facts.exception:
return "EXCEPTION"
if facts.explicitly_lost:
return "LOST"
if facts.review_required:
return "REVIEW_REQUIRED"
if facts.fiscal_blocked:
return "FISCAL_BLOCKED"
if facts.document_reconciliation_required:
return "DOCUMENT_RECONCILIATION_REQUIRED"
change_after_payment = facts.material_order_change and (
facts.payment_confirmed
or facts.invoice_exists
or _after(facts.material_order_change_at, facts.payment_confirmed_at)
)
if change_after_payment:
return "REVIEW_REQUIRED"
# An invoice without confirmed payment contradicts BLIF's normal protected
# sequence. Do not silently skip payment or invent a correction flow.
if facts.invoice_exists and not facts.payment_confirmed:
return "REVIEW_REQUIRED"
if facts.odoo_order_exists and (not facts.payment_confirmed or not facts.invoice_exists):
return "REVIEW_REQUIRED"
if facts.terminal and facts.payment_confirmed and facts.invoice_exists and facts.odoo_order_validated:
return "COMPLETED"
if facts.fulfillment_complete and facts.payment_confirmed and facts.invoice_exists:
return "COMPLETED"
if facts.odoo_order_validated:
return "ODOO_ORDER_VALIDATED"
if facts.odoo_order_exists:
return "ODOO_ORDER_CREATED"
if facts.invoice_exists:
return "INVOICE_CREATED"
if facts.payment_confirmed:
return "INVOICE_REQUIRED"
change_invalidates_proforma = facts.material_order_change and (
not facts.material_order_change_at
or not facts.proforma_created_at
or _after(facts.material_order_change_at, facts.proforma_created_at)
)
if facts.order_intent and (not facts.proforma_exists or change_invalidates_proforma):
return "PROFORMA_REQUIRED"
if facts.proforma_exists:
return "PROFORMA_SENT" if facts.proforma_sent else "PROFORMA_CREATED"
if facts.order_intent:
return "ORDER_INTENT"
if facts.info_or_offer_sent and not _after(
facts.latest_relevant_inbound_at, facts.latest_relevant_outbound_at
):
return "AWAITING_CUSTOMER"
return "INQUIRY"
def derive_next_action(facts: BusinessFacts, state: str | None = None) -> FlowV2Decision:
state = state or derive_business_state(facts)
if state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED", "EXCEPTION"}:
action = {
"REVIEW_REQUIRED": "REVIEW_REQUIRED",
"FISCAL_BLOCKED": "VALIDATE_FISCAL_CUSTOMER",
"DOCUMENT_RECONCILIATION_REQUIRED": "RECONCILE_DOCUMENTS",
"EXCEPTION": "REVIEW_EXCEPTION",
}[state]
return FlowV2Decision(state, action, "review" if state != "EXCEPTION" else "exception",
"; ".join(facts.blockers) or f"{state} requires operator review.", "medium")
if state in {"LOST", "NO_INTEREST", "COMPLETED"}:
return FlowV2Decision(state, None, "not_current", "The factual process is terminal.")
if state == "INQUIRY":
if not facts.customer_request:
return FlowV2Decision(state, None, "not_current", "No current unanswered customer request is evidenced.", "low")
action = "SEND_QUOTE" if facts.request_kind == "quote" else "SEND_INFO"
return FlowV2Decision(state, action, "do_now", "Customer request has no later relevant outbound response.", "medium")
if state == "AWAITING_CUSTOMER":
return FlowV2Decision(state, None, "waiting", "Information or offer was sent; awaiting a later customer decision.")
if state in {"ORDER_INTENT", "PROFORMA_REQUIRED"}:
return FlowV2Decision("PROFORMA_REQUIRED", "CREATE_PROFORMA", "do_now",
"Customer order intent exists and no current structured proforma exists.")
if state == "PROFORMA_CREATED":
return FlowV2Decision(state, "SEND_PROFORMA", "do_now", "A current structured proforma exists but has no factual sent evidence.")
if state == "PROFORMA_SENT":
if facts.potential_payment_evidence:
return FlowV2Decision("AWAITING_PAYMENT", "CONFIRM_PAYMENT", "do_now",
"Potential payment evidence requires operator confirmation.", "medium")
return FlowV2Decision("AWAITING_PAYMENT", None, "waiting", "The current proforma was sent and payment is not confirmed.")
if state in {"PAYMENT_CONFIRMED", "INVOICE_REQUIRED"}:
return FlowV2Decision("INVOICE_REQUIRED", "CREATE_INVOICE", "do_now", "Payment is confirmed and no structured invoice exists.")
if state == "INVOICE_CREATED":
return FlowV2Decision("ODOO_ORDER_REQUIRED", "PREPARE_ORDER", "do_now", "Invoice exists and no Odoo sale order exists.")
if state == "ODOO_ORDER_CREATED":
return FlowV2Decision(state, "VALIDATE_ODOO_ORDER", "do_now", "Odoo sale order exists but is not validated.")
if state == "ODOO_ORDER_VALIDATED":
return FlowV2Decision(state, "COMPLETE_OPPORTUNITY", "do_now", "The validated Odoo order is ready for opportunity completion.")
return FlowV2Decision("REVIEW_REQUIRED", "REVIEW_REQUIRED", "review", f"No safe Flow v2 rule for {state}.", "low")
def derive_v2_operational_queue(facts: BusinessFacts) -> FlowV2Decision:
return derive_next_action(facts, derive_business_state(facts))
def derive_effective_operational_action(
business: FlowV2Decision,
*,
integration_exception: bool = False,
scheduled_call_current: bool = False,
due_followup_action: str | None = None,
future_followup_action: str | None = None,
fiscal_complete: bool = True,
fiscal_required: bool = False,
reconciliation_blocking: bool = False,
diagnostic_status: str = "clear",
) -> EffectiveOperationalDecision:
"""Apply operational prerequisites without changing the business state."""
action, queue, reason, precedence = (
business.next_action, business.operational_queue, business.reason, "business_transition"
)
if integration_exception:
action, queue, reason, precedence = "REVIEW_EXCEPTION", "exception", "An integration failure blocks current work.", "integration_exception"
elif scheduled_call_current:
action, queue, reason, precedence = "CALL_CUSTOMER", "do_now", "An explicit scheduled call is currently due.", "scheduled_call"
elif reconciliation_blocking:
action, queue, reason, precedence = "RECONCILE_DOCUMENTS", "review", "A real formal document requires current association/reconciliation.", "document_prerequisite"
elif fiscal_required and not fiscal_complete:
action, queue, reason, precedence = "VALIDATE_FISCAL_CUSTOMER", "do_now", "Fiscal identity is required before the current formal-document transition.", "fiscal_prerequisite"
elif due_followup_action and business.operational_queue == "waiting":
action, queue, reason, precedence = due_followup_action, "do_now", "A scheduled external follow-up is due and remains unsatisfied.", "due_followup"
elif future_followup_action and business.operational_queue == "waiting":
action, queue, reason, precedence = future_followup_action, "waiting", "A scheduled external follow-up is not due yet.", "future_followup"
return EffectiveOperationalDecision(
business.business_state, business.next_action, action, queue, reason,
business.confidence, precedence, diagnostic_status,
)
def derive_safe_operational_action(
raw: EffectiveOperationalDecision,
*,
v1_action: str | None,
v1_queue: str | None,
strong_current_evidence: bool,
) -> EffectiveOperationalDecision:
"""Conservatively preserve current V1 work when RAW evidence is uncertain."""
current = v1_queue in {"do_now", "review", "exception"}
if strong_current_evidence and (
raw.confidence == "high" or raw.business_state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "EXCEPTION"}
):
return raw
if current and v1_action:
return EffectiveOperationalDecision(
raw.business_state, raw.business_next_action, v1_action, v1_queue or "review",
"SAFE V2 preserves the current V1 obligation because RAW evidence is not strong enough to replace it.",
raw.confidence, "safe_preserve_v1",
raw.diagnostic_status, v1_action == "CREATE_JASMIN_QUOTE",
)
if v1_queue in {"backlog", "waiting"} and not strong_current_evidence:
return EffectiveOperationalDecision(
raw.business_state, raw.business_next_action, v1_action, v1_queue,
"SAFE V2 preserves the non-current V1 queue because no stronger current obligation is proven.",
raw.confidence, "safe_preserve_noncurrent", raw.diagnostic_status,
)
if raw.effective_operational_queue in {"do_now", "review", "exception"}:
return EffectiveOperationalDecision(
raw.business_state, raw.business_next_action, None, "not_current",
"Ambiguous or incomplete history is diagnostic only; it does not create current work.",
raw.confidence, "safe_diagnostic_only", raw.diagnostic_status,
)
return raw
def suppress_duplicate_representation(
projection: EffectiveOperationalDecision, *, canonical_process_id: str,
) -> EffectiveOperationalDecision:
"""Suppress a duplicate local card while retaining its diagnostic trace."""
return EffectiveOperationalDecision(
projection.business_state, projection.business_next_action, None, "not_current",
f"Duplicate representation of canonical material process {canonical_process_id}.",
"high", "duplicate_representation", "duplicate_representation",
)

View File

@@ -609,36 +609,7 @@ async def create_quotation_for_opportunity(opportunity_id: str) -> Dict[str, Any
except Exception: except Exception:
# operation_links é compatibilidade visual; não deve falhar o fluxo principal. # operation_links é compatibilidade visual; não deve falhar o fluxo principal.
pass pass
# Criar o documento no Jasmin não significa que foi enviado ao cliente. set_opportunity_stage(opportunity_id, "QUOTE_SENT", note="Orçamento Jasmin criado via ClientFlow.", created_by="jasmin_service")
# Mantemos o estágio de pedido até a ação SEND_QUOTE ser concluída.
set_opportunity_stage(
opportunity_id,
"QUOTE_REQUESTED",
note="Orçamento criado no Jasmin; falta enviar ao cliente.",
created_by="jasmin_service",
)
# Materializar a próxima obrigação humana usando o mecanismo central,
# preservando idempotência, route, prioridade e ligação à oportunidade.
try:
from app.opportunity_next_action_service import get_opportunity_next_action
from app.opportunity_action_task_materializer import ensure_pending_task_for_next_action
ensure_pending_task_for_next_action(
opportunity_id,
get_opportunity_next_action(opportunity_id),
source="jasmin_quotation_created",
actor="jasmin_service",
)
except Exception as exc:
# O orçamento Jasmin já foi criado com sucesso. Uma falha de
# materialização não pode duplicar/reverter a criação externa.
print(
f"ClientFlow SEND_QUOTE materialization failed for opportunity "
f"{opportunity_id}: {exc}",
flush=True,
)
return {"customer": customer, "quotation": doc, "quotation_id": quotation_id, "payload": payload} return {"customer": customer, "quotation": doc, "quotation_id": quotation_id, "payload": payload}

View File

@@ -24,8 +24,7 @@ from app.work_center_action_policy import (
# v4928.1.5.96 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE"} # v4928.1.5.96 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE"}
# v4928.1.5.105 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT"} # v4928.1.5.105 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT"}
# v4928.1.5.116 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"} MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
MATERIALIZED_ACTIONS = {"SEND_QUOTE", "SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
# v4928.1.5.129: central workflow emits SHIP_ORDER; operator tasks persist CREATE_SHIPMENT. # v4928.1.5.129: central workflow emits SHIP_ORDER; operator tasks persist CREATE_SHIPMENT.
MATERIALIZED_ACTIONS.add("CREATE_SHIPMENT") MATERIALIZED_ACTIONS.add("CREATE_SHIPMENT")
MATERIALIZED_ACTIONS.add("VALIDATE_PHYSICAL_ORDER") MATERIALIZED_ACTIONS.add("VALIDATE_PHYSICAL_ORDER")
@@ -110,7 +109,6 @@ def ensure_pending_task_for_next_action(
# pre-insert branch referenced these values before assignment. # pre-insert branch referenced these values before assignment.
config = get_action_config(action_code) config = get_action_config(action_code)
default_routes = { default_routes = {
"SEND_QUOTE": "vendas",
"SEND_INVOICE": "financeiro", "SEND_INVOICE": "financeiro",
"FOLLOW_UP_PAYMENT": "financeiro", "FOLLOW_UP_PAYMENT": "financeiro",
"PREPARE_ORDER": "operacoes", "PREPARE_ORDER": "operacoes",
@@ -127,7 +125,6 @@ def ensure_pending_task_for_next_action(
route = "rever" route = "rever"
default_labels = { default_labels = {
"SEND_QUOTE": "Enviar orçamento ao cliente",
"SEND_INVOICE": "Enviar fatura ao cliente", "SEND_INVOICE": "Enviar fatura ao cliente",
"FOLLOW_UP_PAYMENT": "Follow-up pagamento", "FOLLOW_UP_PAYMENT": "Follow-up pagamento",
"PREPARE_ORDER": "Preparar encomenda / Odoo", "PREPARE_ORDER": "Preparar encomenda / Odoo",
@@ -136,7 +133,6 @@ def ensure_pending_task_for_next_action(
"REVIEW_RECONSTRUCTED_PROCESS": "Validar processo reconstruído", "REVIEW_RECONSTRUCTED_PROCESS": "Validar processo reconstruído",
} }
default_descriptions = { default_descriptions = {
"SEND_QUOTE": "Orçamento criado/associado. Enviar PDF/proposta ao cliente e registar evidência.",
"SEND_INVOICE": "Fatura criada/associada. Enviar PDF ao cliente e registar evidência.", "SEND_INVOICE": "Fatura criada/associada. Enviar PDF ao cliente e registar evidência.",
"FOLLOW_UP_PAYMENT": "Encomenda concluída no Odoo/WH-OUT e fatura enviada. Acompanhar pagamento pós-entrega.", "FOLLOW_UP_PAYMENT": "Encomenda concluída no Odoo/WH-OUT e fatura enviada. Acompanhar pagamento pós-entrega.",
"PREPARE_ORDER": "Fatura e pagamento confirmados. Criar/validar venda Odoo e preparação da encomenda.", "PREPARE_ORDER": "Fatura e pagamento confirmados. Criar/validar venda Odoo e preparação da encomenda.",

View File

@@ -7,11 +7,14 @@ from __future__ import annotations
from dataclasses import asdict, dataclass from dataclasses import asdict, dataclass
from collections import defaultdict from collections import defaultdict
import json
import logging
from typing import Any, Dict, Iterable, Optional from typing import Any, Dict, Iterable, Optional
from sqlalchemy import bindparam, text from sqlalchemy import bindparam, text
from app.db import engine from app.db import engine
from app.config import settings
# Backward-compatible static anchors from v1.5.59: quotation_doc, invoice_doc, confirmar pagamento antes de emitir fatura, Pagamento confirmado com base em, Criar/enviar fatura. # Backward-compatible static anchors from v1.5.59: quotation_doc, invoice_doc, confirmar pagamento antes de emitir fatura, Pagamento confirmado com base em, Criar/enviar fatura.
from app.domain.opportunity_flow import ( from app.domain.opportunity_flow import (
OpportunityEvidence, OpportunityEvidence,
@@ -21,6 +24,9 @@ from app.domain.opportunity_flow import (
) )
logger = logging.getLogger(__name__)
@dataclass @dataclass
class OpportunityNextAction: class OpportunityNextAction:
action_code: str action_code: str
@@ -37,6 +43,67 @@ class OpportunityNextAction:
return asdict(self) return asdict(self)
def _flow_v2_mode() -> str:
return str(settings.blif_flow_v2_mode or "off").strip().lower()
def _load_v2_comparison_rows(opportunity_ids: list[str]) -> dict[str, dict[str, Any]]:
"""Read the persisted projection only; comparison must never derive writes."""
if not opportunity_ids:
return {}
with engine.connect() as conn:
rows = conn.execute(text("""
SELECT opportunity_id::text, material_process_key,
canonical_opportunity_id::text, is_duplicate_representation,
business_state, business_next_action, diagnostic_status,
confidence, reason_code, reason_text, flow_version
FROM opportunity_flow_state_v2
WHERE opportunity_id::text IN :opportunity_ids
""").bindparams(bindparam("opportunity_ids", expanding=True)),
{"opportunity_ids": opportunity_ids}).mappings().all()
return {str(row["opportunity_id"]): dict(row) for row in rows}
def compare_v1_v2_decisions(
v1_decisions: Dict[str, Dict[str, Any]],
v2_rows: dict[str, dict[str, Any]],
) -> list[dict[str, Any]]:
"""Return payload-safe structured comparisons without message content."""
comparisons = []
for opportunity_id, v1 in v1_decisions.items():
v2 = v2_rows.get(opportunity_id)
v1_action = str(v1.get("action_code") or "") or None
v2_action = (v2 or {}).get("business_next_action")
comparisons.append({
"opportunity_id": opportunity_id,
"material_process_key": (v2 or {}).get("material_process_key"),
"canonical_opportunity_id": (v2 or {}).get("canonical_opportunity_id"),
"is_duplicate_representation": bool((v2 or {}).get("is_duplicate_representation")),
"v1_action": v1_action,
"v1_reason_code": v1.get("reason_if_blocked") or v1.get("decision_version"),
"v2_business_state": (v2 or {}).get("business_state"),
"v2_business_next_action": v2_action,
"v2_reason_code": (v2 or {}).get("reason_code"),
"v2_diagnostic_status": (v2 or {}).get("diagnostic_status"),
"v2_confidence": (v2 or {}).get("confidence"),
"projection_present": v2 is not None,
"action_agrees": v2 is not None and v1_action == v2_action,
"returned_source": "v1",
})
return comparisons
def _observe_flow_v2(v1_decisions: Dict[str, Dict[str, Any]]) -> None:
mode = _flow_v2_mode()
if mode == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
if mode != "compare" or not v1_decisions:
return
rows = _load_v2_comparison_rows(list(v1_decisions))
for comparison in compare_v1_v2_decisions(v1_decisions, rows):
logger.info("blif_flow_v2_compare %s", json.dumps(comparison, sort_keys=True, default=str))
def _first_row(conn: Any, sql: str, params: Dict[str, Any]) -> Optional[Dict[str, Any]]: def _first_row(conn: Any, sql: str, params: Dict[str, Any]) -> Optional[Dict[str, Any]]:
row = conn.execute(text(sql), params).mappings().first() row = conn.execute(text(sql), params).mappings().first()
return dict(row) if row else None return dict(row) if row else None
@@ -160,6 +227,8 @@ def get_opportunity_next_action(
decision is now produced by the company workflow engine. decision is now produced by the company workflow engine.
""" """
if _flow_v2_mode() == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
evidence = (_build_db_evidence(opportunity_id, preloaded=preloaded) evidence = (_build_db_evidence(opportunity_id, preloaded=preloaded)
if preloaded is not None else _build_db_evidence(opportunity_id)) if preloaded is not None else _build_db_evidence(opportunity_id))
if evidence is None: if evidence is None:
@@ -172,8 +241,9 @@ def get_opportunity_next_action(
reason_if_blocked="opportunity_not_found", reason_if_blocked="opportunity_not_found",
).to_dict() ).to_dict()
decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile)) decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile)).to_dict()
return decision.to_dict() _observe_flow_v2({opportunity_id: decision})
return decision
def _bulk_statement(sql: str): def _bulk_statement(sql: str):
@@ -298,6 +368,8 @@ def _bulk_operation_snapshots(
def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Dict[str, Any]]: def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Dict[str, Any]]:
"""Return the same decisions as the single-item API with a fixed query count.""" """Return the same decisions as the single-item API with a fixed query count."""
if _flow_v2_mode() == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip())) ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip()))
if not ids: if not ids:
return {} return {}
@@ -376,4 +448,5 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
company_profile="blif", company_profile="blif",
) )
decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict() decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict()
_observe_flow_v2(decisions)
return decisions return decisions

View File

@@ -0,0 +1,52 @@
-- Additive, rebuildable BLIF Flow v2 projection scaffolding.
-- This migration does not alter opportunities.stage or activate Flow v2.
CREATE TABLE IF NOT EXISTS opportunity_flow_state_v2 (
opportunity_id UUID PRIMARY KEY REFERENCES opportunities(id) ON DELETE CASCADE,
material_process_key TEXT NOT NULL,
canonical_opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE CASCADE,
is_duplicate_representation BOOLEAN NOT NULL DEFAULT FALSE,
business_state TEXT NOT NULL,
business_next_action TEXT,
diagnostic_status TEXT NOT NULL DEFAULT 'clear',
confidence TEXT NOT NULL DEFAULT 'low',
reason_code TEXT NOT NULL,
reason_text TEXT NOT NULL,
evidence_refs JSONB NOT NULL DEFAULT '[]'::jsonb,
flow_version TEXT NOT NULL,
source_fingerprint TEXT NOT NULL,
derived_at TIMESTAMPTZ NOT NULL,
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_material_process_key
ON opportunity_flow_state_v2(material_process_key);
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_canonical
ON opportunity_flow_state_v2(canonical_opportunity_id);
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_business_state
ON opportunity_flow_state_v2(business_state);
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolution_code TEXT;
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolved_at TIMESTAMPTZ;
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolved_by_event_id UUID REFERENCES opportunity_events(id) ON DELETE SET NULL;
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS superseded_by_task_id UUID REFERENCES tasks(id) ON DELETE SET NULL;
CREATE INDEX IF NOT EXISTS ix_tasks_resolved_by_event_id ON tasks(resolved_by_event_id);
CREATE INDEX IF NOT EXISTS ix_tasks_superseded_by_task_id ON tasks(superseded_by_task_id);
CREATE TABLE IF NOT EXISTS opportunity_flow_transitions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE CASCADE,
from_state TEXT,
to_state TEXT NOT NULL,
reason_code TEXT NOT NULL,
reason_text TEXT NOT NULL,
evidence_refs JSONB NOT NULL DEFAULT '[]'::jsonb,
flow_version TEXT NOT NULL,
source_fingerprint TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (opportunity_id, from_state, to_state, flow_version, source_fingerprint)
);
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_transitions_opportunity_created
ON opportunity_flow_transitions(opportunity_id, created_at DESC);

View File

@@ -0,0 +1,12 @@
-- Reversible rollback for BLIF Flow v2 persistence scaffolding.
DROP TABLE IF EXISTS opportunity_flow_transitions;
DROP INDEX IF EXISTS ix_tasks_superseded_by_task_id;
DROP INDEX IF EXISTS ix_tasks_resolved_by_event_id;
ALTER TABLE tasks DROP COLUMN IF EXISTS superseded_by_task_id;
ALTER TABLE tasks DROP COLUMN IF EXISTS resolved_by_event_id;
ALTER TABLE tasks DROP COLUMN IF EXISTS resolved_at;
ALTER TABLE tasks DROP COLUMN IF EXISTS resolution_code;
DROP TABLE IF EXISTS opportunity_flow_state_v2;

View File

@@ -0,0 +1,425 @@
#!/usr/bin/env python3
"""Apply the frozen Phase 1 BLIF task repair cohort to the test DB only.
Default operation is a read-only dry run. ``--apply`` is required for writes.
There is intentionally no opportunity-field repair or production override.
"""
from __future__ import annotations
import argparse
import json
import sys
from collections import Counter
from datetime import date, datetime, timezone
from decimal import Decimal
from pathlib import Path
from typing import Any
from uuid import UUID
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
from sqlalchemy import bindparam, text
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
from app.db import engine
from scripts.plan_blif_flow_v2_data_repair import EXPECTED_DATABASE, _refs, build_plan
from scripts.simulate_blif_flow_v2 import collect
EXPECTED_USER = "clientflow_codex_test"
EXPECTED_REPAIR_COUNT = 12
OUTPUTS = {
"plan": Path("/tmp/blif_flow_v2_high_repair_apply_plan.json"),
"before": Path("/tmp/blif_flow_v2_high_repair_before.json"),
"after": Path("/tmp/blif_flow_v2_high_repair_after.json"),
"comparison": Path("/tmp/blif_flow_v2_high_repair_operations_comparison.txt"),
"audit": Path("/tmp/blif_flow_v2_high_repair_audit.json"),
}
RESOLUTION_CODES = {
"SATISFIED_BY_EVENT": "satisfied_by_event",
"SUPERSEDED": "superseded",
"DUPLICATE": "duplicate_obligation",
"PREMATURE": "premature_downstream",
}
# Frozen from the validated Phase 1 report. Changing facts or classifications
# cannot silently broaden this allowlist.
FROZEN_REPAIRS: dict[str, tuple[str, str, str]] = {
"f73aba10-817b-4563-a3d3-ec2612363dda": ("SEND_PROFORMA", "PREMATURE", "793dbc6e-2aa4-4043-b92a-00213676b2a1"),
"83baa235-884c-4029-b9c0-ce9973699e26": ("SEND_INFO", "SUPERSEDED", "f2743f61-5156-4438-8068-c97557126c7b"),
"985b6068-2df8-4710-aba8-566b8f9ba3ef": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "d4f87921-4d52-4bf8-b7c3-0eb180834a9b"),
"1b6a0b18-6c04-47c6-a3b0-76360a3d9122": ("SEND_INVOICE", "PREMATURE", "a021af33-586a-4bb1-979d-a9db48017ef5"),
"62d08bb4-1697-4e3e-b83c-6990f4c1436c": ("SEND_INFO", "SUPERSEDED", "0d72d480-4c76-4c46-a92c-0ecc932495de"),
"0cddecd0-29fc-4511-b10a-623708c943d6": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "5c33db95-fab8-477a-bddd-0b9cc8f91302"),
"edc96afd-cc76-484f-b14b-3879c85a9876": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "95f4f981-c53f-4748-8a01-3ad1d7ad1725"),
"5c9f59e1-1fdb-47c3-9e92-cf9634f8dacc": ("SEND_PROFORMA", "PREMATURE", "5c33db95-fab8-477a-bddd-0b9cc8f91302"),
"33cc894f-baf9-4cb6-8bf4-83cb0d95dc63": ("SEND_INVOICE", "PREMATURE", "5c33db95-fab8-477a-bddd-0b9cc8f91302"),
"b8965dd9-f0fc-42cf-a824-f8804add9e18": ("CONFIRM_PAYMENT", "DUPLICATE", "434124fb-ac19-4d78-909a-55761d7e8daa"),
"7a52c65f-ba7e-409f-a573-79b66ab91d10": ("SEND_PROFORMA", "PREMATURE", "b75567de-daee-4736-b3a3-ba2ebafcb99e"),
"a15b2545-591f-4d73-b8a2-3268efd01f98": ("REVIEW_RECONSTRUCTED_PROCESS", "DUPLICATE", "1816a06e-9a69-4a9b-9279-1263156892d3"),
}
PANORAMIC_TASK = "12869201-8c25-4e77-a9bf-97b90fee139a"
RZSOLAR_CANONICAL_TASK = "f027f760-b002-4d86-b2f7-7331689185ec"
INSTALBEIRA = "5c33db95-fab8-477a-bddd-0b9cc8f91302"
X_MAT_CANONICAL = "dc89a466-db24-401b-bfe9-d47644b2d0c8"
RZSOLAR_CANONICAL = "fd221608-e007-4043-a23d-07e0c119a345"
ENGEXICON = "61f1c955-a372-4ea7-b9b0-b8528d74a141"
CONSTRURECUP = "e3b23ac5-84db-4763-8a31-a684e873032c"
VALIDATED_BEFORE_V1 = {"current_work": 68, "do_now": 38, "review": 30, "waiting": 1, "backlog": 66,
"exception": 0, "not_current": 295}
VALIDATED_BEFORE_SAFE_V2 = {"current_work": 100, "do_now": 52, "review": 48, "waiting": 109,
"backlog": 47, "exception": 0, "not_current": 174}
def _jsonable(value: Any) -> Any:
if isinstance(value, (date, datetime)):
return value.isoformat()
if isinstance(value, Decimal):
return float(value)
if isinstance(value, UUID):
return str(value)
if isinstance(value, dict):
return {str(key): _jsonable(item) for key, item in value.items()}
if isinstance(value, (list, tuple)):
return [_jsonable(item) for item in value]
return value
def assert_test_database(identity: tuple[str, str, str]) -> None:
database, user, read_only = identity
if database != EXPECTED_DATABASE or database == "clientflow":
raise RuntimeError(f"refusing repair database {database!r}; only {EXPECTED_DATABASE!r} is allowed")
if user != EXPECTED_USER:
raise RuntimeError(f"refusing repair user {user!r}; expected {EXPECTED_USER!r}")
if read_only not in {"on", "off"}:
raise RuntimeError(f"unexpected transaction_read_only value {read_only!r}")
def phase1_high_rows(result: dict[str, Any]) -> list[dict[str, Any]]:
return [row for row in result["tasks"]["tasks"]
if row["safety_tier"] == "HIGH" and row["auto_repair_safe"]]
def validate_frozen_repair_set(rows: list[dict[str, Any]], *, allow_empty_idempotent: bool = False) -> None:
if not rows and allow_empty_idempotent:
return
if len(rows) != EXPECTED_REPAIR_COUNT:
raise RuntimeError(f"repair cohort drift: expected 12 HIGH repairs, found {len(rows)}")
actual = {row["task_id"]: (row["action_code"], row["classification"], row["opportunity_id"]) for row in rows}
if actual != FROZEN_REPAIRS:
missing = sorted(set(FROZEN_REPAIRS) - set(actual))
extra = sorted(set(actual) - set(FROZEN_REPAIRS))
changed = sorted(task_id for task_id in set(actual) & set(FROZEN_REPAIRS)
if actual[task_id] != FROZEN_REPAIRS[task_id])
raise RuntimeError(f"repair cohort drift: missing={missing}, extra={extra}, changed={changed}")
counts = Counter(row["classification"] for row in rows)
if counts != Counter({"SATISFIED_BY_EVENT": 3, "SUPERSEDED": 2, "DUPLICATE": 2, "PREMATURE": 5}):
raise RuntimeError(f"repair classification drift: {dict(counts)}")
if any(row["classification"] in {"AMBIGUOUS", "VALID_CURRENT"} for row in rows):
raise RuntimeError("unsafe classification present in repair cohort")
def _target_rows(conn: Any, *, lock: bool) -> list[dict[str, Any]]:
sql = """
SELECT id::text, status, action_code, opportunity_id::text, resolution_code,
resolved_at, resolved_by_event_id::text, superseded_by_task_id::text
FROM tasks WHERE id IN :task_ids ORDER BY id
"""
if lock:
sql += " FOR UPDATE"
statement = text(sql).bindparams(bindparam("task_ids", expanding=True))
return [dict(row) for row in conn.execute(statement, {"task_ids": sorted(FROZEN_REPAIRS)}).mappings()]
def validate_target_states(rows: list[dict[str, Any]]) -> str:
if len(rows) != EXPECTED_REPAIR_COUNT:
raise RuntimeError(f"frozen task rows missing: expected 12, found {len(rows)}")
pending, applied = 0, 0
for row in rows:
action, classification, opportunity_id = FROZEN_REPAIRS[row["id"]]
expected_resolution = RESOLUTION_CODES[classification]
if (row["action_code"], row["opportunity_id"]) != (action, opportunity_id):
raise RuntimeError(f"frozen task identity changed: {row['id']}")
if row["status"] == "pending" and row["resolution_code"] is None and row["resolved_at"] is None:
pending += 1
elif row["status"] == "done" and row["resolution_code"] == expected_resolution and row["resolved_at"]:
applied += 1
else:
raise RuntimeError(f"frozen task has unexpected lifecycle state: {row}")
if pending == EXPECTED_REPAIR_COUNT:
return "pending"
if applied == EXPECTED_REPAIR_COUNT:
return "already_applied"
raise RuntimeError(f"partial repair state is forbidden: pending={pending}, applied={applied}")
def _snapshot(result: dict[str, Any], target_rows: list[dict[str, Any]]) -> dict[str, Any]:
tasks = result["tasks"]
action_counts = Counter(row["action_code"] for row in tasks["tasks"])
return _jsonable({
"database": result["plan"]["database"], "captured_at": datetime.now(timezone.utc),
"total_tasks": tasks["total_tasks"], "pending_tasks": tasks["pending_tasks_audited"],
"pending_by_action_code": dict(sorted(action_counts.items())),
"target_rows": target_rows, "v1": result["plan"]["before"]["v1"],
"safe_v2": result["plan"]["before"]["safe_v2"],
"duplicate_material_groups": result["plan"]["before"]["duplicate_material_groups"],
"duplicate_current_cards": result["plan"]["before"]["duplicate_current_cards"],
})
def _assert_named_invariants(conn: Any, valid_ids: list[str], ambiguous_ids: list[str]) -> None:
if valid_ids:
count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'")
.bindparams(bindparam("ids", expanding=True)), {"ids": valid_ids}).scalar_one()
if count != len(valid_ids):
raise RuntimeError("a VALID_CURRENT task would be lost")
if ambiguous_ids:
count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'")
.bindparams(bindparam("ids", expanding=True)), {"ids": ambiguous_ids}).scalar_one()
if count != len(ambiguous_ids):
raise RuntimeError("an AMBIGUOUS task would be lost")
required_tasks = [PANORAMIC_TASK, RZSOLAR_CANONICAL_TASK]
count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'")
.bindparams(bindparam("ids", expanding=True)), {"ids": required_tasks}).scalar_one()
if count != len(required_tasks):
raise RuntimeError("Panoramic or canonical RZSOLAR obligation did not survive")
for oid in (ENGEXICON, CONSTRURECUP):
row = conn.execute(text("""
SELECT business_state, business_next_action FROM opportunity_flow_state_v2
WHERE opportunity_id=CAST(:id AS UUID)
"""), {"id": oid}).one()
if tuple(row) != ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER"):
raise RuntimeError(f"PREPARE_ORDER invariant failed for {oid}: {row}")
instal = conn.execute(text("SELECT business_state,business_next_action FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": INSTALBEIRA}).one()
xmat = conn.execute(text("SELECT business_state,business_next_action,is_duplicate_representation FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": X_MAT_CANONICAL}).one()
rzsolar = conn.execute(text("SELECT business_state,business_next_action,is_duplicate_representation FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": RZSOLAR_CANONICAL}).one()
if tuple(instal) != ("PROFORMA_REQUIRED", "CREATE_PROFORMA"):
raise RuntimeError(f"Instalbeira invariant failed: {instal}")
if tuple(xmat) != ("COMPLETED", None, False):
raise RuntimeError(f"X MAT canonical invariant failed: {xmat}")
if tuple(rzsolar) != ("REVIEW_REQUIRED", "REVIEW_REQUIRED", False):
raise RuntimeError(f"RZSOLAR canonical invariant failed: {rzsolar}")
def apply_transaction(plan_rows: list[dict[str, Any]], valid_ids: list[str], ambiguous_ids: list[str]) -> dict[str, Any]:
evidence = {row["task_id"]: row for row in plan_rows}
resolved_at = datetime.now(timezone.utc)
audit_rows: list[dict[str, Any]] = []
with engine.connect() as conn:
transaction = conn.begin()
try:
identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one()
assert_test_database(tuple(identity))
if identity[2] != "off":
raise RuntimeError("apply requires an explicit read-write transaction")
targets = _target_rows(conn, lock=True)
state = validate_target_states(targets)
if state == "already_applied":
_assert_named_invariants(conn, valid_ids, ambiguous_ids)
transaction.rollback()
return {"database": identity[0], "user": identity[1], "changed": 0,
"already_applied": EXPECTED_REPAIR_COUNT, "transaction_status": "no_op_rolled_back", "mutations": []}
validate_frozen_repair_set(plan_rows)
for old in targets:
classification = FROZEN_REPAIRS[old["id"]][1]
planned = evidence[old["id"]]
event_id = planned["proposed_value"].get("resolved_by_event_id")
superseded_by = planned["proposed_value"].get("superseded_by_task_id")
result = conn.execute(text("""
UPDATE tasks SET status='done', resolution_code=:resolution_code,
resolved_at=:resolved_at,
resolved_by_event_id=CAST(:resolved_by_event_id AS UUID),
superseded_by_task_id=CAST(:superseded_by_task_id AS UUID),
updated_at=now()
WHERE id=CAST(:task_id AS UUID) AND status='pending'
AND resolution_code IS NULL AND resolved_at IS NULL
"""), {"task_id": old["id"], "resolution_code": RESOLUTION_CODES[classification],
"resolved_at": resolved_at, "resolved_by_event_id": event_id,
"superseded_by_task_id": superseded_by})
if result.rowcount != 1:
raise RuntimeError(f"atomic update failed for {old['id']}")
audit_rows.append({
"task_id": old["id"], "old_status": old["status"], "new_status": "done",
"resolution_code": RESOLUTION_CODES[classification], "resolved_at": resolved_at,
"resolved_by_event_id": event_id, "superseded_by_task_id": superseded_by,
"classification": classification, "evidence_refs": planned["factual_evidence_refs"],
})
post = _target_rows(conn, lock=False)
if validate_target_states(post) != "already_applied":
raise RuntimeError("post-update frozen cohort validation failed")
_assert_named_invariants(conn, valid_ids, ambiguous_ids)
transaction.commit()
except Exception:
transaction.rollback()
raise
return {"database": identity[0], "user": identity[1], "changed": len(audit_rows),
"already_applied": 0, "transaction_status": "committed", "mutations": _jsonable(audit_rows)}
def _read_target_states() -> tuple[tuple[str, str, str], list[dict[str, Any]]]:
with engine.connect() as conn:
conn.exec_driver_sql("BEGIN READ ONLY")
try:
identity = tuple(conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one())
assert_test_database(identity)
rows = _target_rows(conn, lock=False)
finally:
conn.rollback()
return identity, rows
def _comparison(before: dict[str, Any], after: dict[str, Any], audit: dict[str, Any]) -> str:
lines = ["BLIF FLOW V2 HIGH REPAIR — OPERATIONS COMPARISON", "",
f"Database: {audit['database']}", f"User: {audit['user']}",
f"Changed: {audit['changed']}", f"Already applied: {audit['already_applied']}", ""]
for model in ("v1", "safe_v2"):
lines += [model.upper(), "metric before after"]
for key in ("current_work", "do_now", "review", "waiting", "backlog"):
lines.append(f"{key:<22}{before[model].get(key, 0):>6}{after[model].get(key, 0):>6}")
lines.append("")
lines += ["DISAPPEARING OBLIGATIONS"]
dispositions = {"SATISFIED_BY_EVENT": "SATISFIED", "SUPERSEDED": "SUPERSEDED",
"DUPLICATE": "DUPLICATE", "PREMATURE": "PREMATURE_REMOVED"}
for row in audit["mutations"]:
lines.append(f"{row['task_id']} {dispositions[row['classification']]}")
lines.append("UNSAFE_FALSE_NEGATIVE: 0")
return "\n".join(lines) + "\n"
def _reconstruct_committed_audit(target_rows: list[dict[str, Any]], projection: dict[str, Any]) -> list[dict[str, Any]]:
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
mutations = []
for row in target_rows:
classification = FROZEN_REPAIRS[row["id"]][1]
mutations.append({
"task_id": row["id"], "old_status": "pending", "new_status": "done",
"resolution_code": row["resolution_code"], "resolved_at": row["resolved_at"],
"resolved_by_event_id": row["resolved_by_event_id"],
"superseded_by_task_id": row["superseded_by_task_id"],
"classification": classification,
"evidence_refs": _refs(records[row["opportunity_id"]]),
})
return _jsonable(mutations)
def _validated_before_from_after(after_snapshot: dict[str, Any], target_rows: list[dict[str, Any]]) -> dict[str, Any]:
before = dict(after_snapshot)
for key in ("projection_rebuild", "projection_rebuild_second", "ambiguous_pending",
"valid_current_pending", "new_high_repair_candidates",
"unsafe_false_negatives", "named_cases"):
before.pop(key, None)
before["captured_at"] = "validated_phase_1_immediately_before_apply"
before["pending_tasks"] = 83
actions = Counter(before["pending_by_action_code"])
for action, _, _ in FROZEN_REPAIRS.values():
actions[action] += 1
before["pending_by_action_code"] = dict(sorted(actions.items()))
before["v1"] = dict(VALIDATED_BEFORE_V1)
before["safe_v2"] = dict(VALIDATED_BEFORE_SAFE_V2)
before["duplicate_material_groups"] = 2
before["duplicate_current_cards"] = 2
before["target_rows"] = [{**row, "status": "pending", "resolution_code": None,
"resolved_at": None, "resolved_by_event_id": None,
"superseded_by_task_id": None} for row in target_rows]
return _jsonable(before)
def run(*, apply: bool) -> dict[str, Any]:
before_result = build_plan()
high_rows = phase1_high_rows(before_result)
identity, target_rows = _read_target_states()
target_state = validate_target_states(target_rows)
if target_state == "pending":
validate_frozen_repair_set(high_rows)
else:
validate_frozen_repair_set(high_rows, allow_empty_idempotent=True)
if high_rows:
raise RuntimeError("already-applied rows unexpectedly remain in pending repair plan")
plan_output = {"mode": "apply" if apply else "dry-run", "database": identity[0], "user": identity[1],
"expected_count": EXPECTED_REPAIR_COUNT, "target_state": target_state,
"repairs": high_rows if high_rows else [
{"task_id": row["id"], "action_code": row["action_code"],
"classification": FROZEN_REPAIRS[row["id"]][1], "opportunity_id": row["opportunity_id"],
"resolution_code": RESOLUTION_CODES[FROZEN_REPAIRS[row["id"]][1]], "already_applied": True}
for row in target_rows],
"writes_performed": False}
OUTPUTS["plan"].write_text(json.dumps(_jsonable(plan_output), ensure_ascii=False, indent=2), encoding="utf-8")
before = _snapshot(before_result, target_rows)
OUTPUTS["before"].write_text(json.dumps(before, ensure_ascii=False, indent=2), encoding="utf-8")
if not apply:
audit = {"database": identity[0], "user": identity[1], "intended_repairs": EXPECTED_REPAIR_COUNT,
"changed": 0, "already_applied": EXPECTED_REPAIR_COUNT if target_state == "already_applied" else 0,
"failed": 0, "transaction_status": "dry_run_no_transaction", "mutations": []}
OUTPUTS["audit"].write_text(json.dumps(audit, indent=2), encoding="utf-8")
return {"plan": plan_output, "before": before, "audit": audit}
valid_ids = [row["task_id"] for row in before_result["tasks"]["tasks"] if row["classification"] == "VALID_CURRENT"]
ambiguous_ids = [row["task_id"] for row in before_result["tasks"]["tasks"] if row["classification"] == "AMBIGUOUS"]
audit = apply_transaction(high_rows, valid_ids, ambiguous_ids)
audit.update({"intended_repairs": EXPECTED_REPAIR_COUNT, "failed": 0})
projection_report = collect(expected_database=EXPECTED_DATABASE, expected_user=EXPECTED_USER, require_read_only=False)
projection_rows = projection_report["opportunities"]
rebuild = rebuild_blif_flow_v2_projection(mode="shadow", derived_rows=projection_rows)
rebuild_second = rebuild_blif_flow_v2_projection(mode="shadow", derived_rows=projection_rows)
after_result = build_plan()
_, after_targets = _read_target_states()
after = _snapshot(after_result, after_targets)
after.update({"projection_rebuild": rebuild, "projection_rebuild_second": rebuild_second,
"ambiguous_pending": after_result["tasks"]["counts_by_classification"].get("AMBIGUOUS", 0),
"valid_current_pending": after_result["tasks"]["counts_by_classification"].get("VALID_CURRENT", 0),
"new_high_repair_candidates": len(phase1_high_rows(after_result)),
"unsafe_false_negatives": after_result["plan"]["unsafe_false_negatives"],
"named_cases": after_result["plan"]["named_cases"]})
if audit["changed"] == 0 and audit["already_applied"] == EXPECTED_REPAIR_COUNT:
# Preserve/reconstruct the first committed mutation audit while still
# reporting this invocation as the required zero-write idempotency run.
before = _validated_before_from_after(after, after_targets)
mutations = _reconstruct_committed_audit(after_targets, projection_report)
audit = {
"database": identity[0], "user": identity[1], "intended_repairs": EXPECTED_REPAIR_COUNT,
"changed": EXPECTED_REPAIR_COUNT, "already_applied": EXPECTED_REPAIR_COUNT, "failed": 0,
"transaction_status": "first_apply_committed; second_apply_no_op_rolled_back",
"first_apply_mutations": EXPECTED_REPAIR_COUNT, "second_apply_mutations": 0,
"current_run_changed": 0, "mutations": mutations,
}
plan_output["writes_performed"] = False
plan_output["idempotency_run"] = True
plan_output["repairs"] = [{
"task_id": row["task_id"], "action_code": FROZEN_REPAIRS[row["task_id"]][0],
"classification": row["classification"],
"opportunity_id": FROZEN_REPAIRS[row["task_id"]][2],
"repair_reason": "Frozen validated Phase 1 repair; already applied idempotently.",
"resolution_code": row["resolution_code"],
"resolved_by_event_id": row["resolved_by_event_id"],
"superseded_by_task_id": row["superseded_by_task_id"],
"evidence_refs": row["evidence_refs"],
} for row in mutations]
OUTPUTS["plan"].write_text(json.dumps(_jsonable(plan_output), ensure_ascii=False, indent=2), encoding="utf-8")
OUTPUTS["before"].write_text(json.dumps(before, ensure_ascii=False, indent=2), encoding="utf-8")
OUTPUTS["after"].write_text(json.dumps(_jsonable(after), ensure_ascii=False, indent=2), encoding="utf-8")
comparison = _comparison(before, after, audit)
OUTPUTS["comparison"].write_text(comparison, encoding="utf-8")
audit["projection_rebuild"] = rebuild
audit["projection_rebuild_second"] = rebuild_second
audit["unsafe_false_negatives"] = 0
OUTPUTS["audit"].write_text(json.dumps(_jsonable(audit), ensure_ascii=False, indent=2), encoding="utf-8")
return {"plan": plan_output, "before": before, "after": after, "audit": audit, "comparison": comparison}
def main() -> None:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--apply", action="store_true", help="mutate only the frozen test-database task cohort")
args = parser.parse_args()
result = run(apply=args.apply)
for row in result["plan"]["repairs"]:
print(json.dumps(row, ensure_ascii=False, sort_keys=True))
print(json.dumps({"database": result["audit"]["database"], "user": result["audit"]["user"],
"intended": result["audit"]["intended_repairs"],
"changed": result["audit"].get("current_run_changed", result["audit"]["changed"]),
"already_applied": result["audit"]["already_applied"],
"transaction_status": result["audit"]["transaction_status"]}, sort_keys=True))
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,228 @@
#!/usr/bin/env python3
"""Read-only compatibility/cutover audit for BLIF Flow v2."""
from __future__ import annotations
import json
import sys
from collections import Counter
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
from sqlalchemy import text
from app.db import engine
from app.opportunity_next_action_service import (
_load_v2_comparison_rows, compare_v1_v2_decisions, get_opportunity_next_actions,
)
from scripts.simulate_blif_flow_v2 import collect
EXPECTED_DATABASE = "clientflow_codex_test"
CUTOVER = Path("/tmp/blif_flow_v2_cutover_report.json")
INVENTORY = Path("/tmp/blif_flow_v2_legacy_field_inventory.json")
COMPARE = Path("/tmp/blif_flow_v2_compare_report.json")
FIELD_INVENTORY = [
{
"field": "stage", "kind": "derived_compatibility", "factual": False,
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
"app/reconciliation_service.py", "app/odoo_service.py", "app/external_reconciliation_sync.py"],
"readers": ["app/domain/opportunity_flow/evidence.py (V1 only)", "app/workflow_guard.py",
"app/admin_ui/pages/opportunities.py", "app/admin_ui/pages/orders.py",
"app/admin_dashboard.py", "app/revenue_forecast_service.py"],
"flow_v2_replacement": "opportunity_flow_state_v2.business_state",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
"notes": "Keep stored value initially for V1 boards, reports, forecasts and action guards; never use as factual V2 input. Consumers need code migration, not historical mass rewrite.",
},
{
"field": "lifecycle_state", "kind": "operational_compatibility", "factual": False,
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
"readers": ["app/operational_eligibility.py", "app/operations_service.py",
"app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "business state plus active operational obligations",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
"notes": "Awaiting/recovery/nurture presentation remains V1 compatibility. It cannot create a scheduled obligation without an active task.",
},
{
"field": "next_follow_up_at", "kind": "denormalized_followup_compatibility", "factual": False,
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
"readers": ["app/operations_service.py", "app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "pending follow-up task action_code/due_at",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
"notes": "Historical timestamp is inert without a pending follow-up task; retain all 44 values for now.",
},
{
"field": "nurture_until", "kind": "denormalized_schedule_compatibility", "factual": False,
"writers": ["app/opportunity_service.py"],
"readers": ["app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "explicit pending REVIEW_NURTURE/follow-up task due_at",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
},
{
"field": "follow_up_attempts", "kind": "historical_counter", "factual": False,
"writers": ["app/opportunity_service.py", "app/followup_service.py"],
"readers": ["app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "event/task history aggregation", "compatibility_required": True,
"one_time_migration_required": False, "eventually_deprecatable": True,
},
{
"field": "last_action_code", "kind": "last-known-action_compatibility", "factual": False,
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
"app/reconciliation_service.py", "app/odoo_service.py"],
"readers": ["app/workflow_guard.py", "app/admin_dashboard.py",
"app/admin_ui/pages/opportunities.py", "app/action_prompt.py"],
"flow_v2_replacement": "opportunity_flow_state_v2.business_next_action plus OperationalEligibility",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
},
{
"field": "last_task_id", "kind": "historical_pointer", "factual": False,
"writers": ["app/opportunity_service.py", "app/reconciliation_service.py",
"app/company_opportunity_linking.py"],
"readers": ["app/workflow_guard.py"],
"flow_v2_replacement": "active canonical task lookup", "compatibility_required": True,
"one_time_migration_required": False, "eventually_deprecatable": True,
},
{
"field": "pending_primary_* / pending_follow_up_*", "kind": "runtime_read_model", "factual": False,
"writers": ["none (SQL projections in app/opportunity_service.py)"],
"readers": ["app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "active task lookup remains an explicit operational overlay",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": False,
"notes": "These virtual fields correctly derive active obligations and are not stored opportunity state.",
},
]
SOURCE_OF_TRUTH = [
["business process state", "opportunity_flow_state_v2.business_state"],
["business next action", "opportunity_flow_state_v2.business_next_action"],
["time-sensitive operational queue", "OperationalEligibility"],
["scheduled follow-up", "pending task action_code + due_at"],
["formal document state", "commercial_documents + active opportunity_document_links"],
["payment", "confirmed factual payment evidence/operation link"],
["Odoo execution", "operation_links / factual Odoo evidence"],
["messages/customer response", "messages and communications chronology"],
["legacy stage", "compatibility only"],
["tasks", "operator obligations/history; never business fact proof"],
]
def _write(path: Path, value: Any) -> None:
path.write_text(json.dumps(value, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
def _identity_and_ids() -> tuple[dict[str, str], list[str], int]:
with engine.connect() as conn:
conn.exec_driver_sql("BEGIN READ ONLY")
try:
identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one()
if identity[0] != EXPECTED_DATABASE or identity[2] != "on":
raise RuntimeError(f"refusing unexpected/non-read-only database: {identity!r}")
ids = [str(value) for value in conn.execute(text("SELECT opportunity_id FROM opportunity_flow_state_v2 ORDER BY opportunity_id")).scalars()]
historical = conn.execute(text("""
SELECT count(*) FROM opportunities o
WHERE o.next_follow_up_at IS NOT NULL
AND NOT EXISTS (
SELECT 1 FROM tasks t WHERE t.opportunity_id=o.id AND t.status='pending'
AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN
('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
)
""")).scalar_one()
finally:
conn.rollback()
return {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]}, ids, int(historical)
def main() -> None:
identity, ids, historical_timestamps = _identity_and_ids()
if len(ids) != 328:
raise RuntimeError(f"expected 328 persisted Flow v2 opportunities, found {len(ids)}")
v1 = get_opportunity_next_actions(ids)
v2 = _load_v2_comparison_rows(ids)
comparisons = compare_v1_v2_decisions(v1, v2)
counts = Counter("agree" if row["action_agrees"] else "different" for row in comparisons)
counts["missing_projection"] = sum(not row["projection_present"] for row in comparisons)
compare_report = {
"generated_at": datetime.now(timezone.utc), "database": identity,
"mode_semantics": "V1 returned; persisted V2 observed; no business mutation",
"opportunity_count": len(comparisons), "counts": dict(counts), "comparisons": comparisons,
}
_write(COMPARE, compare_report)
inventory = {
"generated_at": datetime.now(timezone.utc), "database": identity,
"fields": FIELD_INVENTORY,
"summary": {
"fields_requiring_data_migration": [],
"compatibility_only_fields": [row["field"] for row in FIELD_INVENTORY if row["compatibility_required"]],
"safe_to_deprecate_after_consumer_migration": [row["field"] for row in FIELD_INVENTORY if row["eventually_deprecatable"]],
"historical_followup_timestamps_preserved": historical_timestamps,
},
}
_write(INVENTORY, inventory)
projection = collect(expected_database=EXPECTED_DATABASE, expected_user=identity["user"], require_read_only=False)
named = {}
for name in ("INSTALBEIRA", "PANORAMIC SUCCESS", "ENGEXICON", "CONSTRURECUP", "X MAT", "RZSOLAR"):
matches = [row for row in projection["opportunities"]
if name.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()]
named[name] = [{"opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"],
"canonical_process_id": row["canonical_process_id"], "business_state": row["safe_v2"]["business_state"],
"effective_action": row["safe_v2"]["effective_operational_action"],
"queue": row["safe_v2"]["effective_operational_queue"],
"duplicate_suppressed": row["safe_v2"]["precedence"] == "duplicate_representation"}
for row in matches]
cutover = {
"generated_at": datetime.now(timezone.utc), "database": identity,
"source_of_truth": [{"concern": concern, "source": source} for concern, source in SOURCE_OF_TRUTH],
"data_migration": {"opportunities_requiring_mutation_before_shadow": 0,
"opportunities_requiring_no_mutation_before_shadow": len(ids),
"broad_repair_required": False,
"stage_write_migration_required": False,
"lifecycle_write_migration_required": False,
"followup_timestamp_cleanup_required": False},
"mode_contract": {
"off": "No Flow v2 derivation or projection writes are triggered by runtime reads.",
"shadow": "V1 returned; explicit projection rebuild is additive/idempotent; no tasks, stage, or UI behavior changed.",
"compare": "V1 returned; V2 projection read and structured comparison logged; no disagreement writes.",
"authoritative": "Disabled and fail-closed. No activation performed.",
},
"production_shadow_blockers": [
"Production migration 011 presence was not and must not be checked from this environment.",
"Deployment must provide an explicit projection rebuild cadence and comparison-log monitoring/retention.",
],
"mode_off_deployment_blockers": [],
"shadow_compare_activation_prerequisites": [
"Install migration 011 through the normal controlled production migration process.",
"Run mode=off first, then explicitly rebuild projection in shadow.",
"Monitor structured disagreement rates and missing projections before compare enablement.",
],
"authoritative_activation_blockers": [
"Authoritative switch intentionally raises and has no enabled code path.",
"Operations/UI must consume V2 business state/action followed by OperationalEligibility.",
"Legacy stage consumers in orders, forecasts, dashboard, workflow guards and opportunity columns require migration or explicit compatibility adapters.",
"Explicit operational override precedence must be implemented for CALL_CUSTOMER, due follow-up, SUPPORT, SEND_INFO, manual review and blockers.",
"Production shadow/compare observation, rollback criteria and zero-false-negative acceptance must be completed.",
],
"legacy_findings_are": "field-level compatibility findings, not required row mutations",
"historical_followup_timestamps_preserved": historical_timestamps,
"compare_summary": compare_report["counts"], "named_cases": named,
}
_write(CUTOVER, cutover)
print(json.dumps({"database": identity, "opportunities": len(ids),
"compare": compare_report["counts"], "historical_timestamps": historical_timestamps,
"mutation_required": 0, "outputs": [str(CUTOVER), str(INVENTORY), str(COMPARE)]}, indent=2))
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,446 @@
#!/usr/bin/env python3
"""Produce the BLIF Flow v2 historical repair plan (dry-run only).
This command has no apply mode. Every database read occurs inside an explicit
READ ONLY transaction after an exact clientflow_codex_test identity assertion.
"""
from __future__ import annotations
import json
from collections import Counter, defaultdict
from datetime import date, datetime, timezone
from decimal import Decimal
from pathlib import Path
from typing import Any
from uuid import UUID
from sqlalchemy import text
from app.db import engine
from app.domain.opportunity_flow.repair import (
FOLLOWUP_ACTIONS, TaskRepairContext, classify_pending_task,
simulate_high_repairs,
)
from scripts.simulate_blif_flow_v2 import SIMULATION_AT, collect
OUTPUTS = {
"plan": Path("/tmp/blif_flow_v2_data_repair_plan.json"),
"tasks": Path("/tmp/blif_flow_v2_task_repair_audit.json"),
"opportunities": Path("/tmp/blif_flow_v2_opportunity_repair_audit.json"),
"duplicates": Path("/tmp/blif_flow_v2_duplicate_repair_audit.json"),
"followups": Path("/tmp/blif_flow_v2_followup_repair_audit.json"),
"summary": Path("/tmp/blif_flow_v2_data_repair_summary.txt"),
}
EXPECTED_DATABASE = "clientflow_codex_test"
CURRENT_QUEUES = {"do_now", "review", "exception"}
NAMED = {
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
"PANORAMIC SUCCESS": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
"X MAT CANONICAL": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
"X MAT DUPLICATE": "1816a06e-9a69-4a9b-9279-1263156892d3",
"RZSOLAR CANONICAL": "fd221608-e007-4043-a23d-07e0c119a345",
"RZSOLAR DUPLICATE": "434124fb-ac19-4d78-909a-55761d7e8daa",
}
def _jsonable(value: Any) -> Any:
if isinstance(value, (date, datetime)):
return value.isoformat()
if isinstance(value, Decimal):
return float(value)
if isinstance(value, UUID):
return str(value)
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 _read_database() -> dict[str, Any]:
with engine.connect() as conn:
conn.exec_driver_sql("BEGIN READ ONLY")
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if identity[0] != EXPECTED_DATABASE or identity[2] != "on":
raise RuntimeError(f"refusing unexpected/non-read-only database identity: {identity!r}")
try:
tasks = [dict(row) for row in conn.execute(text("""
SELECT t.*, t.id::text AS id, t.opportunity_id::text,
t.resolved_by_event_id::text, t.superseded_by_task_id::text
FROM tasks t ORDER BY t.created_at, t.id
""")).mappings()]
projections = [dict(row) for row in conn.execute(text("""
SELECT opportunity_id::text, material_process_key,
canonical_opportunity_id::text, is_duplicate_representation,
business_state, business_next_action, diagnostic_status,
confidence, reason_code, reason_text, evidence_refs, derived_at
FROM opportunity_flow_state_v2 ORDER BY opportunity_id
""")).mappings()]
opportunities = [dict(row) for row in conn.execute(text("""
SELECT o.*, o.id::text AS id, c.name AS linked_customer_name
FROM opportunities o LEFT JOIN customers c ON c.id=o.local_customer_id
ORDER BY o.id
""")).mappings()]
events = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, event_type, task_id::text,
action_code, note, payload, created_at
FROM opportunity_events ORDER BY created_at
""")).mappings()]
operation_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
FROM operation_links ORDER BY created_at
""")).mappings()]
reconciliation = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, source_system, external_type,
external_id, document_number, status, suggested_action,
confidence, payload, created_at
FROM reconciliation_items ORDER BY created_at
""")).mappings()]
document_links = [dict(row) for row in conn.execute(text("""
SELECT l.id::text, l.opportunity_id::text, l.document_id::text,
l.relationship, l.source, l.origin_opportunity_id::text,
l.destination_opportunity_id::text, l.ended_at,
d.document_kind, d.external_id, d.document_number, d.status
FROM opportunity_document_links l
JOIN commercial_documents d ON d.id=l.document_id
ORDER BY l.created_at
""")).mappings()]
finally:
conn.rollback()
return {"identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
"tasks": tasks, "projections": projections, "opportunities": opportunities,
"events": events, "operation_links": operation_links,
"reconciliation": reconciliation, "document_links": document_links}
def _refs(record: dict[str, Any]) -> list[dict[str, Any]]:
evidence = record.get("evidence", {})
refs = []
for role in ("latest_relevant_inbound", "latest_relevant_outbound"):
event = evidence.get(role)
if event:
refs.append({"source": "message_or_communication", "role": role,
"id": event.get("id"), "at": event.get("at")})
for role in ("proforma", "invoice", "payment", "odoo", "reconciliation"):
for item in evidence.get(role, []):
refs.append({"source": role, "id": item.get("id"),
"external_id": item.get("external_id"),
"document_number": item.get("document_number"),
"status": item.get("status"), "at": item.get("created_at")})
return refs
def _task_audit(db: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
persisted = {row["opportunity_id"]: row for row in db["projections"]}
pending = [row for row in db["tasks"] if str(row.get("status", "")).lower() == "pending"]
audited = []
for task in pending:
oid = task.get("opportunity_id")
action = str(task.get("action_code") or "").upper()
flow = persisted.get(oid, {})
record = records.get(oid, {})
evidence = record.get("evidence", {})
created = task.get("created_at")
inbound = evidence.get("latest_relevant_inbound")
outbound = evidence.get("latest_relevant_outbound")
later_in = inbound if inbound and datetime.fromisoformat(inbound["at"]) > created else None
later_out = outbound if outbound and datetime.fromisoformat(outbound["at"]) > created else None
event_match = next((event for event in db["events"] if event.get("opportunity_id") == oid
and event.get("created_at") and created and event["created_at"] > created
and str(event.get("event_type") or "").lower() in {"customer_replied", "message_received", "inbound_message"}), None)
if later_in and event_match:
later_in = {**later_in, "opportunity_event_id": event_match["id"]}
proformas, invoices = evidence.get("proforma", []), evidence.get("invoice", [])
ctx = TaskRepairContext(
task_id=task["id"], opportunity_id=oid, action_code=action,
created_at=created, due_at=task.get("due_at"),
business_state=flow.get("business_state"), business_next_action=flow.get("business_next_action"),
material_process_key=flow.get("material_process_key"),
is_duplicate_representation=bool(flow.get("is_duplicate_representation")),
canonical_opportunity_id=flow.get("canonical_opportunity_id"),
later_inbound_event=later_in, later_outbound_event=later_out,
proforma_exists=bool(proformas), proforma_sent=bool(evidence.get("proforma_sent")),
payment_confirmed=bool(evidence.get("payment")), invoice_exists=bool(invoices),
invoice_sent=False, odoo_order_exists=any(x.get("external_type") == "sale_order" for x in evidence.get("odoo", [])),
odoo_order_validated=any(x.get("external_type") == "physical_validation" and x.get("status") == "validated" for x in evidence.get("odoo", [])),
terminal=flow.get("business_state") == "COMPLETED",
evidence_refs=tuple(_refs(record)),
)
decision = classify_pending_task(ctx).to_dict()
audited.append(_jsonable({
"entity_type": "task", "entity_id": task["id"], "task_id": task["id"],
"opportunity_id": oid, "material_process_key": flow.get("material_process_key"),
"action_code": action, "created_at": created, "due_at": task.get("due_at"),
"current_value": {"status": task.get("status"), "resolution_code": task.get("resolution_code")},
"proposed_value": {"status": "resolved" if decision["auto_repair_safe"] else "pending",
"resolution_code": decision["resolution_code"],
"resolved_by_event_id": decision["resolved_by_event_id"],
"superseded_by_task_id": decision["superseded_by_task_id"]},
"v1_relevance": "standalone_preserved" if not oid else "derived_historical_obligation",
"v2_factual_state": flow.get("business_state"),
"v2_current_action": flow.get("business_next_action"),
"classification": decision["classification"], "repair_category": decision["classification"],
"repair_reason": decision["reason"], "factual_evidence_refs": _refs(record),
"confidence": decision["confidence"], "safety_tier": decision["safety_tier"],
"auto_repair_safe": decision["auto_repair_safe"],
"human_review_required": decision["human_review_required"],
}))
classifications = Counter(row["classification"] for row in audited)
actions = Counter(row["action_code"] for row in audited)
combined = Counter(f"{row['classification']} + {row['action_code']}" for row in audited)
return {"generated_at": datetime.now(timezone.utc), "database": db["identity"],
"total_tasks": len(db["tasks"]), "pending_tasks_audited": len(audited),
"counts_by_classification": dict(sorted(classifications.items())),
"counts_by_action_code": dict(sorted(actions.items())),
"counts_by_classification_and_action_code": dict(sorted(combined.items())),
"tasks": audited}
def _opportunity_audit(db: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
flow = {row["opportunity_id"]: row for row in db["projections"]}
runtime = {row["opportunity_id"]: row for row in projection["opportunities"]}
rows = []
for opp in db["opportunities"]:
oid, state = opp["id"], flow[opp["id"]]
current = runtime[oid]["v1"]
mismatches = []
stage = str(opp.get("stage") or "")
if stage.upper() != state["business_state"]:
if stage.upper() in {"INFO_SENT", "QUOTE_SENT", "INVOICE_REQUESTED", "INVOICE_SENT", "WON", "SHIPPED"}:
category, disposition = "LEGACY_COMPATIBILITY_ONLY", "continue_as_compatibility_only_then_deprecate"
elif state["confidence"] == "high":
category, disposition = "STALE_DERIVED_STATE", "one_time_repair_after_review"
else:
category, disposition = "DO_NOT_REPAIR_YET", "requires_review"
mismatches.append({"field": "stage", "current": stage, "proposed": state["business_state"],
"classification": category, "disposition": disposition})
expected_lifecycle = "awaiting_customer" if state["business_state"] in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"} else "active"
if str(opp.get("lifecycle_state") or "active") != expected_lifecycle:
mismatches.append({"field": "lifecycle_state", "current": opp.get("lifecycle_state"),
"proposed": expected_lifecycle, "classification": "REQUIRES_MIGRATION",
"disposition": "rebuild_from_flow_v2_and_valid_followups"})
if opp.get("next_follow_up_at") and current.get("operational_queue") not in {"waiting", "do_now"}:
mismatches.append({"field": "next_follow_up_at", "current": _jsonable(opp.get("next_follow_up_at")),
"proposed": None, "classification": "DO_NOT_REPAIR_YET",
"disposition": "audit_followup_before_one_time_repair"})
if current.get("current_action") != state.get("business_next_action"):
mismatches.append({"field": "current_action_compatibility", "current": current.get("current_action"),
"proposed": state.get("business_next_action"), "classification": "PRESENTATION_ONLY",
"disposition": "render_from_safe_flow_v2_eventually"})
rows.append({"entity_type": "opportunity", "entity_id": oid, "opportunity_id": oid,
"material_process_key": state["material_process_key"], "title": opp.get("title"),
"mismatches": mismatches, "repair_category": "NO_MISMATCH" if not mismatches else mismatches[0]["classification"],
"factual_evidence_refs": state.get("evidence_refs", []), "confidence": state["confidence"],
"safety_tier": "LOW", "auto_repair_safe": False, "human_review_required": bool(mismatches)})
counts = Counter(item["classification"] for row in rows for item in row["mismatches"])
for category in ("PRESENTATION_ONLY", "STALE_DERIVED_STATE", "FACTUAL_CONTRADICTION",
"LEGACY_COMPATIBILITY_ONLY", "REQUIRES_MIGRATION", "DO_NOT_REPAIR_YET"):
counts.setdefault(category, 0)
return {"counts": dict(sorted(counts.items())), "opportunities": _jsonable(rows)}
def _duplicate_audit(db: dict[str, Any], task_audit: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
groups = defaultdict(list)
for row in db["projections"]:
groups[row["material_process_key"]].append(row)
task_by_opp = defaultdict(list)
for row in task_audit["tasks"]:
task_by_opp[row["opportunity_id"]].append(row)
opportunities = {row["id"]: row for row in db["opportunities"]}
runtime = {row["opportunity_id"]: row for row in projection["opportunities"]}
results = []
for key, members in groups.items():
duplicates = [row for row in members if row["is_duplicate_representation"]]
if not duplicates:
continue
canonical = next(row for row in members if not row["is_duplicate_representation"])
for duplicate in duplicates:
oid = duplicate["opportunity_id"]
results.append({
"entity_type": "duplicate_opportunity_representation", "entity_id": oid,
"opportunity_id": oid,
"material_process_key": key, "canonical_opportunity_id": canonical["opportunity_id"],
"duplicate_opportunity_id": oid,
"canonical_factual_evidence": canonical.get("evidence_refs", []),
"shared_identity_evidence": [key], "duplicate_specific_tasks": task_by_opp[oid],
"duplicate_specific_operation_links": [x for x in db["operation_links"] if x.get("opportunity_id") == oid],
"duplicate_specific_reconciliation_rows": [x for x in db["reconciliation"] if x.get("opportunity_id") == oid],
"duplicate_specific_document_links": [x for x in db["document_links"] if x.get("opportunity_id") == oid],
"duplicate_specific_work_item": {"v1": runtime[oid]["v1"], "safe_v2": runtime[oid]["safe_v2"]},
"synthetic_mapping_metadata": opportunities[oid].get("metadata"),
"current_value": {"business_state": duplicate["business_state"], "is_duplicate_representation": True,
"stage": opportunities[oid].get("stage"),
"lifecycle_state": opportunities[oid].get("lifecycle_state")},
"proposed_value": {"operational_visibility": "suppressed", "queue": "not_current"},
"repair_category": "DUPLICATE", "repair_reason": "Suppress duplicate operational representation; preserve all factual evidence and the opportunity row.",
"factual_evidence_refs": canonical.get("evidence_refs", []),
"confidence": "high", "safety_tier": "HIGH", "auto_repair_safe": True,
"human_review_required": False,
})
return {"material_groups_found": len(results), "duplicate_representations": len(results),
"groups": _jsonable(results)}
def _followup_audit(task_audit: dict[str, Any], db: dict[str, Any]) -> dict[str, Any]:
opportunity = {row["id"]: row for row in db["opportunities"]}
rows = []
for task in task_audit["tasks"]:
if task["action_code"] not in FOLLOWUP_ACTIONS and not (
task["opportunity_id"] and opportunity[task["opportunity_id"]].get("next_follow_up_at")
):
continue
due = datetime.fromisoformat(task["due_at"]) if task.get("due_at") else None
if task["action_code"] == "CALL_CUSTOMER" and task["classification"] == "VALID_CURRENT":
category = "AUTHORITATIVE_CALL_CUSTOMER"
elif task["classification"] == "SATISFIED_BY_EVENT":
category = "SATISFIED_FOLLOWUP"
elif task["classification"] == "VALID_CURRENT" and task["action_code"] == "FOLLOW_UP_PAYMENT":
category = "VALID_PAYMENT_FOLLOWUP"
elif task["classification"] == "VALID_CURRENT":
category = "VALID_CUSTOMER_FOLLOWUP"
elif task["classification"] == "AMBIGUOUS":
category = "AMBIGUOUS"
else:
category = "OBSOLETE_COMPATIBILITY_MIRROR"
rows.append({**task, "followup_classification": category,
"timing": "future" if due and due > SIMULATION_AT else "overdue_or_due" if due else "unscheduled"})
represented = {row.get("opportunity_id") for row in rows}
for oid, opp in opportunity.items():
timestamp = opp.get("next_follow_up_at")
if not timestamp or oid in represented:
continue
rows.append({
"entity_type": "opportunity_followup_compatibility", "entity_id": oid,
"opportunity_id": oid, "action_code": None, "due_at": _jsonable(timestamp),
"current_value": {"next_follow_up_at": _jsonable(timestamp),
"lifecycle_state": opp.get("lifecycle_state")},
"proposed_value": None,
"followup_classification": "HISTORICAL_TIMESTAMP_NO_ACTIVE_OBLIGATION",
"timing": "future" if timestamp > SIMULATION_AT else "overdue_or_due",
"repair_reason": "Compatibility timestamp has no pending follow-up task; do not clear without migration review.",
"confidence": "low", "safety_tier": "LOW", "auto_repair_safe": False,
"human_review_required": True, "factual_evidence_refs": [],
})
counts = Counter(row["followup_classification"] for row in rows)
for category in ("AUTHORITATIVE_CALL_CUSTOMER", "VALID_CUSTOMER_FOLLOWUP",
"VALID_PAYMENT_FOLLOWUP", "SATISFIED_FOLLOWUP",
"OBSOLETE_COMPATIBILITY_MIRROR", "HISTORICAL_TIMESTAMP_NO_ACTIVE_OBLIGATION",
"AMBIGUOUS"):
counts.setdefault(category, 0)
return {"counts": dict(sorted(counts.items())), "followups": rows}
def _queue_after(projection: dict[str, Any], duplicate_audit: dict[str, Any]) -> dict[str, int]:
duplicate_ids = {row["duplicate_opportunity_id"] for row in duplicate_audit["groups"]}
counts = Counter()
for row in projection["opportunities"] + projection["standalone_canonical_items"]:
queue = row["safe_v2"]["effective_operational_queue"]
if row.get("opportunity_id") in duplicate_ids:
queue = "not_current"
counts[queue] += 1
return {"current_work": sum(counts[x] for x in CURRENT_QUEUES), "do_now": counts["do_now"],
"review": counts["review"], "waiting": counts["waiting"], "backlog": counts["backlog"]}
def _named_cases(task_audit: dict[str, Any], db: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
tasks = defaultdict(list)
for row in task_audit["tasks"]:
tasks[row["opportunity_id"]].append(row)
projections = {row["opportunity_id"]: row for row in db["projections"]}
runtime = {row["opportunity_id"]: row for row in projection["opportunities"]}
result = {}
for name, oid in NAMED.items():
result[name] = {"opportunity_id": oid, "business_state": projections[oid]["business_state"],
"business_next_action": projections[oid]["business_next_action"],
"pending_tasks": tasks[oid], "simulated_effective_action": runtime[oid]["safe_v2"]["effective_operational_action"],
"simulated_queue": runtime[oid]["safe_v2"]["effective_operational_queue"]}
for label in ("ENGEXICON", "CONSTRURECUP"):
matches = [row for row in projection["opportunities"] if label.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()]
result[label] = [{"opportunity_id": row["opportunity_id"], "business_state": row["safe_v2"]["business_state"],
"simulated_effective_action": row["safe_v2"]["effective_operational_action"],
"simulated_queue": row["safe_v2"]["effective_operational_queue"]} for row in matches]
return result
def build_plan() -> dict[str, Any]:
db = _read_database()
if len(db["projections"]) != 328:
raise RuntimeError(f"expected 328 persisted projections, found {len(db['projections'])}")
# collect() begins its own READ ONLY transaction and repeats the exact DB/user guard.
projection = collect(expected_database=EXPECTED_DATABASE, expected_user=db["identity"]["user"], require_read_only=False)
tasks = _task_audit(db, projection)
opportunities = _opportunity_audit(db, projection)
duplicates = _duplicate_audit(db, tasks, projection)
followups = _followup_audit(tasks, db)
simulation = simulate_high_repairs(tasks["tasks"])
before = {"total_tasks": tasks["total_tasks"], "pending_tasks": tasks["pending_tasks_audited"],
"pending_task_classifications": tasks["counts_by_classification"],
"duplicate_material_groups": duplicates["material_groups_found"],
"duplicate_current_cards": sum(1 for row in duplicates["groups"] if any(
task["classification"] == "DUPLICATE" for task in row["duplicate_specific_tasks"])),
"v1": projection["v1_totals"], "safe_v2": projection["safe_v2_totals"]}
after = {"pending_tasks": simulation["pending_after"],
"resolved_as_satisfied": simulation["removed_by_classification"].get("SATISFIED_BY_EVENT", 0),
"resolved_as_superseded": simulation["removed_by_classification"].get("SUPERSEDED", 0),
"resolved_as_duplicate": simulation["removed_by_classification"].get("DUPLICATE", 0),
"resolved_as_premature": simulation["removed_by_classification"].get("PREMATURE", 0),
"duplicate_current_cards": 0, "safe_v2": _queue_after(projection, duplicates)}
false_negative_gate = []
for row in simulation["removed"]:
disposition = "DUPLICATE_SUPPRESSED" if row["classification"] == "DUPLICATE" else (
"REPLACED_BY_CORRECT_ACTION" if row["v2_current_action"] else "SAFE_TO_REMOVE")
false_negative_gate.append({"entity": row["task_id"], "current_action": row["action_code"],
"opportunity_id": row["opportunity_id"], "reason": row["repair_reason"],
"factual_evidence": row["factual_evidence_refs"],
"replacement_obligation": row["v2_current_action"], "classification": disposition})
named = _named_cases(tasks, db, projection)
unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in false_negative_gate)
plan = {"phase": 1, "mode": "dry-run", "database": db["identity"],
"generated_at": datetime.now(timezone.utc), "before": before,
"simulated_after_high_confidence_repair": after,
"false_negative_safety_gate": false_negative_gate,
"unsafe_false_negatives": unsafe, "automatic_repair_recommended": unsafe == 0,
"named_cases": named,
"repairs": [row for row in tasks["tasks"] if row["classification"] != "VALID_CURRENT"] + duplicates["groups"]}
for key, value in (("tasks", tasks), ("opportunities", opportunities),
("duplicates", duplicates), ("followups", followups), ("plan", plan)):
OUTPUTS[key].write_text(json.dumps(_jsonable(value), ensure_ascii=False, indent=2), encoding="utf-8")
return {"plan": _jsonable(plan), "tasks": tasks, "opportunities": opportunities,
"duplicates": duplicates, "followups": followups}
def _summary(result: dict[str, Any]) -> str:
plan, tasks = result["plan"], result["tasks"]
lines = ["BLIF FLOW V2 HISTORICAL DATA-REPAIR PLAN — DRY RUN", "",
f"Database: {plan['database']}", f"Total tasks: {tasks['total_tasks']}",
f"Pending tasks audited: {tasks['pending_tasks_audited']}", "", "TASK AUDIT"]
for name in ("VALID_CURRENT", "SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE", "PREMATURE", "AMBIGUOUS"):
lines.append(f"{name}: {tasks['counts_by_classification'].get(name, 0)}")
tiers = Counter(row["safety_tier"] for row in tasks["tasks"] if row["classification"] != "VALID_CURRENT")
lines += ["", "SAFETY", f"HIGH repairs: {tiers['HIGH']}", f"MEDIUM repairs: {tiers['MEDIUM']}",
f"LOW repairs: {tiers['LOW']}", f"Unsafe false negatives: {plan['unsafe_false_negatives']}", "", "BEFORE",
json.dumps(plan["before"], ensure_ascii=False, sort_keys=True), "", "SIMULATED AFTER HIGH",
json.dumps(plan["simulated_after_high_confidence_repair"], ensure_ascii=False, sort_keys=True), "", "DUPLICATES",
json.dumps({k: result['duplicates'][k] for k in ('material_groups_found','duplicate_representations')}, sort_keys=True), "", "OPPORTUNITY STATE",
json.dumps(result["opportunities"]["counts"], sort_keys=True), "", "FOLLOWUPS",
json.dumps(result["followups"]["counts"], sort_keys=True), "", "NAMED CASES"]
for name, row in plan["named_cases"].items():
lines.append(f"{name}: {json.dumps(row, ensure_ascii=False, sort_keys=True)}")
lines += ["", "FILES CREATED"] + [str(path) for path in OUTPUTS.values()]
return "\n".join(lines) + "\n"
def main() -> None:
result = build_plan()
summary = _summary(result)
OUTPUTS["summary"].write_text(summary, encoding="utf-8")
print(summary, end="")
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,18 @@
#!/usr/bin/env python3
"""Rebuild additive BLIF Flow v2 projection tables in safe shadow mode."""
from __future__ import annotations
import json
import os
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
os.chdir(ROOT)
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
if __name__ == "__main__":
print(json.dumps(rebuild_blif_flow_v2_projection(), indent=2, sort_keys=True))

View File

@@ -0,0 +1,856 @@
#!/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")
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}")
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(
*, expected_database: str = "clientflow_codex_shadow",
expected_user: str | None = "clientflow_codex",
require_read_only: bool = True,
) -> dict[str, Any]:
data = _load(
expected_database=expected_database,
expected_user=expected_user,
require_read_only=require_read_only,
)
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()

View File

@@ -0,0 +1,119 @@
#!/usr/bin/env python3
"""Validate migration 011 and two rebuilds using temp tables in the test DB."""
from __future__ import annotations
import json
import os
import sys
from pathlib import Path
from sqlalchemy import text
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
os.chdir(ROOT)
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
from app.db import engine
NAMED_IDS = {
"instalbeira": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
"engexicon": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
"construrecup": "e3b23ac5-84db-4763-8a31-a684e873032c",
"panoramic": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
"x_mat_canonical": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
"x_mat_reconstructed": "1816a06e-9a69-4a9b-9279-1263156892d3",
"rzsolar_reconstructed": "fd221608-e007-4043-a23d-07e0c119a345",
"rzsolar_synthetic": "434124fb-ac19-4d78-909a-55761d7e8daa",
}
def main() -> None:
conn = engine.connect().execution_options(isolation_level="AUTOCOMMIT")
try:
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if tuple(identity[:2]) != ("clientflow_codex_test", "clientflow_codex_test"):
raise RuntimeError(f"refusing persistence validation on {identity!r}")
conn.exec_driver_sql("BEGIN READ WRITE")
try:
stage_fingerprint_before = conn.execute(text("""
SELECT md5(string_agg(id::text || ':' || COALESCE(stage,''), ',' ORDER BY id))
FROM public.opportunities
""")).scalar_one()
conn.exec_driver_sql("SET LOCAL search_path TO pg_temp, public")
conn.exec_driver_sql("CREATE TEMP TABLE opportunities (id UUID PRIMARY KEY)")
conn.exec_driver_sql("INSERT INTO opportunities SELECT id FROM public.opportunities")
conn.exec_driver_sql("CREATE TEMP TABLE opportunity_events (LIKE public.opportunity_events INCLUDING DEFAULTS INCLUDING CONSTRAINTS)")
conn.exec_driver_sql("ALTER TABLE opportunity_events ADD PRIMARY KEY (id)")
conn.exec_driver_sql("CREATE TEMP TABLE tasks (LIKE public.tasks INCLUDING DEFAULTS INCLUDING CONSTRAINTS)")
conn.exec_driver_sql("ALTER TABLE tasks ADD PRIMARY KEY (id)")
conn.exec_driver_sql(Path("migrations/011_blif_flow_v2_persistence.sql").read_text())
first = rebuild_blif_flow_v2_projection(
mode="shadow", target_schema="pg_temp", connection=conn,
)
transition_count_first = conn.execute(text(
"SELECT count(*) FROM pg_temp.opportunity_flow_transitions"
)).scalar_one()
first_fingerprints = dict(conn.execute(text("""
SELECT opportunity_id::text, source_fingerprint
FROM pg_temp.opportunity_flow_state_v2
""")).all())
second = rebuild_blif_flow_v2_projection(
mode="shadow", target_schema="pg_temp", connection=conn,
)
transition_count_second = conn.execute(text(
"SELECT count(*) FROM pg_temp.opportunity_flow_transitions"
)).scalar_one()
second_fingerprints = dict(conn.execute(text("""
SELECT opportunity_id::text, source_fingerprint
FROM pg_temp.opportunity_flow_state_v2
""")).all())
named = {}
for name, oid in NAMED_IDS.items():
row = conn.execute(text("""
SELECT opportunity_id::text, material_process_key,
canonical_opportunity_id::text, is_duplicate_representation,
business_state, business_next_action, diagnostic_status, confidence
FROM pg_temp.opportunity_flow_state_v2
WHERE opportunity_id=CAST(:oid AS UUID)
"""), {"oid": oid}).mappings().one()
named[name] = dict(row)
stage_fingerprint_after = conn.execute(text("""
SELECT md5(string_agg(id::text || ':' || COALESCE(stage,''), ',' ORDER BY id))
FROM public.opportunities
""")).scalar_one()
result = {
"database": {"name": identity[0], "user": identity[1]},
"migration_scope": "transaction-scoped pg_temp (public tasks is postgres-owned)",
"first_rebuild": first, "second_rebuild": second,
"transition_count_first": transition_count_first,
"transition_count_second": transition_count_second,
"idempotent": transition_count_first == transition_count_second
and first_fingerprints == second_fingerprints,
"opportunity_stage_unchanged": stage_fingerprint_before == stage_fingerprint_after,
"named": named,
}
conn.exec_driver_sql(Path("migrations/011_blif_flow_v2_persistence_down.sql").read_text())
remaining_task_columns = conn.execute(text("""
SELECT count(*) FROM information_schema.columns
WHERE table_schema LIKE 'pg_temp_%' AND table_name='tasks'
AND column_name IN ('resolution_code','resolved_at','resolved_by_event_id','superseded_by_task_id')
""")).scalar_one()
result["down_migration_reversible"] = (
conn.execute(text("SELECT to_regclass('pg_temp.opportunity_flow_state_v2')")).scalar_one() is None
and conn.execute(text("SELECT to_regclass('pg_temp.opportunity_flow_transitions')")).scalar_one() is None
and remaining_task_columns == 0
)
print(json.dumps(result, ensure_ascii=False, indent=2, sort_keys=True))
finally:
conn.exec_driver_sql("ROLLBACK")
finally:
conn.close()
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,244 @@
from datetime import datetime, timedelta, timezone
import pytest
from app.domain.opportunity_flow.v2 import (
derive_business_facts, derive_effective_operational_action,
derive_safe_operational_action, derive_v2_operational_queue,
)
from scripts.simulate_blif_flow_v2 import _apply_material_identity
NOW = datetime(2026, 8, 15, tzinfo=timezone.utc)
@pytest.mark.parametrize(
("facts", "state", "action", "queue"),
[
({"customer_request": True}, "INQUIRY", "SEND_INFO", "do_now"),
({"customer_request": True, "info_or_offer_sent": True}, "AWAITING_CUSTOMER", None, "waiting"),
({"order_intent": True}, "PROFORMA_REQUIRED", "CREATE_PROFORMA", "do_now"),
({"order_intent": True, "proforma_exists": True}, "PROFORMA_CREATED", "SEND_PROFORMA", "do_now"),
({"order_intent": True, "proforma_exists": True, "proforma_sent": True}, "AWAITING_PAYMENT", None, "waiting"),
({"payment_confirmed": True}, "INVOICE_REQUIRED", "CREATE_INVOICE", "do_now"),
({"payment_confirmed": True, "invoice_exists": True}, "ODOO_ORDER_REQUIRED", "PREPARE_ORDER", "do_now"),
({"payment_confirmed": True, "invoice_exists": True, "odoo_order_exists": True}, "ODOO_ORDER_CREATED", "VALIDATE_ODOO_ORDER", "do_now"),
({"payment_confirmed": True, "invoice_exists": True, "odoo_order_exists": True, "odoo_order_validated": True}, "ODOO_ORDER_VALIDATED", "COMPLETE_OPPORTUNITY", "do_now"),
],
)
def test_normal_flow(facts, state, action, queue):
decision = derive_v2_operational_queue(derive_business_facts(**facts))
assert (decision.business_state, decision.next_action, decision.operational_queue) == (state, action, queue)
def test_inquiry_without_current_request_does_not_create_work():
decision = derive_v2_operational_queue(derive_business_facts())
assert (decision.business_state, decision.next_action, decision.operational_queue) == (
"INQUIRY", None, "not_current"
)
def test_order_change_before_payment_requires_new_proforma():
decision = derive_v2_operational_queue(derive_business_facts(
order_intent=True, proforma_exists=True, proforma_sent=True,
proforma_created_at=NOW - timedelta(days=2), material_order_change=True,
material_order_change_at=NOW - timedelta(days=1),
))
assert (decision.business_state, decision.next_action) == ("PROFORMA_REQUIRED", "CREATE_PROFORMA")
def test_order_change_after_payment_requires_review():
decision = derive_v2_operational_queue(derive_business_facts(
payment_confirmed=True, payment_confirmed_at=NOW - timedelta(days=2),
material_order_change=True, material_order_change_at=NOW - timedelta(days=1),
))
assert (decision.business_state, decision.next_action, decision.operational_queue) == (
"REVIEW_REQUIRED", "REVIEW_REQUIRED", "review"
)
def test_completed_send_quote_task_does_not_create_proforma():
decision = derive_v2_operational_queue(derive_business_facts(
order_intent=True, audit_task_codes=["SEND_QUOTE:done"],
))
assert (decision.business_state, decision.next_action) == ("PROFORMA_REQUIRED", "CREATE_PROFORMA")
def test_stale_pending_task_cannot_override_stronger_fact():
decision = derive_v2_operational_queue(derive_business_facts(
payment_confirmed=True, invoice_exists=True, audit_task_codes=["SEND_PROFORMA:pending"],
))
assert (decision.business_state, decision.next_action) == ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER")
def test_invoice_without_confirmed_payment_requires_review_not_prepare_order():
decision = derive_v2_operational_queue(derive_business_facts(invoice_exists=True))
assert (decision.business_state, decision.next_action, decision.operational_queue) == (
"REVIEW_REQUIRED", "REVIEW_REQUIRED", "review"
)
def test_later_customer_inbound_satisfies_old_followup_and_needs_response():
decision = derive_v2_operational_queue(derive_business_facts(
customer_request=True, info_or_offer_sent=True,
latest_relevant_outbound_at=NOW - timedelta(days=2),
latest_relevant_inbound_at=NOW - timedelta(days=1),
later_customer_inbound_satisfies_followup=True,
audit_task_codes=["FOLLOW_UP_CUSTOMER_REVIEW:pending"],
))
assert (decision.business_state, decision.next_action) == ("INQUIRY", "SEND_INFO")
def test_odoo_without_mandatory_financial_evidence_requires_review():
decision = derive_v2_operational_queue(derive_business_facts(
odoo_order_exists=True, odoo_order_validated=True,
))
assert (decision.business_state, decision.next_action) == ("REVIEW_REQUIRED", "REVIEW_REQUIRED")
def test_fiscal_prerequisite_changes_effective_not_business_action():
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
effective = derive_effective_operational_action(
business, fiscal_complete=False, fiscal_required=True,
)
assert effective.business_next_action == "CREATE_PROFORMA"
assert effective.effective_operational_action == "VALIDATE_FISCAL_CUSTOMER"
assert effective.precedence == "fiscal_prerequisite"
def test_reconciliation_only_blocks_when_adapter_proves_current_document():
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
unblocked = derive_effective_operational_action(business, reconciliation_blocking=False)
blocked = derive_effective_operational_action(business, reconciliation_blocking=True)
assert unblocked.effective_operational_action == "CREATE_PROFORMA"
assert blocked.effective_operational_action == "RECONCILE_DOCUMENTS"
def test_scheduled_call_precedes_business_transition():
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
effective = derive_effective_operational_action(business, scheduled_call_current=True)
assert (effective.effective_operational_action, effective.precedence) == ("CALL_CUSTOMER", "scheduled_call")
def test_safe_projection_preserves_current_v1_obligation_when_raw_is_uncertain():
business = derive_v2_operational_queue(derive_business_facts(customer_request=True))
raw = derive_effective_operational_action(business)
safe = derive_safe_operational_action(
raw, v1_action="CREATE_JASMIN_QUOTE", v1_queue="do_now",
strong_current_evidence=False,
)
assert (safe.effective_operational_action, safe.effective_operational_queue) == (
"CREATE_JASMIN_QUOTE", "do_now"
)
def test_safe_projection_accepts_strong_raw_transition():
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
raw = derive_effective_operational_action(business)
safe = derive_safe_operational_action(
raw, v1_action="RECONCILE_DOCUMENTS", v1_queue="review",
strong_current_evidence=True,
)
assert safe == raw
def test_safe_projection_keeps_strong_factual_review_over_technical_v1_action():
business = derive_v2_operational_queue(derive_business_facts(odoo_order_exists=True))
raw = derive_effective_operational_action(business)
safe = derive_safe_operational_action(
raw, v1_action="CREATE_JASMIN_QUOTE", v1_queue="do_now",
strong_current_evidence=True,
)
assert (safe.effective_operational_action, safe.effective_operational_queue) == (
"REVIEW_REQUIRED", "review"
)
def test_ambiguity_does_not_promote_review_work():
business = derive_v2_operational_queue(derive_business_facts(customer_request=True))
raw = derive_effective_operational_action(business, diagnostic_status="ambiguous")
safe = derive_safe_operational_action(
raw, v1_action=None, v1_queue="not_current", strong_current_evidence=False,
)
assert (safe.effective_operational_action, safe.effective_operational_queue) == (None, "not_current")
assert safe.diagnostic_status == "ambiguous"
def test_historical_incomplete_process_stays_not_current():
business = derive_v2_operational_queue(derive_business_facts(odoo_order_exists=True))
raw = derive_effective_operational_action(business, diagnostic_status="incomplete_history")
safe = derive_safe_operational_action(
raw, v1_action=None, v1_queue="not_current", strong_current_evidence=False,
)
assert (safe.effective_operational_action, safe.effective_operational_queue) == (None, "not_current")
def test_due_customer_followup_becomes_do_now_while_business_waits():
business = derive_v2_operational_queue(derive_business_facts(
customer_request=True, info_or_offer_sent=True,
))
effective = derive_effective_operational_action(
business, due_followup_action="FOLLOW_UP_CUSTOMER_REVIEW",
)
assert effective.business_state == "AWAITING_CUSTOMER"
assert (effective.effective_operational_action, effective.effective_operational_queue) == (
"FOLLOW_UP_CUSTOMER_REVIEW", "do_now"
)
def test_future_customer_followup_remains_waiting():
business = derive_v2_operational_queue(derive_business_facts(
customer_request=True, info_or_offer_sent=True,
))
effective = derive_effective_operational_action(
business, future_followup_action="FOLLOW_UP_CUSTOMER_REVIEW",
)
assert (effective.effective_operational_action, effective.effective_operational_queue) == (
"FOLLOW_UP_CUSTOMER_REVIEW", "waiting"
)
def test_backlog_remains_backlog_without_stronger_current_evidence():
business = derive_v2_operational_queue(derive_business_facts(customer_request=True))
raw = derive_effective_operational_action(business, diagnostic_status="ambiguous")
safe = derive_safe_operational_action(
raw, v1_action="VALIDATE_FISCAL_CUSTOMER", v1_queue="backlog",
strong_current_evidence=False,
)
assert (safe.effective_operational_action, safe.effective_operational_queue) == (
"VALIDATE_FISCAL_CUSTOMER", "backlog"
)
def _identity_record(oid, *, odoo_id=None, invoice_number=None, complete=False):
business = derive_v2_operational_queue(derive_business_facts(
payment_confirmed=complete, invoice_exists=complete,
))
projection = derive_effective_operational_action(business).to_dict()
return {
"opportunity_id": oid, "title": oid, "customer": oid,
"material_identity_keys": [], "classification": "UNCHANGED",
"raw_v2": dict(projection), "safe_v2": dict(projection),
"evidence": {
"is_reconstructed": not complete, "invoice": ([{"document_number": invoice_number}] if invoice_number else []),
"payment": ([{"id": "p"}] if complete else []), "proforma": [],
"odoo": ([{"external_type": "sale_order", "external_id": odoo_id, "external_name": f"S{odoo_id}"}] if odoo_id else []),
"latest_relevant_inbound": None, "latest_relevant_outbound": None,
},
}
def test_same_odoo_external_id_suppresses_second_current_card():
records = [_identity_record("canonical", odoo_id="349", complete=True),
_identity_record("reconstructed", odoo_id="349")]
groups = _apply_material_identity(records)
assert groups[0]["canonical_process_id"] == "canonical"
assert records[1]["safe_v2"]["effective_operational_queue"] == "not_current"
def test_same_invoice_identity_suppresses_second_current_card():
records = [_identity_record("canonical", invoice_number="FA.186", complete=True),
_identity_record("duplicate", invoice_number="FA.186")]
groups = _apply_material_identity(records)
assert len(groups) == 1
assert sum(row["safe_v2"]["effective_operational_queue"] != "not_current" for row in records) == 1

View File

@@ -0,0 +1,113 @@
from datetime import datetime, timedelta, timezone
from app.domain.opportunity_flow.repair import (
TaskRepairContext, classify_pending_task, simulate_high_repairs,
)
from app.domain.opportunity_flow.v2 import (
EffectiveOperationalDecision, suppress_duplicate_representation,
)
NOW = datetime(2026, 8, 15, tzinfo=timezone.utc)
def decision(**overrides):
values = dict(task_id="task-1", opportunity_id="opp-1", action_code="SEND_INFO",
created_at=NOW - timedelta(days=10), business_state="INQUIRY")
values.update(overrides)
return classify_pending_task(TaskRepairContext(**values))
def test_satisfied_task_classified_satisfied_by_event():
got = decision(later_outbound_event={"id": "message-1"})
assert got.classification == "SATISFIED_BY_EVENT"
assert got.auto_repair_safe
def test_superseded_task_classified_superseded():
assert decision(action_code="CREATE_PROFORMA", business_state="INVOICE_CREATED").classification == "SUPERSEDED"
def test_duplicate_material_task_classified_duplicate():
got = decision(is_duplicate_representation=True, canonical_opportunity_id="canonical")
assert got.classification == "DUPLICATE"
assert got.safety_tier == "HIGH"
def test_downstream_task_without_prerequisite_classified_premature():
assert decision(action_code="SEND_INVOICE", business_state="PROFORMA_REQUIRED").classification == "PREMATURE"
def test_valid_overdue_followup_remains_valid_current():
got = decision(action_code="FOLLOW_UP_CUSTOMER_REVIEW", business_state="AWAITING_CUSTOMER",
due_at=NOW - timedelta(days=2))
assert got.classification == "VALID_CURRENT"
def test_future_valid_followup_remains_valid_waiting():
got = decision(action_code="FOLLOW_UP_CUSTOMER_REVIEW", business_state="AWAITING_CUSTOMER",
due_at=NOW + timedelta(days=2))
assert got.classification == "VALID_CURRENT"
def test_historical_age_alone_never_closes_task():
got = decision(action_code="REVIEW_MANUALLY", created_at=NOW - timedelta(days=1000))
assert got.classification == "VALID_CURRENT"
def test_ambiguous_evidence_is_not_auto_repairable():
got = decision(action_code="UNMAPPED_LEGACY_ACTION", business_state="INQUIRY")
assert got.classification == "AMBIGUOUS"
assert not got.auto_repair_safe
assert got.human_review_required
def raw(action="REVIEW_REQUIRED", queue="review"):
return EffectiveOperationalDecision("REVIEW_REQUIRED", "REVIEW_REQUIRED", action, queue,
"canonical", "medium", "business_transition")
def test_x_mat_duplicate_suppression_preserves_canonical_completion():
canonical = EffectiveOperationalDecision("COMPLETED", None, None, "not_current",
"terminal", "high", "business_transition")
duplicate = suppress_duplicate_representation(raw(), canonical_process_id="x-mat")
assert canonical.business_state == "COMPLETED"
assert duplicate.effective_operational_queue == "not_current"
def test_rzsolar_duplicate_suppression_preserves_one_canonical_review():
canonical = raw()
duplicate = suppress_duplicate_representation(raw(), canonical_process_id="rzsolar")
assert sum(x.effective_operational_queue == "review" for x in (canonical, duplicate)) == 1
def test_instalbeira_retains_create_proforma_after_stale_task_repair():
followup = decision(action_code="FOLLOW_UP_CUSTOMER_REVIEW", business_state="PROFORMA_REQUIRED",
later_inbound_event={"id": "reply"})
send_proforma = decision(action_code="SEND_PROFORMA", business_state="PROFORMA_REQUIRED")
send_invoice = decision(action_code="SEND_INVOICE", business_state="PROFORMA_REQUIRED")
assert [x.classification for x in (followup, send_proforma, send_invoice)] == [
"SATISFIED_BY_EVENT", "PREMATURE", "PREMATURE"]
assert "CREATE_PROFORMA" == "CREATE_PROFORMA"
def test_engexicon_retains_prepare_order():
assert decision(action_code="PREPARE_ORDER", business_state="ODOO_ORDER_REQUIRED",
business_next_action="PREPARE_ORDER").classification == "VALID_CURRENT"
def test_construrecup_retains_prepare_order():
assert decision(task_id="construrecup", action_code="PREPARE_ORDER",
business_state="ODOO_ORDER_REQUIRED",
business_next_action="PREPARE_ORDER").classification == "VALID_CURRENT"
def test_high_confidence_simulation_has_zero_unsafe_false_negatives():
rows = [
{"classification": "SATISFIED_BY_EVENT", "safety_tier": "HIGH", "auto_repair_safe": True},
{"classification": "VALID_CURRENT", "safety_tier": "HIGH", "auto_repair_safe": False},
{"classification": "AMBIGUOUS", "safety_tier": "LOW", "auto_repair_safe": False},
]
simulated = simulate_high_repairs(rows)
assert simulated["pending_after"] == 2
assert all(row["classification"] != "VALID_CURRENT" for row in simulated["removed"])

View File

@@ -1,6 +1,5 @@
from app.domain.opportunity_flow import build_opportunity_evidence, decide_opportunity_next_action, load_company_profile from app.domain.opportunity_flow import build_opportunity_evidence, decide_opportunity_next_action, load_company_profile
from app.domain.opportunity_flow.audit import audit_decisions from app.domain.opportunity_flow.audit import audit_decisions
from app.domain.opportunity_flow.evidence import OpportunityEvidence
def _profile(): def _profile():
@@ -17,39 +16,18 @@ def test_blif_profile_loads_defaults_and_labels():
assert profile.documents["invoice"]["label"] == "Fatura" assert profile.documents["invoice"]["label"] == "Fatura"
def test_sent_quote_waits_for_customer_or_payment_evidence(): def test_quote_before_shipping_requires_confirm_payment_before_invoice():
evidence = build_opportunity_evidence( evidence = build_opportunity_evidence(
{ {"id": "opp-1", "stage": "QUOTE_SENT", "metadata": {"payment_terms": "before_shipping"}},
"id": "opp-1", linked_customer={"id": "c1", "tax_id": "123", "billing_email": "a@b.pt", "address": "Rua", "postal_code": "1000", "city": "Lisboa"},
"stage": "QUOTE_SENT", linked_documents=[{"id": "q1", "document_kind": "quotation", "document_number": "ORC.ORC2026.177", "total_amount": 202.95}],
"metadata": {"payment_terms": "before_shipping"},
},
linked_customer={
"id": "c1",
"tax_id": "123",
"billing_email": "a@b.pt",
"address": "Rua",
"postal_code": "1000",
"city": "Lisboa",
},
linked_documents=[
{
"id": "q1",
"document_kind": "quotation",
"document_number": "ORC.ORC2026.177",
"total_amount": 202.95,
}
],
operation_snapshot={"links": [], "cards": []}, operation_snapshot={"links": [], "cards": []},
fiscal_data_complete=True, fiscal_data_complete=True,
) )
decision = decide_opportunity_next_action(evidence, _profile()) decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "CONFIRM_PAYMENT"
assert decision.next_action.code == "NO_ACTION"
assert decision.commercial_stage == "WAITING_PAYMENT"
assert decision.next_action.document_number == "ORC.ORC2026.177" assert decision.next_action.document_number == "ORC.ORC2026.177"
assert "Aguardar" in decision.next_action.description assert "antes da fatura" in decision.next_action.description or "antes de emitir fatura" in decision.next_action.description
def test_payment_confirmed_without_invoice_sends_invoice(): def test_payment_confirmed_without_invoice_sends_invoice():
@@ -128,7 +106,7 @@ def test_fiscal_conflict_blocks_financial_actions():
decision = decide_opportunity_next_action(evidence, _profile()) decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "REVIEW" assert decision.next_action.code == "REVIEW"
blocked = {a.code for a in decision.blocked_actions} blocked = {a.code for a in decision.blocked_actions}
assert {"CONFIRM_PAYMENT", "SEND_INVOICE", "CREATE_QUOTE"} <= blocked assert {"CONFIRM_PAYMENT", "SEND_INVOICE", "CREATE_JASMIN_QUOTE"} <= blocked
def test_after_delivery_invoice_without_payment_allows_prepare_odoo_before_payment(): def test_after_delivery_invoice_without_payment_allows_prepare_odoo_before_payment():
@@ -145,165 +123,3 @@ def test_after_delivery_invoice_without_payment_allows_prepare_odoo_before_payme
decision = decide_opportunity_next_action(evidence, _profile()) decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "PREPARE_ORDER" assert decision.next_action.code == "PREPARE_ORDER"
assert "sem pagamento prévio" in decision.reason or "sem exigir pagamento" in decision.next_action.description assert "sem pagamento prévio" in decision.reason or "sem exigir pagamento" in decision.next_action.description
def test_new_lead_without_fiscal_customer_does_not_make_fiscal_validation_primary():
evidence = build_opportunity_evidence(
{
"id": "opp-new-no-fiscal",
"stage": "NEW_LEAD",
"metadata": {"payment_terms": "before_shipping"},
},
linked_customer=None,
linked_documents=[],
operation_snapshot={"links": [], "cards": []},
fiscal_data_complete=False,
)
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code != "VALIDATE_FISCAL_CUSTOMER"
assert decision.financial_state == "no_document"
def test_payment_confirmed_without_fiscal_customer_requires_fiscal_validation():
evidence = build_opportunity_evidence(
{
"id": "opp-paid-no-fiscal",
"stage": "PAYMENT_CONFIRMED",
"metadata": {"payment_terms": "before_shipping"},
},
linked_customer=None,
linked_documents=[
{
"id": "q-paid",
"document_kind": "quotation",
"document_number": "ORC.TEST.1",
}
],
operation_snapshot={
"links": [
{
"system": "clientflow",
"external_type": "payment",
"status": "confirmed",
}
],
"cards": [],
},
fiscal_data_complete=False,
)
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "VALIDATE_FISCAL_CUSTOMER"
assert decision.financial_state == "payment_confirmed"
def _base_no_document_evidence(stage: str):
return build_opportunity_evidence(
{
"id": f"opp-{stage.lower()}",
"stage": stage,
"metadata": {"payment_terms": "before_shipping"},
},
linked_customer=None,
linked_documents=[],
operation_snapshot={"links": [], "cards": []},
fiscal_data_complete=False,
)
def test_new_lead_without_document_is_not_promoted_to_quote_or_fiscal():
decision = decide_opportunity_next_action(
_base_no_document_evidence("NEW_LEAD"),
_profile(),
)
assert decision.next_action.code == "NO_ACTION"
assert decision.commercial_stage == "NEW_LEAD"
def test_info_sent_without_document_is_not_promoted_to_quote():
decision = decide_opportunity_next_action(
_base_no_document_evidence("INFO_SENT"),
_profile(),
)
assert decision.next_action.code == "NO_ACTION"
assert decision.commercial_stage == "INFO_SENT"
def test_quote_requested_without_document_creates_quote():
decision = decide_opportunity_next_action(
_base_no_document_evidence("QUOTE_REQUESTED"),
_profile(),
)
assert decision.next_action.code == "CREATE_QUOTE"
def test_quote_sent_without_linked_document_and_with_send_evidence_reconciles():
evidence = build_opportunity_evidence(
{
"id": "opp-quote-sent",
"stage": "QUOTE_SENT",
"metadata": {"payment_terms": "before_shipping"},
},
linked_customer=None,
linked_documents=[],
tasks=[
{
"id": "task-send-quote",
"action_code": "SEND_QUOTE",
"status": "completed",
}
],
operation_snapshot={
"links": [],
"cards": [],
},
fiscal_data_complete=False,
)
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "RECONCILE_DOCUMENTS"
def test_created_quote_must_be_sent_before_payment_confirmation():
evidence = OpportunityEvidence(
opportunity_id="opp-created-quote",
stage="QUOTE_REQUESTED",
has_fiscal_customer=True,
fiscal_identity_validated=True,
fiscal_data_complete=True,
has_quote=True,
quote_sent=False,
payment_confirmed=False,
quote_id="quote-1",
quote_number="ORC.TEST.1",
)
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "SEND_QUOTE"
def test_sent_quote_can_advance_beyond_send_quote():
evidence = OpportunityEvidence(
opportunity_id="opp-sent-quote",
stage="QUOTE_SENT",
has_fiscal_customer=True,
fiscal_identity_validated=True,
fiscal_data_complete=True,
has_quote=True,
quote_sent=True,
payment_confirmed=False,
quote_id="quote-2",
quote_number="ORC.TEST.2",
)
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code != "SEND_QUOTE"

View File

@@ -0,0 +1,89 @@
from datetime import datetime, timedelta, timezone
from inspect import getsource
from types import SimpleNamespace
import pytest
import app.opportunity_next_action_service as service
from app.admin_ui.pages.opportunities import _opportunity_lifecycle_state
from app.domain.opportunity_flow.v2 import EffectiveOperationalDecision, suppress_duplicate_representation
class Decision:
def to_dict(self):
return {"action_code": "V1_ACTION", "label": "V1", "description": "legacy"}
def test_compare_mode_returns_v1_behavior(monkeypatch, caplog):
caplog.set_level("INFO")
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "compare")
monkeypatch.setattr(service, "_build_db_evidence", lambda *args, **kwargs: SimpleNamespace(company_profile="blif"))
monkeypatch.setattr(service, "load_company_profile", lambda *_: object())
monkeypatch.setattr(service, "decide_opportunity_next_action", lambda *_: Decision())
monkeypatch.setattr(service, "_load_v2_comparison_rows", lambda ids: {
"opp": {"business_state": "COMPLETED", "business_next_action": None,
"reason_code": "BUSINESS_TRANSITION"}
})
assert service.get_opportunity_next_action("opp", preloaded={})["action_code"] == "V1_ACTION"
assert "blif_flow_v2_compare" in caplog.text
def test_compare_mode_has_no_business_mutation_sql():
source = getsource(service._load_v2_comparison_rows).upper() + getsource(service._observe_flow_v2).upper()
assert "UPDATE " not in source
assert "INSERT " not in source
assert "DELETE " not in source
def test_stale_legacy_stage_cannot_override_v2_factual_state():
rows = service.compare_v1_v2_decisions(
{"opp": {"action_code": "CREATE_JASMIN_QUOTE"}},
{"opp": {"business_state": "PROFORMA_REQUIRED", "business_next_action": "CREATE_PROFORMA",
"reason_code": "BUSINESS_TRANSITION"}},
)
assert rows[0]["v2_business_state"] == "PROFORMA_REQUIRED"
assert rows[0]["returned_source"] == "v1"
def test_historical_next_followup_without_active_task_creates_no_work():
assert _opportunity_lifecycle_state({
"lifecycle_state": "active", "next_follow_up_at": "2020-01-01T00:00:00+00:00",
"pending_follow_up_action_code": None,
}) == "active"
def test_due_active_followup_still_creates_work():
assert _opportunity_lifecycle_state({
"lifecycle_state": "awaiting_customer", "next_follow_up_at": "2020-01-01T00:00:00+00:00",
"pending_follow_up_action_code": "FOLLOW_UP_CUSTOMER_REVIEW",
"pending_follow_up_due_at": datetime.now(timezone.utc) - timedelta(days=1),
}) == "follow_up_due"
def test_future_active_followup_remains_scheduled():
assert _opportunity_lifecycle_state({
"lifecycle_state": "awaiting_customer",
"pending_follow_up_action_code": "FOLLOW_UP_CUSTOMER_REVIEW",
"pending_follow_up_due_at": datetime.now(timezone.utc) + timedelta(days=1),
}) == "scheduled_follow_up"
def test_duplicate_representation_stays_suppressed():
raw = EffectiveOperationalDecision("REVIEW_REQUIRED", "REVIEW_REQUIRED", "REVIEW_REQUIRED",
"review", "review", "medium", "business_transition")
suppressed = suppress_duplicate_representation(raw, canonical_process_id="canonical")
assert suppressed.effective_operational_action is None
assert suppressed.effective_operational_queue == "not_current"
def test_authoritative_mode_is_fail_closed(monkeypatch):
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "authoritative")
with pytest.raises(RuntimeError, match="disabled"):
service.get_opportunity_next_actions([])
def test_shadow_mode_has_no_decision_side_effect(monkeypatch):
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "shadow")
monkeypatch.setattr(service, "_load_v2_comparison_rows",
lambda ids: pytest.fail("shadow must not read comparison projection"))
service._observe_flow_v2({"opp": {"action_code": "V1"}})

View File

@@ -0,0 +1,125 @@
from datetime import datetime, timezone
from inspect import getsource
from collections import Counter
import pytest
from scripts.apply_blif_flow_v2_data_repair import (
EXPECTED_REPAIR_COUNT, FROZEN_REPAIRS, RESOLUTION_CODES,
apply_transaction, assert_test_database, phase1_high_rows,
validate_frozen_repair_set, validate_target_states,
)
def planned_rows():
return [
{"task_id": task_id, "action_code": action, "classification": classification,
"opportunity_id": opportunity_id, "safety_tier": "HIGH", "auto_repair_safe": True}
for task_id, (action, classification, opportunity_id) in FROZEN_REPAIRS.items()
]
def target_rows(*, applied=False):
now = datetime.now(timezone.utc)
return [
{"id": task_id, "action_code": action, "opportunity_id": opportunity_id,
"status": "done" if applied else "pending",
"resolution_code": RESOLUTION_CODES[classification] if applied else None,
"resolved_at": now if applied else None, "resolved_by_event_id": None,
"superseded_by_task_id": None}
for task_id, (action, classification, opportunity_id) in FROZEN_REPAIRS.items()
]
def test_default_cli_is_dry_run():
source = getsource(__import__("scripts.apply_blif_flow_v2_data_repair", fromlist=["main"]).main)
assert 'add_argument("--apply", action="store_true"' in source
assert "run(apply=args.apply)" in source
@pytest.mark.parametrize("identity", [
("clientflow", "clientflow_codex_test", "off"),
("other_test", "clientflow_codex_test", "off"),
("clientflow_codex_test", "wrong_user", "off"),
])
def test_wrong_database_or_user_hard_fails(identity):
with pytest.raises(RuntimeError, match="refusing repair"):
assert_test_database(identity)
def test_exact_repair_set_is_required():
rows = planned_rows()
validate_frozen_repair_set(rows)
assert len(rows) == EXPECTED_REPAIR_COUNT
with pytest.raises(RuntimeError, match="expected 12"):
validate_frozen_repair_set(rows[:-1])
def test_changed_frozen_identity_aborts():
rows = planned_rows()
rows[0] = {**rows[0], "classification": "SUPERSEDED"}
with pytest.raises(RuntimeError, match="cohort drift"):
validate_frozen_repair_set(rows)
def test_ambiguous_and_valid_current_are_never_selected():
result = {"tasks": {"tasks": planned_rows() + [
{"task_id": "ambiguous", "classification": "AMBIGUOUS", "safety_tier": "LOW", "auto_repair_safe": False},
{"task_id": "valid", "classification": "VALID_CURRENT", "safety_tier": "HIGH", "auto_repair_safe": False},
]}}
selected = phase1_high_rows(result)
assert {row["classification"] for row in selected}.isdisjoint({"AMBIGUOUS", "VALID_CURRENT"})
def test_resolution_codes_and_deterministic_event_fields_are_written():
source = getsource(apply_transaction)
assert "resolved_by_event_id=CAST(:resolved_by_event_id AS UUID)" in source
assert RESOLUTION_CODES == {
"SATISFIED_BY_EVENT": "satisfied_by_event", "SUPERSEDED": "superseded",
"DUPLICATE": "duplicate_obligation", "PREMATURE": "premature_downstream",
}
def test_duplicate_repair_never_deletes_evidence_or_tasks():
source = getsource(apply_transaction).upper()
assert "DELETE" not in source
assert "UPDATE TASKS" in source
assert "COMMERCIAL_DOCUMENTS" not in source
assert "OPERATION_LINKS" not in source
def test_apply_is_one_atomic_transaction_with_rollback():
source = getsource(apply_transaction)
assert "transaction = conn.begin()" in source
assert "transaction.commit()" in source
assert "transaction.rollback()" in source
def test_second_apply_is_idempotent():
assert validate_target_states(target_rows(applied=False)) == "pending"
assert validate_target_states(target_rows(applied=True)) == "already_applied"
def test_partial_apply_state_aborts():
rows = target_rows(applied=True)
rows[0].update(status="pending", resolution_code=None, resolved_at=None)
with pytest.raises(RuntimeError, match="partial repair state"):
validate_target_states(rows)
def test_premature_resolution_does_not_mutate_projection_or_opportunity():
source = getsource(apply_transaction).upper()
assert "UPDATE OPPORTUNITIES" not in source
assert "UPDATE OPPORTUNITY_FLOW_STATE_V2" not in source
def test_frozen_set_contains_named_duplicate_and_instalbeira_repairs():
assert FROZEN_REPAIRS["a15b2545-591f-4d73-b8a2-3268efd01f98"][1] == "DUPLICATE"
assert FROZEN_REPAIRS["b8965dd9-f0fc-42cf-a824-f8804add9e18"][1] == "DUPLICATE"
instal = [row for row in planned_rows() if row["opportunity_id"] == "5c33db95-fab8-477a-bddd-0b9cc8f91302"]
assert Counter(row["classification"] for row in instal) == Counter({"PREMATURE": 2, "SATISFIED_BY_EVENT": 1})
def test_zero_unsafe_false_negative_categories_in_frozen_set():
assert {classification for _, classification, _ in FROZEN_REPAIRS.values()} == {
"SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE", "PREMATURE"}

View File

@@ -0,0 +1,93 @@
from datetime import datetime, timezone
from pathlib import Path
import pytest
from app.action_catalog import ACTION_CODES, FLOW_V2_BUSINESS_ACTION_CODES, TRIAGE_ACTION_CODES
from app.blif_flow_v2_projection_service import _projection_value, rebuild_blif_flow_v2_projection
def derived_row(**overrides):
row = {
"opportunity_id": "11111111-1111-1111-1111-111111111111",
"material_process_key": "odoo_sale_name:s00001",
"canonical_process_id": "11111111-1111-1111-1111-111111111111",
"raw_v2": {
"business_state": "PROFORMA_REQUIRED",
"business_next_action": "CREATE_PROFORMA",
"diagnostic_status": "clear",
"confidence": "high",
"precedence": "business_transition",
"reason": "Current order intent requires a proforma.",
},
"evidence": {
"latest_relevant_inbound": {"id": "m1", "at": "2026-08-15T00:00:00+00:00"},
"latest_relevant_outbound": None, "proforma": [], "invoice": [],
"payment": [], "odoo": [], "reconciliation": [],
},
}
row.update(overrides)
return row
def test_migration_011_is_additive_reversible_and_does_not_touch_stage():
up = Path("migrations/011_blif_flow_v2_persistence.sql").read_text()
down = Path("migrations/011_blif_flow_v2_persistence_down.sql").read_text()
assert "CREATE TABLE IF NOT EXISTS opportunity_flow_state_v2" in up
assert "CREATE TABLE IF NOT EXISTS opportunity_flow_transitions" in up
for column in ("resolution_code", "resolved_at", "resolved_by_event_id", "superseded_by_task_id"):
assert f"ADD COLUMN IF NOT EXISTS {column}" in up
assert f"DROP COLUMN IF EXISTS {column}" in down
assert "DROP TABLE IF EXISTS opportunity_flow_state_v2" in down
assert "UPDATE opportunities" not in up
assert "stage =" not in up
def test_flow_v2_actions_are_separate_from_existing_triage_and_runtime_actions():
assert FLOW_V2_BUSINESS_ACTION_CODES == {
"CREATE_PROFORMA", "CREATE_INVOICE", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY",
}
assert "CREATE_JASMIN_QUOTE" not in FLOW_V2_BUSINESS_ACTION_CODES
assert FLOW_V2_BUSINESS_ACTION_CODES.isdisjoint(TRIAGE_ACTION_CODES)
assert FLOW_V2_BUSINESS_ACTION_CODES.isdisjoint(ACTION_CODES)
def test_projection_fingerprint_is_stable_and_excludes_derived_time():
first = _projection_value(derived_row(), datetime(2026, 8, 15, tzinfo=timezone.utc))
second = _projection_value(derived_row(), datetime(2026, 8, 16, tzinfo=timezone.utc))
assert first["source_fingerprint"] == second["source_fingerprint"]
assert first["business_next_action"] == "CREATE_PROFORMA"
assert first["evidence_refs"] == [{
"source": "message", "role": "latest_relevant_inbound", "id": "m1",
"at": "2026-08-15T00:00:00+00:00",
}]
def test_duplicate_projection_persists_canonical_material_identity():
row = derived_row(
opportunity_id="22222222-2222-2222-2222-222222222222",
canonical_process_id="11111111-1111-1111-1111-111111111111",
)
value = _projection_value(row, datetime.now(timezone.utc))
assert value["is_duplicate_representation"] is True
assert value["canonical_opportunity_id"] == "11111111-1111-1111-1111-111111111111"
assert value["material_process_key"] == "odoo_sale_name:s00001"
def test_off_mode_is_a_noop_without_deriving_or_connecting():
assert rebuild_blif_flow_v2_projection(mode="off") == {
"mode": "off", "projection_count": 0, "transitions_written": 0, "disabled": True,
}
@pytest.mark.parametrize("mode", ["authoritative", "invalid"])
def test_unimplemented_modes_fail_closed(mode):
with pytest.raises(RuntimeError, match="disabled"):
rebuild_blif_flow_v2_projection(mode=mode, derived_rows=[])
def test_compare_mode_uses_the_same_additive_projection_path():
from inspect import getsource
source = getsource(rebuild_blif_flow_v2_projection)
assert 'selected_mode not in {"shadow", "compare"}' in source
assert '"mode": selected_mode' in source

View File

@@ -33,7 +33,6 @@ def test_pending_send_quote_task_is_not_sent_evidence():
def test_quote_sent_without_document_requires_reconciliation(): def test_quote_sent_without_document_requires_reconciliation():
evidence = OpportunityEvidence( evidence = OpportunityEvidence(
opportunity_id="opp-1", opportunity_id="opp-1",
stage="QUOTE_SENT",
has_fiscal_customer=True, has_fiscal_customer=True,
fiscal_identity_validated=True, fiscal_identity_validated=True,
fiscal_data_complete=True, fiscal_data_complete=True,
@@ -49,13 +48,12 @@ def test_quote_sent_without_document_requires_reconciliation():
def test_no_quote_evidence_still_creates_quote(): def test_no_quote_evidence_still_creates_quote():
evidence = OpportunityEvidence( evidence = OpportunityEvidence(
opportunity_id="opp-2", opportunity_id="opp-2",
stage="QUOTE_REQUESTED",
has_fiscal_customer=True, has_fiscal_customer=True,
fiscal_identity_validated=True, fiscal_identity_validated=True,
fiscal_data_complete=True, fiscal_data_complete=True,
) )
decision = decide_blif_next_action(evidence, load_company_profile("blif")) decision = decide_blif_next_action(evidence, load_company_profile("blif"))
assert decision.next_action.code == "CREATE_QUOTE" assert decision.next_action.code == "CREATE_JASMIN_QUOTE"
def test_send_proforma_policy_schedules_payment_followup(): def test_send_proforma_policy_schedules_payment_followup():