8 Commits

30 changed files with 4139 additions and 390 deletions

View File

@@ -40,6 +40,16 @@ 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 = {
"CALL_CUSTOMER": {
"route": "vendas",

View File

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

View File

@@ -0,0 +1,236 @@
"""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"})
PROJECTION_WRITE_TABLES = frozenset({"opportunity_flow_state_v2", "opportunity_flow_transitions"})
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(
*,
expected_database: str = "clientflow_codex_test",
expected_user: str | None = "clientflow_codex_test",
expected_opportunity_count: int | None = 328,
require_opportunities: bool = False,
) -> 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=expected_database,
expected_user=expected_user,
# 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,
expected_opportunity_count=expected_opportunity_count,
require_opportunities=require_opportunities,
)
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,
derive_expected_database: str = "clientflow_codex_test",
derive_expected_user: str | None = "clientflow_codex_test",
expected_opportunity_count: int | None = 328,
require_opportunities: bool = False,
) -> 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(
expected_database=derive_expected_database,
expected_user=derive_expected_user,
expected_opportunity_count=expected_opportunity_count,
require_opportunities=require_opportunities,
)
if require_opportunities and not rows:
raise RuntimeError("Flow v2 projection requires at least one opportunity")
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.
database_url: str
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_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy"

View File

@@ -249,15 +249,7 @@ def build_opportunity_evidence(
quote = _find_doc(docs, QUOTE_KINDS)
invoice = _find_doc(docs, INVOICE_KINDS)
# A existência do documento comercial prova apenas que o orçamento foi
# 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)
)
quote_sent = bool(quote) or _completed_send_quote_task_evidence(tasks)
pending_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_CONFIRM_PAYMENT,
ACTION_CREATE_QUOTE,
ACTION_SEND_QUOTE,
ACTION_FOLLOW_UP,
ACTION_FOLLOW_UP_PAYMENT,
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)
if not e.has_fiscal_customer:
# A ausência de cliente fiscal é prontidão operacional, não intenção
# comercial. Não deve substituir a próxima ação da oportunidade.
# O bloqueio fiscal é aplicado apenas mais abaixo quando uma transição
# 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."
)
blocked_actions.extend(_blocked(profile, code, "cliente fiscal por associar") for code in SENSITIVE_DOCUMENT_ACTIONS)
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)
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)
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")
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:
# Não promover automaticamente qualquer oportunidade para orçamento.
# 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:
if e.quote_sent:
next_action = _action(
profile,
ACTION_RECONCILE_DOCUMENTS,
@@ -162,26 +151,9 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
profile_name=profile.name,
decision_version=profile.version,
)
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)
return OpportunityDecision(
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,
)
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)
available_actions.append(next_action)
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)
if e.payment_terms == PAYMENT_AFTER_DELIVERY:
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)
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 not e.has_invoice and (
not e.has_fiscal_customer or not e.fiscal_data_complete
):
if not e.has_invoice and e.has_fiscal_customer and not e.fiscal_data_complete:
blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos"))
next_action = _action(
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)
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.
if (
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
):
if e.payment_terms in {PAYMENT_BEFORE_SHIPPING, "", "undefined", "agreement"} and e.has_quote and not e.payment_confirmed:
next_action = _action(
profile,
ACTION_NO_ACTION,
f"Orçamento {e.quote_number or ''} enviado. Aguardar decisão do cliente ou evidência de pagamento.",
force_label="Aguardar cliente / pagamento",
ACTION_CONFIRM_PAYMENT,
f"Orçamento {e.quote_number or ''} associado. Confirmar pagamento antes de emitir fatura.",
priority="alta",
target_url=f"/opportunities/{e.opportunity_id}#operacao" if e.opportunity_id else None,
document_id=e.quote_id,
document_number=e.quote_number,
)
return OpportunityDecision(
next_action,
"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,
)
available_actions.append(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)
if e.payment_confirmed and not e.has_invoice and (
not e.has_fiscal_customer or not e.fiscal_data_complete
):
if e.payment_confirmed and not e.has_invoice and e.has_fiscal_customer and not e.fiscal_data_complete:
blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos"))
next_action = _action(
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)
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_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,
)
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=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_VALIDATE_FISCAL_CUSTOMER = "VALIDATE_FISCAL_CUSTOMER"
ACTION_RECONCILE_DOCUMENTS = "RECONCILE_DOCUMENTS"
ACTION_CREATE_QUOTE = "CREATE_QUOTE"
ACTION_SEND_QUOTE = "SEND_QUOTE"
ACTION_CREATE_QUOTE = "CREATE_JASMIN_QUOTE"
ACTION_CONFIRM_PAYMENT = "CONFIRM_PAYMENT"
ACTION_SEND_INVOICE = "SEND_INVOICE"
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:
# operation_links é compatibilidade visual; não deve falhar o fluxo principal.
pass
# Criar o documento no Jasmin não significa que foi enviado ao cliente.
# 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,
)
set_opportunity_stage(opportunity_id, "QUOTE_SENT", note="Orçamento Jasmin criado via ClientFlow.", created_by="jasmin_service")
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.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_QUOTE", "SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
# v4928.1.5.129: central workflow emits SHIP_ORDER; operator tasks persist CREATE_SHIPMENT.
MATERIALIZED_ACTIONS.add("CREATE_SHIPMENT")
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.
config = get_action_config(action_code)
default_routes = {
"SEND_QUOTE": "vendas",
"SEND_INVOICE": "financeiro",
"FOLLOW_UP_PAYMENT": "financeiro",
"PREPARE_ORDER": "operacoes",
@@ -127,7 +125,6 @@ def ensure_pending_task_for_next_action(
route = "rever"
default_labels = {
"SEND_QUOTE": "Enviar orçamento ao cliente",
"SEND_INVOICE": "Enviar fatura ao cliente",
"FOLLOW_UP_PAYMENT": "Follow-up pagamento",
"PREPARE_ORDER": "Preparar encomenda / Odoo",
@@ -136,7 +133,6 @@ def ensure_pending_task_for_next_action(
"REVIEW_RECONSTRUCTED_PROCESS": "Validar processo reconstruído",
}
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.",
"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.",

View File

@@ -7,11 +7,14 @@ from __future__ import annotations
from dataclasses import asdict, dataclass
from collections import defaultdict
import json
import logging
from typing import Any, Dict, Iterable, Optional
from sqlalchemy import bindparam, text
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.
from app.domain.opportunity_flow import (
OpportunityEvidence,
@@ -21,6 +24,9 @@ from app.domain.opportunity_flow import (
)
logger = logging.getLogger(__name__)
@dataclass
class OpportunityNextAction:
action_code: str
@@ -37,6 +43,67 @@ class OpportunityNextAction:
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]]:
row = conn.execute(text(sql), params).mappings().first()
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.
"""
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)
if preloaded is not None else _build_db_evidence(opportunity_id))
if evidence is None:
@@ -172,8 +241,9 @@ def get_opportunity_next_action(
reason_if_blocked="opportunity_not_found",
).to_dict()
decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile))
return decision.to_dict()
decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile)).to_dict()
_observe_flow_v2({opportunity_id: decision})
return decision
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]]:
"""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()))
if not ids:
return {}
@@ -376,4 +448,5 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
company_profile="blif",
)
decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict()
_observe_flow_v2(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,328 @@
#!/usr/bin/env python3
"""Run a strictly read-only BLIF Flow v2 shadow/cutover audit."""
from __future__ import annotations
import argparse
import json
import os
import sys
from collections import Counter
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Sequence
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
os.chdir(ROOT)
PRODUCTION_DATABASE = "clientflow"
TEST_DATABASE = "clientflow_codex_test"
AUDIT = Path("/tmp/blif_flow_v2_production_shadow_audit.json")
SEMANTIC = Path("/tmp/blif_flow_v2_production_semantic_compare.json")
SUMMARY = Path("/tmp/blif_flow_v2_production_shadow_summary.txt")
CLASSIFICATIONS = {
"SEMANTICALLY_EQUIVALENT", "V1_OPERATIONAL_OVERRIDE", "V2_CORRECTS_V1",
"LEGACY_ONLY", "REAL_CONFLICT", "MISSING_PROJECTION",
}
CURRENT_QUEUES = {"do_now", "review", "exception"}
OVERRIDE_PRECEDENCE = {
"scheduled_call", "due_followup", "future_followup", "integration_exception",
"document_prerequisite", "fiscal_prerequisite", "safe_preserve_v1",
}
LEGACY_ACTIONS = {
"CREATE_JASMIN_QUOTE", "NO_ACTION", "WAIT_CUSTOMER", "WAIT_PAYMENT",
"WAIT_PRODUCTION", "WAIT_LOGISTICS", "WAIT_SUPPLIER", "WAIT_SCHEDULED_DATE",
}
REVIEW_ACTIONS = {
"REVIEW", "REVIEW_REQUIRED", "REVIEW_MANUALLY", "REVIEW_RECONSTRUCTED_PROCESS",
"RECONCILE_DOCUMENTS", "VALIDATE_FISCAL_CUSTOMER", "REVIEW_EXCEPTION",
}
NAMED_IDS = {
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
"PANORAMIC": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
"ENGEXICON": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
"CONSTRURECUP": "e3b23ac5-84db-4763-8a31-a684e873032c",
"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 build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument(
"--production-readonly-audit", action="store_true",
help="explicitly authorize a read-only audit of database clientflow",
)
return parser
def validate_execution(
*, production_readonly_audit: bool, database: str, transaction_read_only: str,
mode: str,
) -> None:
normalized_mode = str(mode or "").strip().lower()
if production_readonly_audit:
if database != PRODUCTION_DATABASE:
raise RuntimeError(
f"--production-readonly-audit requires database {PRODUCTION_DATABASE!r}, found {database!r}"
)
if transaction_read_only != "on":
raise RuntimeError("production audit requires transaction_read_only=on")
if normalized_mode not in {"shadow", "compare"}:
raise RuntimeError("production audit requires BLIF_FLOW_V2_MODE=shadow or compare")
return
if database == PRODUCTION_DATABASE:
raise RuntimeError("production database requires explicit --production-readonly-audit opt-in")
if database != TEST_DATABASE:
raise RuntimeError(f"default audit requires database {TEST_DATABASE!r}, found {database!r}")
if transaction_read_only != "on":
raise RuntimeError("cutover audit requires an explicit READ ONLY transaction")
def _code(value: Any) -> str:
return str(value or "").strip().upper()
def classify_semantic_difference(
*, v1_state: str | None, v1_action: str | None,
v2_state: str | None, v2_action: str | None,
v2_projection_present: bool = True,
is_duplicate_representation: bool = False,
operational_action: str | None = None,
operational_precedence: str | None = None,
operational_queue: str | None = None,
v2_confidence: str | None = None,
v2_diagnostic_status: str | None = None,
) -> tuple[str, str]:
"""Conservatively compare meanings rather than raw action vocabulary."""
if not v2_projection_present:
return "MISSING_PROJECTION", "No persisted Flow v2 projection exists for this opportunity."
v1, v2, effective = _code(v1_action), _code(v2_action), _code(operational_action)
precedence = str(operational_precedence or "").strip().lower()
queue = str(operational_queue or "").strip().lower()
if is_duplicate_representation:
return "V2_CORRECTS_V1", "Material identity suppresses a duplicate representation without deleting evidence."
if precedence in OVERRIDE_PRECEDENCE and effective and (v1 == effective or queue in CURRENT_QUEUES | {"waiting"}):
return "V1_OPERATIONAL_OVERRIDE", f"Explicit operational precedence {precedence} validly overlays the V2 business transition."
if v1 == v2 and v1:
return "SEMANTICALLY_EQUIVALENT", "V1 and V2 select the same action."
if not v1 and not v2:
return "SEMANTICALLY_EQUIVALENT", "Neither model has a current business action."
if v1 in REVIEW_ACTIONS and (v2 in REVIEW_ACTIONS or _code(v2_state) in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED"}):
return "SEMANTICALLY_EQUIVALENT", "Both decisions require review/blocker handling."
if v1 in {"NO_ACTION", ""} and not v2 and _code(v2_state) in {"COMPLETED", "LOST", "NO_INTEREST"}:
return "SEMANTICALLY_EQUIVALENT", "Both decisions represent a terminal/non-current process."
if v1.startswith("FOLLOW_UP_") and _code(v2_state) in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}:
return "V1_OPERATIONAL_OVERRIDE", "A current follow-up obligation overlays a waiting V2 business state."
if v1 == "CONFIRM_PAYMENT" and _code(v2_state) == "AWAITING_PAYMENT" and not v2:
return "SEMANTICALLY_EQUIVALENT", "Both decisions mean payment remains outstanding; V1 names the compatibility action."
if v1 in {"VALIDATE_FISCAL_CUSTOMER", "RECONCILE_DOCUMENTS"} and not effective and (
queue in {"not_current", "backlog"} or _code(v2_state) in {"INQUIRY", "AWAITING_CUSTOMER", "AWAITING_PAYMENT"}
):
return "LEGACY_ONLY", "Legacy data-hygiene/blocker vocabulary is not a current factual V2 obligation."
if precedence in {"safe_diagnostic_only", "safe_ambiguous_review"}:
return "V2_CORRECTS_V1", "SAFE V2 prevents ambiguous historical compatibility state from creating current work."
if _code(v2_state) in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED"}:
return "V2_CORRECTS_V1", "V2 converts conflicting factual history into an explicit protected review state."
if effective and effective == v2 and v1 != v2:
return "V2_CORRECTS_V1", "The safe operational action follows the factual V2 transition rather than the legacy action."
if v1 in LEGACY_ACTIONS:
if v2 and _code(v2_state) not in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}:
return "V2_CORRECTS_V1", "V2 replaces a legacy compatibility action with a factual business transition."
return "LEGACY_ONLY", "V1 action is compatibility vocabulary with no native V2 business transition."
if v2 and str(v2_confidence or "").lower() == "high" and str(v2_diagnostic_status or "").lower() == "clear":
return "V2_CORRECTS_V1", "High-confidence factual V2 transition corrects a different legacy action."
if effective and v1 == effective:
return "V1_OPERATIONAL_OVERRIDE", "V1 matches the safe operational overlay rather than the business action."
return "REAL_CONFLICT", "The available evidence does not establish equivalence, a valid override, or a safe V2 correction."
def _write(path: Path, value: Any) -> None:
path.write_text(json.dumps(value, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
def _queue_totals(rows: list[dict[str, Any]], key: str) -> dict[str, int]:
counts = Counter(str(row[key].get("operational_queue") or row[key].get("effective_operational_queue") or "not_current") for row in rows)
return {
"current": sum(counts[name] for name in CURRENT_QUEUES),
"do_now": counts["do_now"], "review": counts["review"],
"waiting": counts["waiting"], "backlog": counts["backlog"],
"not_current": counts["not_current"],
}
def _operation_totals(source: dict[str, Any]) -> dict[str, int]:
return {"current": int(source.get("current_work") or 0),
"do_now": int(source.get("do_now") or 0),
"review": int(source.get("review") or 0),
"waiting": int(source.get("waiting") or 0),
"backlog": int(source.get("backlog") or 0),
"not_current": int(source.get("not_current") or 0)}
def _database_snapshot(engine: Any) -> tuple[dict[str, str], list[str], dict[str, dict[str, Any]]]:
from sqlalchemy import text
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()
opportunity_ids = [str(value) for value in conn.execute(text(
"SELECT id FROM opportunities ORDER BY id"
)).scalars()]
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, evidence_refs, flow_version
FROM opportunity_flow_state_v2 ORDER BY opportunity_id
""")).mappings().all()
finally:
conn.rollback()
return (
{"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
opportunity_ids, {str(row["opportunity_id"]): dict(row) for row in rows},
)
def run_audit(*, production_readonly_audit: bool) -> dict[str, Any]:
# Imports occur only after argparse, so --help cannot initialize DB code.
from app.config import settings
from app.db import engine
from scripts.simulate_blif_flow_v2 import collect
identity, opportunity_ids, persisted = _database_snapshot(engine)
mode = str(settings.blif_flow_v2_mode or "").strip().lower()
validate_execution(
production_readonly_audit=production_readonly_audit,
database=identity["database"], transaction_read_only=identity["transaction_read_only"],
mode=mode,
)
preflight = {**identity, "BLIF_FLOW_V2_MODE": mode,
"production_readonly_opt_in": production_readonly_audit}
print(json.dumps(preflight, sort_keys=True), flush=True)
expected_count = None if production_readonly_audit else 328
projection = collect(
expected_database=identity["database"], expected_user=identity["user"],
require_read_only=production_readonly_audit,
expected_opportunity_count=expected_count, require_opportunities=True,
)
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
if set(records) != set(opportunity_ids):
raise RuntimeError("factual collector did not return the complete opportunity universe")
comparisons = []
for opportunity_id in opportunity_ids:
record = records[opportunity_id]
v1, operational = record["v1"], record["safe_v2"]
v2 = persisted.get(opportunity_id)
classification, reason = classify_semantic_difference(
v1_state=v1.get("commercial_stage"), v1_action=v1.get("current_action"),
v2_state=(v2 or {}).get("business_state"), v2_action=(v2 or {}).get("business_next_action"),
v2_projection_present=v2 is not None,
is_duplicate_representation=bool((v2 or {}).get("is_duplicate_representation")),
operational_action=operational.get("effective_operational_action"),
operational_precedence=operational.get("precedence"),
operational_queue=operational.get("effective_operational_queue"),
v2_confidence=(v2 or {}).get("confidence"),
v2_diagnostic_status=(v2 or {}).get("diagnostic_status"),
)
comparisons.append({
"opportunity_id": opportunity_id, "title": record.get("title"),
"customer_name": record.get("customer"), "classification": classification,
"classification_reason": reason,
"v1": {"state": v1.get("commercial_stage"), "action": v1.get("current_action"),
"queue": v1.get("operational_queue"), "reason_code": v1.get("reason")},
"v2": {"business_state": (v2 or {}).get("business_state"),
"business_next_action": (v2 or {}).get("business_next_action"),
"reason_code": (v2 or {}).get("reason_code"),
"diagnostic_status": (v2 or {}).get("diagnostic_status"),
"confidence": (v2 or {}).get("confidence")},
"operational_override": {"action": operational.get("effective_operational_action"),
"queue": operational.get("effective_operational_queue"),
"precedence": operational.get("precedence")},
"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")),
"evidence_refs": (v2 or {}).get("evidence_refs", []),
})
counts = Counter(row["classification"] for row in comparisons)
for name in CLASSIFICATIONS:
counts.setdefault(name, 0)
real_conflicts = [row for row in comparisons if row["classification"] == "REAL_CONFLICT"]
missing = [row for row in comparisons if row["classification"] == "MISSING_PROJECTION"]
semantic_report = {
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
"classification_counts": dict(sorted(counts.items())),
"real_conflicts": real_conflicts, "missing_projections": missing,
"comparisons": comparisons,
}
_write(SEMANTIC, semantic_report)
named = {}
by_id = {row["opportunity_id"]: row for row in comparisons}
for name, opportunity_id in NAMED_IDS.items():
row = by_id[opportunity_id]
named[name] = row
# Use the complete canonical Operations candidate universe, including
# preserved standalone obligations, rather than opportunity cards alone.
v1_metrics = _operation_totals(projection["v1_totals"])
safe_metrics = _operation_totals(projection["safe_v2_totals"])
audit = {
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
"read_only": True, "business_writes": 0, "projection_writes": 0,
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
"canonical_count": sum(not row.get("is_duplicate_representation") for row in persisted.values()),
"duplicate_representation_count": sum(bool(row.get("is_duplicate_representation")) for row in persisted.values()),
"operations_metrics": {"v1": v1_metrics, "safe_v2": safe_metrics},
"semantic_classification_counts": dict(sorted(counts.items())),
"authoritative_cutover_blockers": {
"real_conflicts": len(real_conflicts), "missing_projections": len(missing),
"blocked": bool(real_conflicts or missing),
},
"real_conflicts": real_conflicts, "missing_projections": missing,
"named_cases": named,
}
_write(AUDIT, audit)
lines = [
"BLIF FLOW V2 PRODUCTION SHADOW READ-ONLY AUDIT", "",
f"database: {identity['database']}", f"user: {identity['user']}",
f"transaction_read_only: {identity['transaction_read_only']}",
f"BLIF_FLOW_V2_MODE: {mode}",
f"production_readonly_opt_in: {production_readonly_audit}", "",
f"opportunities: {len(opportunity_ids)}", f"projections: {len(persisted)}",
f"canonical: {audit['canonical_count']}",
f"duplicate representations: {audit['duplicate_representation_count']}", "",
f"V1 metrics: {json.dumps(v1_metrics, sort_keys=True)}",
f"SAFE V2 metrics: {json.dumps(safe_metrics, sort_keys=True)}", "",
f"semantic classifications: {json.dumps(dict(sorted(counts.items())), sort_keys=True)}",
f"REAL_CONFLICT blockers: {len(real_conflicts)}",
f"MISSING_PROJECTION blockers: {len(missing)}",
"business writes: 0", "projection writes: 0",
]
SUMMARY.write_text("\n".join(lines) + "\n", encoding="utf-8")
return audit
def main(argv: Sequence[str] | None = None) -> int:
# --help exits here before application/database imports or connections.
args = build_parser().parse_args(argv)
result = run_audit(production_readonly_audit=args.production_readonly_audit)
print(json.dumps({
"database": result["preflight"]["database"],
"opportunities": result["opportunity_count"], "projections": result["projection_count"],
"semantic_classifications": result["semantic_classification_counts"],
"authoritative_cutover_blockers": result["authoritative_cutover_blockers"],
"outputs": [str(AUDIT), str(SEMANTIC), str(SUMMARY)],
}, indent=2, sort_keys=True))
return 0
if __name__ == "__main__":
raise SystemExit(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,84 @@
#!/usr/bin/env python3
"""Rebuild additive BLIF Flow v2 projection tables with explicit safeguards."""
from __future__ import annotations
import argparse
import json
import os
import sys
from pathlib import Path
from typing import Sequence
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
os.chdir(ROOT)
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument(
"--production-shadow", action="store_true",
help="explicitly authorize projection-only shadow writes to database clientflow",
)
return parser
def validate_execution(*, production_shadow: bool, mode: str, database: str) -> None:
"""Validate CLI intent independently of DATABASE_URL inference."""
normalized_mode = str(mode or "").strip().lower()
if production_shadow:
if normalized_mode != "shadow":
raise RuntimeError("--production-shadow requires BLIF_FLOW_V2_MODE=shadow")
if database != "clientflow":
raise RuntimeError(
f"--production-shadow requires database 'clientflow', found {database!r}"
)
return
if database == "clientflow":
raise RuntimeError("production database requires explicit --production-shadow opt-in")
def main(argv: Sequence[str] | None = None) -> int:
# argparse handles --help and exits before any application/DB import below.
args = build_parser().parse_args(argv)
from sqlalchemy import text
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
from app.config import settings
from app.db import engine
mode = str(settings.blif_flow_v2_mode or "").strip().lower()
with engine.connect() as conn:
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
database, user, transaction_read_only = identity
validate_execution(
production_shadow=args.production_shadow, mode=mode, database=database,
)
print(json.dumps({
"database": database, "user": user, "blif_flow_v2_mode": mode,
"transaction_read_only": transaction_read_only,
"production_shadow_opt_in": args.production_shadow,
}, sort_keys=True), flush=True)
if args.production_shadow:
result = rebuild_blif_flow_v2_projection(
mode="shadow",
# Production is deliberately scoped to this invocation; the
# module-level default allowlist remains test-only.
allowed_databases=frozenset({"clientflow"}),
derive_expected_database="clientflow",
derive_expected_user=user,
expected_opportunity_count=None,
require_opportunities=True,
)
else:
result = rebuild_blif_flow_v2_projection()
print(json.dumps(result, indent=2, sort_keys=True))
return 0
if __name__ == "__main__":
raise SystemExit(main())

View File

@@ -0,0 +1,877 @@
#!/usr/bin/env python3
"""Read-only BLIF Flow v2 projection against the isolated shadow snapshot."""
from __future__ import annotations
import json
import re
from collections import Counter, defaultdict
from dataclasses import replace
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterable
from sqlalchemy import text
from app.db import engine
from app.domain.opportunity_flow.v2 import (
EffectiveOperationalDecision,
derive_business_facts, derive_effective_operational_action,
derive_safe_operational_action, derive_v2_operational_queue,
suppress_duplicate_representation,
)
from app.operations_service import get_operations_summary
from app.opportunity_next_action_service import get_opportunity_next_actions
from app.opportunity_service import list_opportunities
PROJECTION = Path("/tmp/blif_flow_v2_projection.json")
COMPARISON = Path("/tmp/blif_flow_v2_operations_comparison.txt")
AMBIGUOUS = Path("/tmp/blif_flow_v2_ambiguous_cases.json")
PROMOTIONS = Path("/tmp/blif_flow_v2_promotions_audit.json")
INVOICE_WITHOUT_PAYMENT = Path("/tmp/blif_flow_v2_invoice_without_payment.json")
REVIEW_AUDIT = Path("/tmp/blif_flow_v2_review_audit.json")
BACKLOG_DELTA = Path("/tmp/blif_flow_v2_backlog_delta.json")
CURRENT_DELTA = Path("/tmp/blif_flow_v2_current_delta.json")
MATERIAL_IDENTITY = Path("/tmp/blif_flow_v2_material_identity.json")
TERMINAL = {"WON", "LOST", "NO_INTEREST", "ARCHIVED", "COMPLETED", "CLOSED"}
ORDER_INTENT = re.compile(r"\b(quero|queremos|pretendo|pretendemos|aceito|aceitamos|adjudic|encomendar|encomenda|avançar|avancar|proceder)\b", re.I)
ORDER_CHANGE_VERB = re.compile(r"\b(alterar|alteração|alteracao|mudar|mudança|mudanca|trocar|substituir|corrigir|retificar)\b", re.I)
ORDER_CHANGE_SUBJECT = re.compile(r"\b(produto|modelo|quantidade|morada|entrega|nif|fiscal|faturação|faturacao|condições|condicoes)\b", re.I)
QUOTE_REQUEST = re.compile(r"\b(preço|preco|orçamento|orcamento|cotação|cotacao|proposta|quote)\b", re.I)
PAYMENT_PROOF = re.compile(r"\b(comprovativo|transferência|transferencia|pagamento efetuado|pago|liquidado)\b", re.I)
SIMULATION_AT = datetime(2026, 8, 15, tzinfo=timezone.utc)
def _s(value: Any) -> str:
return str(value or "").strip()
def _dt(value: Any) -> datetime | None:
if isinstance(value, datetime):
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
if not value:
return None
try:
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
except ValueError:
return None
def _jsonable(value: Any) -> Any:
if isinstance(value, datetime):
return value.isoformat()
if isinstance(value, dict):
return {key: _jsonable(item) for key, item in value.items()}
if isinstance(value, (list, tuple)):
return [_jsonable(item) for item in value]
return value
def _compact(value: Any, limit: int = 260) -> str:
result = re.sub(r"\s+", " ", _s(value))
return result if len(result) <= limit else result[: limit - 1].rstrip() + "…"
def _payload(value: Any) -> dict[str, Any]:
if isinstance(value, dict):
return value
if isinstance(value, str) and value.strip():
try:
parsed = json.loads(value)
return parsed if isinstance(parsed, dict) else {}
except ValueError:
pass
return {}
def _group(rows: Iterable[dict[str, Any]], key: str = "opportunity_id") -> dict[str, list[dict[str, Any]]]:
result: dict[str, list[dict[str, Any]]] = defaultdict(list)
for row in rows:
result[_s(row.get(key))].append(dict(row))
return result
def _load(
*, expected_database: str = "clientflow_codex_shadow",
expected_user: str | None = "clientflow_codex",
require_read_only: bool = True,
) -> dict[str, Any]:
with engine.connect() as conn:
conn = conn.execution_options(isolation_level="AUTOCOMMIT")
transaction_started = False
try:
# A fresh connection commonly reports transaction_read_only=off.
# Establish the protected transaction on the same connection used
# for every factual read before validating a strict audit.
if require_read_only:
conn.execute(text("BEGIN READ ONLY"))
transaction_started = True
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if identity[0] != expected_database or (expected_user and identity[1] != expected_user):
raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
if require_read_only and identity[2] != "on":
raise RuntimeError(f"read-only simulation requires transaction_read_only=on: {identity!r}")
# Preserve the development/test path's established ordering: check
# its identity first, then protect the factual reads themselves.
if not require_read_only:
conn.execute(text("BEGIN READ ONLY"))
transaction_started = True
opportunities = [dict(row) for row in conn.execute(text("""
SELECT o.*, o.id::text AS id, o.local_customer_id::text,
c.name AS linked_customer_name, c.tax_id, c.email AS fiscal_email,
c.street_name, c.postal_zone, c.city_name, c.phone AS fiscal_phone
FROM opportunities o LEFT JOIN customers c ON c.id=o.local_customer_id
ORDER BY o.created_at, o.id
""")).mappings()]
tasks = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, action_code, action, note, status,
due_at, created_at, done_at, metadata
FROM tasks WHERE opportunity_id IS NOT NULL ORDER BY created_at
""")).mappings()]
messages = [dict(row) for row in conn.execute(text("""
SELECT o.id::text AS opportunity_id, m.id::text, m.direction,
COALESCE(m.clean_body,m.raw_body,'') AS body, m.created_at,
m.source_system, m.metadata
FROM opportunities o JOIN messages m ON m.conversation_id=o.conversation_id
WHERE m.source_system IN ('chatwoot','chatwoot_backfill')
ORDER BY m.created_at
""")).mappings()]
communications = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, direction, classification, subject,
body, status, created_at, metadata
FROM communications WHERE opportunity_id IS NOT NULL ORDER BY created_at
""")).mappings()]
documents = [dict(row) for row in conn.execute(text("""
SELECT DISTINCT ON (COALESCE(l.opportunity_id,d.opportunity_id),d.id)
COALESCE(l.opportunity_id,d.opportunity_id)::text AS opportunity_id,
d.id::text, d.document_kind, d.document_type, d.document_number,
d.external_id, d.status, d.payload, d.created_at, d.updated_at,
COALESCE(l.relationship, CASE WHEN d.is_primary THEN 'PRIMARY' ELSE d.role END, 'PRIMARY') AS relationship,
l.ended_at
FROM commercial_documents d
LEFT JOIN opportunity_document_links l ON l.document_id=d.id AND l.ended_at IS NULL
WHERE COALESCE(l.opportunity_id,d.opportunity_id) IS NOT NULL
ORDER BY COALESCE(l.opportunity_id,d.opportunity_id),d.id,l.updated_at DESC NULLS LAST
""")).mappings()]
links = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, system, external_type, external_id,
external_name, status, payload, created_at, updated_at, last_synced_at
FROM operation_links ORDER BY created_at
""")).mappings()]
reconciliation = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, title, description, document_number,
status, suggested_action, confidence, payload, created_at
FROM reconciliation_items
WHERE status IN ('open','needs_review','conflict') ORDER BY created_at
""")).mappings()]
finally:
if transaction_started:
conn.execute(text("ROLLBACK"))
return {
"identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
"opportunities": opportunities, "tasks": _group(tasks), "messages": _group(messages),
"communications": _group(communications), "documents": _group(documents),
"links": _group(links), "reconciliation": _group(reconciliation),
"all_reconciliation": reconciliation,
}
def _event(row: dict[str, Any]) -> dict[str, Any]:
return {
"id": row.get("id"), "at": _jsonable(row.get("created_at")),
"direction": row.get("direction"), "classification": row.get("classification"),
"text": _compact(row.get("body") or row.get("subject")),
}
def _derive_record(opp: dict[str, Any], data: dict[str, Any], v1: dict[str, Any], v1_item: dict[str, Any] | None) -> dict[str, Any]:
oid = _s(opp["id"])
messages = data["messages"].get(oid, [])
comms = data["communications"].get(oid, [])
tasks = data["tasks"].get(oid, [])
docs = data["documents"].get(oid, [])
links = data["links"].get(oid, [])
recons = data["reconciliation"].get(oid, [])
events = sorted(messages + comms, key=lambda row: _dt(row.get("created_at")) or datetime.min.replace(tzinfo=timezone.utc))
inbound = [row for row in events if _s(row.get("direction")).lower() == "inbound"]
outbound = [row for row in events if _s(row.get("direction")).lower() == "outbound"]
latest_in, latest_out = (inbound[-1] if inbound else None), (outbound[-1] if outbound else None)
inbound_text = "\n".join(_s(row.get("body") or row.get("subject")) for row in inbound)
inbound_classes = {_s(row.get("classification")).upper() for row in inbound}
request_kind = "quote" if inbound_classes & {"SEND_QUOTE", "SEND_PROFORMA"} or QUOTE_REQUEST.search(inbound_text) else "info"
order_rows = [row for row in inbound if _s(row.get("classification")).upper() in {"SEND_PROFORMA", "CONFIRM_PAYMENT", "SEND_INVOICE"}
or ORDER_INTENT.search(_s(row.get("body") or row.get("subject")))]
order_intent_at = _dt(order_rows[-1].get("created_at")) if order_rows else None
current_docs = [row for row in docs if _s(row.get("relationship")).upper() == "PRIMARY"
and _s(row.get("status")).lower() not in {"cancelled", "canceled", "failed"}]
proformas = [row for row in current_docs if _s(row.get("document_kind")).lower() in {"quotation", "quote", "proforma"}
and _s(row.get("status")).lower() != "converted"]
invoices = [row for row in current_docs if _s(row.get("document_kind")).lower() == "invoice"]
proforma = proformas[-1] if proformas else None
invoice = invoices[-1] if invoices else None
doc_number = _s((proforma or {}).get("document_number") or (proforma or {}).get("external_id"))
proforma_payload = _payload((proforma or {}).get("payload"))
sent_outbound = next((row for row in reversed(outbound) if doc_number and doc_number.casefold() in _s(row.get("body") or row.get("subject")).casefold()), None)
proforma_sent = bool(proforma and (
(proforma or {}).get("sent_at") or proforma_payload.get("sent_at") or proforma_payload.get("clientflow_sent_evidence")
or _s((proforma or {}).get("status")).lower() in {"sent", "issued_sent"} or sent_outbound
))
payment_links = [row for row in links if row.get("system") == "clientflow" and row.get("external_type") == "payment"]
payment = next((row for row in reversed(payment_links) if _s(row.get("status")).lower() == "confirmed"), None)
payment_proof_rows = [row for row in inbound if _s(row.get("classification")).upper() == "CONFIRM_PAYMENT"
or PAYMENT_PROOF.search(_s(row.get("body") or row.get("subject")))]
odoo_sales = [row for row in links if row.get("system") == "odoo" and row.get("external_type") == "sale_order"
and _s(row.get("status")).lower() not in {"not_found", "no_order", "cancelled"}]
validation = [row for row in links if row.get("system") == "odoo" and row.get("external_type") in {"physical_validation", "physical_status"}
and _s(row.get("status")).lower() in {"validated", "ready_to_ship", "shipped", "done", "delivered"}]
fulfilled = any(row.get("external_type") in {"physical_status", "delivery"} and _s(row.get("status")).lower() in {"shipped", "done", "delivered"} for row in links)
change_rows = [
row for row in inbound
if ORDER_CHANGE_VERB.search(_s(row.get("body") or row.get("subject")))
and ORDER_CHANGE_SUBJECT.search(_s(row.get("body") or row.get("subject")))
]
change_at = _dt(change_rows[-1].get("created_at")) if change_rows else None
proforma_at = _dt((proforma or {}).get("created_at"))
material_change = bool(change_at and proforma_at and change_at > proforma_at)
fiscal_complete = bool(opp.get("local_customer_id") and opp.get("tax_id") and opp.get("fiscal_email")
and opp.get("street_name") and opp.get("postal_zone") and opp.get("city_name"))
fiscal_conflict = bool((_payload(opp.get("metadata")).get("fiscal_conflict") or _payload(opp.get("metadata")).get("has_nif_conflict")))
blockers = []
if fiscal_conflict:
blockers.append("Conflicting fiscal/NIF evidence.")
conflict_recons = [row for row in recons if _s(row.get("status")).lower() in {"needs_review", "conflict"}]
if conflict_recons:
blockers.append("Unresolved document reconciliation conflict.")
reconstructed = _s(_payload(opp.get("metadata")).get("clientflow_record_mode")) in {
"reconstructed_invoice_review", "historical_reconstructed", "legacy_review"
} or any(marker in _s(opp.get("title")).casefold() for marker in ("processo reconstruído", "sem oportunidade"))
if reconstructed and not (invoice or payment or odoo_sales):
blockers.append("Reconstructed process lacks corroborating structured evidence.")
if invoice and not payment:
blockers.append("Structured invoice exists without confirmed payment evidence; correction/reconstruction flow is unspecified.")
if odoo_sales and (not payment or not invoice):
blockers.append("Odoo execution evidence exists without the mandatory linked payment and invoice evidence.")
status = _s(opp.get("status")).upper()
stage = _s(opp.get("stage")).upper()
lost = status in {"LOST", "NO_INTEREST"} or stage in {"LOST", "NO_INTEREST", "ARCHIVED"}
info_sent = bool(latest_out and (not latest_in or _dt(latest_out.get("created_at")) >= _dt(latest_in.get("created_at"))))
followup_satisfied = any(
_s(task.get("status")).lower() == "pending" and _s(task.get("action_code")).upper().startswith("FOLLOW_UP_")
and latest_in and _dt(latest_in.get("created_at")) > (_dt(task.get("created_at")) or datetime.max.replace(tzinfo=timezone.utc))
for task in tasks
)
sparse = not events and not current_docs and not links
review_required = (
bool(material_change and (payment or invoice)) or (reconstructed and sparse)
or bool(invoice and not payment) or bool(odoo_sales and (not payment or not invoice))
)
facts = derive_business_facts(
opportunity_id=oid, terminal=status in TERMINAL or stage in TERMINAL, explicitly_lost=lost,
review_required=review_required, fiscal_blocked=fiscal_conflict,
document_reconciliation_required=False, customer_request=bool(inbound),
request_kind=request_kind, latest_relevant_inbound_at=_dt((latest_in or {}).get("created_at")),
latest_relevant_outbound_at=_dt((latest_out or {}).get("created_at")), info_or_offer_sent=info_sent,
order_intent=bool(order_rows), order_intent_at=order_intent_at, fiscal_identity_evidence=fiscal_complete,
proforma_exists=bool(proforma), proforma_sent=proforma_sent, proforma_created_at=proforma_at,
proforma_sent_at=_dt((sent_outbound or {}).get("created_at")), potential_payment_evidence=bool(payment_proof_rows and not payment),
payment_confirmed=bool(payment), payment_confirmed_at=_dt((payment or {}).get("created_at")),
invoice_exists=bool(invoice), invoice_created_at=_dt((invoice or {}).get("created_at")),
odoo_order_exists=bool(odoo_sales), odoo_order_validated=bool(validation), fulfillment_complete=fulfilled,
material_order_change=material_change, material_order_change_at=change_at,
later_customer_inbound_satisfies_followup=followup_satisfied, blockers=blockers,
audit_task_codes=[f"{task.get('action_code')}:{task.get('status')}" for task in tasks],
)
decision = derive_v2_operational_queue(facts)
confidence = decision.confidence
ambiguity = []
if sparse and not lost:
confidence = "low"
ambiguity.append("No message, structured document, payment, or Odoo evidence is linked.")
if stage in {"QUOTE_SENT", "PROFORMA_SENT", "WAITING_PAYMENT"} and not proforma:
confidence = "low"
ambiguity.append("V1 stage suggests a formal offer, but no current structured proforma is linked.")
if stage == "PAYMENT_CONFIRMED" and not payment:
confidence = "low"
ambiguity.append("V1 stage says payment confirmed, but no confirmed payment operation link exists.")
if stage in {"WON", "SHIPPED", "ODOO_ORDER_CREATED", "IN_PRODUCTION"} and not odoo_sales:
confidence = "low"
ambiguity.append("V1 stage implies execution, but no Odoo sale-order link exists.")
decision = replace(decision, confidence=confidence)
v1_action = _s((v1_item or {}).get("current_action_code") or v1.get("action_code")) or None
v1_queue = _s((v1_item or {}).get("operational_queue")) or "not_current"
pending = [task for task in tasks if _s(task.get("status")).lower() == "pending"]
call_task = next((task for task in pending if _s(task.get("action_code")).upper() == "CALL_CUSTOMER"), None)
call_due = bool(call_task and (not _dt(call_task.get("due_at")) or _dt(call_task.get("due_at")) <= SIMULATION_AT))
followup_codes = (
{"FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA"}
if decision.business_state == "AWAITING_CUSTOMER"
else {"FOLLOW_UP_PAYMENT"} if decision.business_state == "AWAITING_PAYMENT" else set()
)
followup_task = next((
task for task in pending if _s(task.get("action_code")).upper() in followup_codes
), None)
followup_satisfied_now = bool(
followup_task and latest_in
and _dt(latest_in.get("created_at")) > (_dt(followup_task.get("created_at")) or SIMULATION_AT)
) or bool(followup_task and payment and _dt(payment.get("created_at")) > (_dt(followup_task.get("created_at")) or SIMULATION_AT))
followup_due = bool(
followup_task and not followup_satisfied_now and _dt(followup_task.get("due_at"))
and _dt(followup_task.get("due_at")) <= SIMULATION_AT
)
followup_future = bool(
followup_task and not followup_satisfied_now and _dt(followup_task.get("due_at"))
and _dt(followup_task.get("due_at")) > SIMULATION_AT
)
formal_doc_for_reconciliation = bool(current_docs)
reconciliation_blocking = bool(
formal_doc_for_reconciliation
and (conflict_recons or v1_action == "RECONCILE_DOCUMENTS")
)
fiscal_required = decision.next_action in {"CREATE_PROFORMA", "CREATE_INVOICE"}
if blockers and any("conflict" in blocker.casefold() or "without confirmed payment" in blocker.casefold() for blocker in blockers):
diagnostic_status = "conflicting_evidence"
elif sparse or (reconstructed and not (invoice and payment)) or (odoo_sales and (not invoice or not payment)):
diagnostic_status = "incomplete_history"
elif ambiguity:
diagnostic_status = "ambiguous"
else:
diagnostic_status = "clear"
raw = derive_effective_operational_action(
decision,
integration_exception=v1_queue == "exception",
scheduled_call_current=call_due,
due_followup_action=_s((followup_task or {}).get("action_code")).upper() if followup_due else None,
future_followup_action=_s((followup_task or {}).get("action_code")).upper() if followup_future else None,
fiscal_complete=fiscal_complete,
fiscal_required=fiscal_required,
reconciliation_blocking=reconciliation_blocking,
diagnostic_status=diagnostic_status,
)
strong_current_evidence = bool(
raw.precedence in {"integration_exception", "scheduled_call", "document_prerequisite", "fiscal_prerequisite"}
or (decision.business_state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "EXCEPTION"}
and diagnostic_status == "conflicting_evidence")
or followup_due
or (decision.next_action in {"SEND_PROFORMA"} and proforma)
or (decision.next_action == "CONFIRM_PAYMENT" and payment_proof_rows)
or (decision.next_action == "CREATE_INVOICE" and payment)
or (decision.next_action in {"PREPARE_ORDER", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY"} and invoice and payment)
or (decision.next_action == "CREATE_PROFORMA" and order_rows and not proforma)
or (decision.next_action in {"SEND_INFO", "SEND_QUOTE"} and latest_in
and (not latest_out or _dt(latest_in.get("created_at")) > _dt(latest_out.get("created_at"))))
or decision.operational_queue in {"waiting", "not_current"}
)
safe = derive_safe_operational_action(
raw, v1_action=v1_action, v1_queue=v1_queue,
strong_current_evidence=strong_current_evidence,
)
raw_v2 = raw.to_dict()
safe_v2 = safe.to_dict()
v2_action = raw.effective_operational_action
if decision.business_state == "REVIEW_REQUIRED":
classification = "REVIEW_REQUIRED"
elif ambiguity:
classification = "AMBIGUOUS"
elif v1_action == v2_action and v1_queue == raw.effective_operational_queue:
classification = "UNCHANGED"
elif v1_queue in {"do_now", "review", "exception"} and raw.effective_operational_queue in {"waiting", "not_current"}:
classification = "DEMOTED_TO_WAITING"
elif v1_queue in {"waiting", "backlog", "not_current"} and raw.effective_operational_queue in {"do_now", "review", "exception"}:
classification = "PROMOTED_TO_CURRENT"
else:
classification = "ACTION_CHANGED"
return _jsonable({
"opportunity_id": oid, "title": opp.get("title"),
"customer": opp.get("linked_customer_name") or opp.get("customer_name") or opp.get("customer_email"),
"v1": {"commercial_stage": opp.get("stage"), "lifecycle_state": opp.get("lifecycle_state"),
"current_action": v1_action, "operational_queue": v1_queue,
"reason": (v1_item or {}).get("eligibility_reason_code") or v1.get("reason")},
"raw_v2": raw_v2, "safe_v2": safe_v2,
# Compatibility alias for first-iteration report consumers.
"v2": raw_v2, "classification": classification,
"evidence": {
"latest_relevant_inbound": _event(latest_in) if latest_in else None,
"latest_relevant_outbound": _event(latest_out) if latest_out else None,
"order_intent": [_event(row) for row in order_rows[-3:]],
"fiscal_customer_identity": {"complete": fiscal_complete, "customer_id": _s(opp.get("local_customer_id")), "tax_id_present": bool(opp.get("tax_id"))},
"proforma": [{key: _jsonable(row.get(key)) for key in ("id", "external_id", "document_kind", "document_number", "status", "relationship", "created_at")} for row in proformas],
"payment": [{key: _jsonable(row.get(key)) for key in ("id", "status", "external_name", "created_at")} for row in payment_links],
"invoice": [{key: _jsonable(row.get(key)) for key in ("id", "external_id", "document_number", "status", "relationship", "created_at")} for row in invoices],
"odoo": [{key: _jsonable(row.get(key)) for key in ("id", "external_type", "external_id", "external_name", "status", "created_at")} for row in links if row.get("system") == "odoo"],
"blockers": blockers, "pending_tasks_for_audit_only": [
{key: _jsonable(task.get(key)) for key in ("id", "action_code", "status", "due_at", "created_at")} for task in tasks if task.get("status") == "pending"
], "ambiguity": ambiguity, "strong_current_evidence": strong_current_evidence,
"reconciliation_blocking": reconciliation_blocking,
"fiscal_required_for_transition": fiscal_required,
"scheduled_call_current": call_due,
"followup_due": followup_due, "followup_future": followup_future,
"followup_satisfied": followup_satisfied_now,
"diagnostic_status": diagnostic_status,
"is_reconstructed": reconstructed,
"reconciliation": [{
"id": row.get("id"), "document_number": row.get("document_number"),
"status": row.get("status"), "created_at": _jsonable(row.get("created_at")),
} for row in recons],
"material_order_change_evidence": [_event(row) for row in change_rows[-3:]],
},
})
def _material_keys(row: dict[str, Any]) -> set[str]:
evidence = row["evidence"]
keys = set()
for link in evidence.get("odoo", []):
if link.get("external_type") != "sale_order":
continue
if _s(link.get("external_id")):
keys.add(f"odoo_sale_id:{_s(link['external_id']).casefold()}")
sale_name = _s(link.get("external_name"))
if sale_name and re.fullmatch(r"[A-Z]{1,4}[-/]?[0-9]{2,}", sale_name, re.I):
keys.add(f"odoo_sale_name:{sale_name.casefold()}")
for kind in ("invoice", "proforma"):
for doc in evidence.get(kind, []):
if _s(doc.get("external_id")):
keys.add(f"jasmin_{kind}_id:{_s(doc['external_id']).casefold()}")
if _s(doc.get("document_number")):
keys.add(f"jasmin_{kind}_number:{_s(doc['document_number']).casefold()}")
return keys
def _apply_material_identity(records: list[dict[str, Any]]) -> list[dict[str, Any]]:
parent = {row["opportunity_id"]: row["opportunity_id"] for row in records}
def find(value: str) -> str:
while parent[value] != value:
parent[value] = parent[parent[value]]
value = parent[value]
return value
def union(left: str, right: str) -> None:
a, b = find(left), find(right)
if a != b:
parent[b] = a
by_key: dict[str, list[str]] = defaultdict(list)
for row in records:
keys = sorted(_material_keys(row))
row["material_identity_keys"] = keys
for key in keys:
by_key[key].append(row["opportunity_id"])
for ids in by_key.values():
for oid in ids[1:]:
union(ids[0], oid)
groups: dict[str, list[dict[str, Any]]] = defaultdict(list)
for row in records:
groups[find(row["opportunity_id"])].append(row)
report = []
for group in groups.values():
if len(group) < 2:
row = group[0]
row["material_process_key"] = next(iter(row["material_identity_keys"]), f"opportunity:{row['opportunity_id']}")
row["canonical_process_id"] = row["opportunity_id"]
row["duplicate_process_ids"] = []
continue
def score(row: dict[str, Any]) -> tuple[int, str]:
evidence = row["evidence"]
value = 0
value += 50 if not evidence.get("is_reconstructed") else 0
value += 20 if evidence.get("invoice") else 0
value += 20 if evidence.get("payment") else 0
value += 15 if evidence.get("proforma") else 0
value += 15 if any(link.get("external_type") == "sale_order" for link in evidence.get("odoo", [])) else 0
value += 10 if evidence.get("latest_relevant_inbound") or evidence.get("latest_relevant_outbound") else 0
value += 8 if "processo reconstruído" in _s(row.get("title")).casefold() else 0
value -= 8 if "sem oportunidade" in _s(row.get("title")).casefold() else 0
return value, row["opportunity_id"]
canonical = max(group, key=score)
common = set(canonical["material_identity_keys"])
for row in group:
common &= set(row["material_identity_keys"])
preferred = sorted(common, key=lambda key: (0 if key.startswith("odoo_sale_id:") else 1, key))
process_key = preferred[0] if preferred else sorted(canonical["material_identity_keys"])[0]
duplicates = [row["opportunity_id"] for row in group if row is not canonical]
canonical["material_process_key"] = process_key
canonical["canonical_process_id"] = canonical["opportunity_id"]
canonical["duplicate_process_ids"] = duplicates
for duplicate in group:
if duplicate is canonical:
continue
duplicate["material_process_key"] = process_key
duplicate["canonical_process_id"] = canonical["opportunity_id"]
duplicate["duplicate_process_ids"] = []
duplicate["raw_v2"] = suppress_duplicate_representation(
EffectiveOperationalDecision(**duplicate["raw_v2"]),
canonical_process_id=canonical["opportunity_id"],
).to_dict()
duplicate["safe_v2"] = suppress_duplicate_representation(
EffectiveOperationalDecision(**duplicate["safe_v2"]),
canonical_process_id=canonical["opportunity_id"],
).to_dict()
duplicate["classification"] = "DUPLICATE_REPRESENTATION"
report.append({
"material_process_key": process_key,
"canonical_process_id": canonical["opportunity_id"],
"duplicate_process_ids": duplicates,
"identity_keys": sorted(set.intersection(*(set(row["material_identity_keys"]) for row in group))),
"canonical_reason": "Highest factual completeness; prefers non-reconstructed process and explicit reconstructed process over an unassociated synthetic record.",
})
return report
def _totals(items: Iterable[dict[str, Any]], queue_key: str) -> dict[str, int]:
counts = Counter(_s(item.get(queue_key)) or "not_current" for item in items)
return {
"current_work": sum(counts[name] for name in ("do_now", "review", "exception")),
"do_now": counts["do_now"], "review": counts["review"], "waiting": counts["waiting"],
"backlog": counts["backlog"], "exception": counts["exception"], "not_current": counts["not_current"],
}
def collect(
*, expected_database: str = "clientflow_codex_shadow",
expected_user: str | None = "clientflow_codex",
require_read_only: bool = True,
expected_opportunity_count: int | None = None,
require_opportunities: bool = False,
) -> dict[str, Any]:
data = _load(
expected_database=expected_database,
expected_user=expected_user,
require_read_only=require_read_only,
)
opportunities = data["opportunities"]
if expected_opportunity_count is not None and len(opportunities) != expected_opportunity_count:
raise RuntimeError(
f"expected {expected_opportunity_count} opportunities, found {len(opportunities)}"
)
if require_opportunities and not opportunities:
raise RuntimeError("Flow v2 projection requires at least one opportunity")
ids = [_s(opp["id"]) for opp in opportunities]
v1_decisions = get_opportunity_next_actions(ids)
operations = get_operations_summary(limit=200)
all_v1_items = []
for key in ("work_items", "waiting_items", "backlog_items", "not_current_items"):
all_v1_items.extend(operations.get(key, []))
by_opp = {_s(item.get("opportunity_id")): item for item in all_v1_items if item.get("opportunity_id")}
records = [_derive_record(opp, data, v1_decisions.get(_s(opp["id"]), {}), by_opp.get(_s(opp["id"]))) for opp in opportunities]
material_identity = _apply_material_identity(records)
classifications = Counter(row["classification"] for row in records)
standalone = []
for item in all_v1_items:
if item.get("opportunity_id"):
continue
action = _s(item.get("current_action_code")) or None
queue = _s(item.get("operational_queue")) or "backlog"
projection = {
"business_state": None, "business_next_action": None,
"effective_operational_action": action,
"effective_operational_queue": queue,
"reason": "Standalone canonical Operations work is outside the standard commercial flow and is preserved.",
"confidence": "high", "precedence": "preserved_non_opportunity",
}
standalone.append({
"candidate_key": item.get("work_item_key") or f"standalone:{item.get('source')}:{item.get('id')}",
"opportunity_id": None, "title": item.get("title"), "customer": item.get("customer_name"),
"v1": {"current_action": action, "operational_queue": queue,
"reason": item.get("eligibility_reason_code")},
"raw_v2": dict(projection), "safe_v2": dict(projection),
"classification": "UNCHANGED", "source": item.get("source"),
})
universe = records + standalone
transitions = Counter((
row["v1"]["current_action"] or "<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.audit import audit_decisions
from app.domain.opportunity_flow.evidence import OpportunityEvidence
def _profile():
@@ -17,39 +16,18 @@ def test_blif_profile_loads_defaults_and_labels():
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(
{
"id": "opp-1",
"stage": "QUOTE_SENT",
"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,
}
],
{"id": "opp-1", "stage": "QUOTE_SENT", "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": []},
fiscal_data_complete=True,
)
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "NO_ACTION"
assert decision.commercial_stage == "WAITING_PAYMENT"
assert decision.next_action.code == "CONFIRM_PAYMENT"
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():
@@ -128,7 +106,7 @@ def test_fiscal_conflict_blocks_financial_actions():
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "REVIEW"
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():
@@ -145,165 +123,3 @@ def test_after_delivery_invoice_without_payment_allows_prepare_odoo_before_payme
decision = decide_opportunity_next_action(evidence, _profile())
assert decision.next_action.code == "PREPARE_ORDER"
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

@@ -0,0 +1,110 @@
from inspect import getsource
import pytest
import scripts.audit_blif_flow_v2_cutover as audit
def classify(**overrides):
values = dict(
v1_state="INFO_SENT", v1_action="SEND_INFO",
v2_state="INQUIRY", v2_action="SEND_INFO",
v2_projection_present=True, is_duplicate_representation=False,
operational_action="SEND_INFO", operational_precedence="business_transition",
operational_queue="do_now", v2_confidence="high", v2_diagnostic_status="clear",
)
values.update(overrides)
return audit.classify_semantic_difference(**values)[0]
def test_help_performs_zero_database_work(monkeypatch, capsys):
monkeypatch.setattr(audit, "run_audit", lambda **kwargs: pytest.fail("help reached DB audit"))
with pytest.raises(SystemExit) as exc:
audit.main(["--help"])
assert exc.value.code == 0
assert "--production-readonly-audit" in capsys.readouterr().out
def test_production_requires_explicit_opt_in():
with pytest.raises(RuntimeError, match="explicit --production-readonly-audit"):
audit.validate_execution(production_readonly_audit=False, database="clientflow",
transaction_read_only="on", mode="shadow")
def test_production_database_must_be_clientflow():
with pytest.raises(RuntimeError, match="requires database 'clientflow'"):
audit.validate_execution(production_readonly_audit=True, database="clientflow_codex_test",
transaction_read_only="on", mode="shadow")
def test_production_transaction_must_be_read_only():
with pytest.raises(RuntimeError, match="transaction_read_only=on"):
audit.validate_execution(production_readonly_audit=True, database="clientflow",
transaction_read_only="off", mode="shadow")
def test_production_audit_contains_no_sql_writes():
source = (getsource(audit._database_snapshot) + getsource(audit.run_audit)).upper()
for verb in ("INSERT ", "UPDATE ", "DELETE ", "CREATE ", "ALTER ", "DROP ", "TRUNCATE "):
assert verb not in source
assert 'BEGIN READ ONLY' in source
@pytest.mark.parametrize("mode", ["shadow", "compare"])
def test_shadow_and_compare_are_allowed(mode):
audit.validate_execution(production_readonly_audit=True, database="clientflow",
transaction_read_only="on", mode=mode)
@pytest.mark.parametrize("mode", ["", "off", "authoritative"])
def test_unsafe_production_modes_are_refused(mode):
with pytest.raises(RuntimeError, match="requires BLIF_FLOW_V2_MODE=shadow or compare"):
audit.validate_execution(production_readonly_audit=True, database="clientflow",
transaction_read_only="on", mode=mode)
def test_missing_projection_is_reported():
assert classify(v2_projection_present=False, v2_state=None, v2_action=None) == "MISSING_PROJECTION"
def test_semantically_equivalent_is_classified():
assert classify() == "SEMANTICALLY_EQUIVALENT"
assert classify(v1_action="NO_ACTION", v2_state="COMPLETED", v2_action=None,
operational_action=None, operational_queue="not_current") == "SEMANTICALLY_EQUIVALENT"
def test_operational_override_is_classified():
assert classify(v1_action="FOLLOW_UP_CUSTOMER_REVIEW", v2_state="AWAITING_CUSTOMER",
v2_action=None, operational_action="FOLLOW_UP_CUSTOMER_REVIEW",
operational_precedence="due_followup") == "V1_OPERATIONAL_OVERRIDE"
def test_v2_correction_is_classified():
assert classify(v1_action="SEND_INVOICE", v2_state="PROFORMA_REQUIRED",
v2_action="CREATE_PROFORMA") == "V2_CORRECTS_V1"
assert classify(v1_action="REVIEW_RECONSTRUCTED_PROCESS", v2_state="REVIEW_REQUIRED",
v2_action="REVIEW_REQUIRED", is_duplicate_representation=True) == "V2_CORRECTS_V1"
def test_legacy_only_is_classified():
assert classify(v1_action="CREATE_JASMIN_QUOTE", v2_state="AWAITING_CUSTOMER",
v2_action=None, v2_confidence="medium") == "LEGACY_ONLY"
def test_real_conflict_is_classified():
assert classify(v1_action="UNMAPPED_ACTION", v2_state="INQUIRY",
v2_action=None, v2_confidence="medium",
v2_diagnostic_status="ambiguous", operational_action=None) == "REAL_CONFLICT"
def test_named_case_ids_and_expected_semantics_are_frozen():
assert audit.NAMED_IDS == {
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
"PANORAMIC": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
"ENGEXICON": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
"CONSTRURECUP": "e3b23ac5-84db-4763-8a31-a684e873032c",
"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",
}

View File

@@ -0,0 +1,104 @@
from inspect import getsource
import pytest
import app.blif_flow_v2_projection_service as projection
import scripts.rebuild_blif_flow_v2_projection as cli
import scripts.simulate_blif_flow_v2 as simulator
def test_help_performs_no_rebuild(monkeypatch, capsys):
monkeypatch.setattr(cli, "validate_execution", lambda **kwargs: pytest.fail("help reached execution"))
with pytest.raises(SystemExit) as exc:
cli.main(["--help"])
assert exc.value.code == 0
assert "--production-shadow" in capsys.readouterr().out
def test_default_execution_does_not_authorize_production():
with pytest.raises(RuntimeError, match="explicit --production-shadow"):
cli.validate_execution(production_shadow=False, mode="shadow", database="clientflow")
@pytest.mark.parametrize("mode", ["", "off", "compare", "authoritative"])
def test_production_shadow_requires_exact_shadow_mode(mode):
with pytest.raises(RuntimeError, match="requires BLIF_FLOW_V2_MODE=shadow"):
cli.validate_execution(production_shadow=True, mode=mode, database="clientflow")
@pytest.mark.parametrize("database", ["clientflow_codex_test", "clientflow_codex_shadow", "other"])
def test_production_shadow_requires_exact_production_database(database):
with pytest.raises(RuntimeError, match="requires database 'clientflow'"):
cli.validate_execution(production_shadow=True, mode="shadow", database=database)
def test_production_shadow_explicit_proof_is_accepted():
cli.validate_execution(production_shadow=True, mode="shadow", database="clientflow")
def test_global_write_allowlist_remains_test_only():
assert projection.WRITE_DATABASE_ALLOWLIST == frozenset({"clientflow_codex_test"})
assert "clientflow" not in projection.WRITE_DATABASE_ALLOWLIST
def test_factual_read_transaction_remains_read_only():
source = getsource(simulator._load)
assert 'conn.execute(text("BEGIN READ ONLY"))' in source
assert 'conn.execute(text("ROLLBACK"))' in source
def test_projection_writer_only_mutates_two_additive_tables():
source = getsource(projection.rebuild_blif_flow_v2_projection).upper()
assert projection.PROJECTION_WRITE_TABLES == {
"opportunity_flow_state_v2", "opportunity_flow_transitions"}
assert "INSERT INTO OPPORTUNITY_FLOW_TRANSITIONS" in source
assert "INSERT INTO OPPORTUNITY_FLOW_STATE_V2" in source
for forbidden in ("OPPORTUNITIES", "TASKS", "MESSAGES", "COMMUNICATIONS",
"COMMERCIAL_DOCUMENTS", "CUSTOMERS", "PAYMENTS", "OPERATION_LINKS"):
assert f"INSERT INTO {forbidden}" not in source
assert f"UPDATE {forbidden}" not in source
assert f"DELETE FROM {forbidden}" not in source
def test_production_derivation_has_no_fixed_328_requirement(monkeypatch):
captured = {}
def fake_collect(**kwargs):
captured.update(kwargs)
return {"opportunities": [{"opportunity_id": "one"}]}
monkeypatch.setattr(simulator, "collect", fake_collect)
rows = projection._derive_all(
expected_database="clientflow", expected_user="runtime-role",
expected_opportunity_count=None, require_opportunities=True,
)
assert rows == [{"opportunity_id": "one"}]
assert captured["expected_opportunity_count"] is None
assert captured["require_opportunities"] is True
def test_snapshot_expected_count_validation_is_still_available(monkeypatch):
monkeypatch.setattr(simulator, "_load", lambda **kwargs: {"opportunities": [object()]})
with pytest.raises(RuntimeError, match="expected 328 opportunities, found 1"):
simulator.collect(expected_opportunity_count=328)
def test_empty_production_universe_is_rejected(monkeypatch):
monkeypatch.setattr(simulator, "collect", lambda **kwargs: {"opportunities": []})
assert projection._derive_all(
expected_database="clientflow", expected_user=None,
expected_opportunity_count=None, require_opportunities=True,
) == []
with pytest.raises(RuntimeError, match="at least one opportunity"):
projection.rebuild_blif_flow_v2_projection(
mode="shadow", derived_rows=[], require_opportunities=True,
)
def test_existing_test_defaults_remain_guarded(monkeypatch):
captured = {}
monkeypatch.setattr(simulator, "collect", lambda **kwargs: captured.update(kwargs) or {"opportunities": []})
projection._derive_all()
assert captured["expected_database"] == "clientflow_codex_test"
assert captured["expected_user"] == "clientflow_codex_test"
assert captured["expected_opportunity_count"] == 328

View File

@@ -0,0 +1,125 @@
from contextlib import nullcontext
import pytest
import scripts.simulate_blif_flow_v2 as simulator
class _Result:
def __init__(self, *, identity=None, rows=()):
self._identity = identity
self._rows = rows
def one(self):
return self._identity
def mappings(self):
return self._rows
class _Connection:
def __init__(self, database="clientflow", user="clientflow"):
self.database = database
self.user = user
self.read_only = False
self.statements = []
self.identity_observations = []
def execution_options(self, **_kwargs):
return self
def execute(self, statement):
sql = " ".join(str(statement).split())
self.statements.append(sql)
upper = sql.upper()
if upper == "BEGIN READ ONLY":
self.read_only = True
return _Result()
if upper == "ROLLBACK":
self.read_only = False
return _Result()
if upper.startswith(("INSERT ", "UPDATE ", "DELETE ")):
if self.read_only:
raise RuntimeError("cannot execute write in a read-only transaction")
return _Result()
if "CURRENT_DATABASE()" in upper:
identity = (self.database, self.user, "on" if self.read_only else "off")
self.identity_observations.append(identity)
return _Result(identity=identity)
return _Result(rows=())
class _Engine:
def __init__(self, connection):
self.connection = connection
def connect(self):
return nullcontext(self.connection)
def _load(monkeypatch, connection, **kwargs):
monkeypatch.setattr(simulator, "engine", _Engine(connection))
return simulator._load(**kwargs)
def test_strict_load_establishes_read_only_before_same_connection_identity_and_reads(monkeypatch):
connection = _Connection()
assert connection.read_only is False # normal fresh-connection state
result = _load(
monkeypatch, connection, expected_database="clientflow",
expected_user="clientflow", require_read_only=True,
)
assert connection.statements[0] == "BEGIN READ ONLY"
assert "CURRENT_DATABASE()" in connection.statements[1].upper()
assert connection.identity_observations == [("clientflow", "clientflow", "on")]
assert result["identity"]["transaction_read_only"] == "on"
assert connection.statements[2].upper().startswith("SELECT O.*")
assert connection.statements[-1] == "ROLLBACK"
@pytest.mark.parametrize(
("database", "user", "expected_database", "expected_user"),
[
("wrong", "clientflow", "clientflow", "clientflow"),
("clientflow", "wrong", "clientflow", "clientflow"),
],
)
def test_strict_load_refuses_wrong_identity_before_factual_reads(
monkeypatch, database, user, expected_database, expected_user,
):
connection = _Connection(database=database, user=user)
with pytest.raises(RuntimeError, match="unexpected database identity"):
_load(
monkeypatch, connection, expected_database=expected_database,
expected_user=expected_user, require_read_only=True,
)
assert connection.statements[0] == "BEGIN READ ONLY"
assert len([sql for sql in connection.statements if sql.upper().startswith("SELECT")]) == 1
assert connection.statements[-1] == "ROLLBACK"
def test_strict_factual_transaction_rejects_writes():
connection = _Connection()
connection.execute("BEGIN READ ONLY")
with pytest.raises(RuntimeError, match="read-only transaction"):
connection.execute("UPDATE opportunities SET stage = 'forbidden'")
assert connection.read_only is True
def test_non_strict_load_preserves_identity_then_read_only_collection_order(monkeypatch):
connection = _Connection(database="clientflow_codex_test", user="clientflow_codex_test")
result = _load(
monkeypatch, connection, expected_database="clientflow_codex_test",
expected_user="clientflow_codex_test", require_read_only=False,
)
assert "CURRENT_DATABASE()" in connection.statements[0].upper()
assert connection.identity_observations == [
("clientflow_codex_test", "clientflow_codex_test", "off")
]
assert connection.statements[1] == "BEGIN READ ONLY"
assert result["identity"]["transaction_read_only"] == "off"
assert connection.statements[-1] == "ROLLBACK"

View File

@@ -33,7 +33,6 @@ def test_pending_send_quote_task_is_not_sent_evidence():
def test_quote_sent_without_document_requires_reconciliation():
evidence = OpportunityEvidence(
opportunity_id="opp-1",
stage="QUOTE_SENT",
has_fiscal_customer=True,
fiscal_identity_validated=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():
evidence = OpportunityEvidence(
opportunity_id="opp-2",
stage="QUOTE_REQUESTED",
has_fiscal_customer=True,
fiscal_identity_validated=True,
fiscal_data_complete=True,
)
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():