10 Commits

32 changed files with 4932 additions and 36 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

@@ -8,6 +8,10 @@ from __future__ import annotations
from typing import Any
ACTION_LABELS = {
"CREATE_PROFORMA": "Criar proforma",
"CREATE_INVOICE": "Criar fatura",
"VALIDATE_ODOO_ORDER": "Validar encomenda",
"REVIEW_REQUIRED": "Rever processo",
"CALL_CUSTOMER": "Ligar ao cliente",
"SEND_INFO": "Enviar informação",
"SEND_QUOTE": "Preparar orçamento",
@@ -38,6 +42,10 @@ ACTION_LABELS = {
}
PRIMARY_ACTION_LABELS = {
"CREATE_PROFORMA": "Criar proforma",
"CREATE_INVOICE": "Criar fatura",
"VALIDATE_ODOO_ORDER": "Validar encomenda",
"REVIEW_REQUIRED": "Rever processo",
"CALL_CUSTOMER": "Ligar ao cliente",
"SEND_INFO": "Preparar resposta",
"SEND_QUOTE": "Preparar orçamento",

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,182 @@
"""Authoritative BLIF Flow v2 -> operator decision adapter.
The persisted projection owns factual business state. Pending tasks are only
operational obligations and may override that state when this policy proves
that they are current. This module is deliberately pure and performs no I/O.
"""
from __future__ import annotations
from dataclasses import asdict, dataclass
from datetime import datetime, timezone
from typing import Any, Iterable, Mapping
from app.admin_ui.labels import action_label
FORMAL_ACTIONS = {"CREATE_PROFORMA", "CREATE_INVOICE"}
REVIEW_ACTIONS = {"REVIEW", "REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS"}
PAYMENT_FOLLOWUPS = {"FOLLOW_UP_PAYMENT", "FOLLOW_UP_PROFORMA"}
VALID_OVERRIDES = REVIEW_ACTIONS | {"SUPPORT", "SEND_INFO", "CALL_CUSTOMER", "FOLLOW_UP_CUSTOMER_REVIEW"} | PAYMENT_FOLLOWUPS
TERMINAL_STATES = {"COMPLETED", "LOST", "NO_INTEREST"}
@dataclass(frozen=True)
class AuthoritativeOperationalDecision:
opportunity_id: str
canonical_opportunity_id: str
business_state: str
business_next_action: str | None
effective_action: str | None
queue: str
eligible: bool
reason_code: str
reason_text: str
blocking_action_code: str | None = None
obligation_source_refs: tuple[dict[str, Any], ...] = ()
confidence: str = "high"
def to_dict(self) -> dict[str, Any]:
result = asdict(self)
result["obligation_source_refs"] = list(self.obligation_source_refs)
# Shared legacy presentation contract consumed by Operations and detail.
result.update({
"action_code": self.effective_action or "NO_ACTION",
"label": action_label(self.effective_action, "Sem ação"),
"description": self.reason_text,
"priority": "alta" if self.queue in {"review", "blocked", "exception"} else "normal",
"can_execute": self.eligible and self.effective_action is not None,
"reason_if_blocked": self.reason_code if not self.eligible else None,
"operational_queue": self.queue,
"authoritative_v2": True,
"suppress_current_card": self.queue == "not_current",
})
return result
def _code(value: Any) -> str:
return str(value or "").strip().upper()
def _dt(value: Any) -> datetime | None:
if not value:
return None
if isinstance(value, datetime):
parsed = value
else:
try:
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except ValueError:
return None
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
def _active_obligations(obligations: Iterable[Mapping[str, Any]]) -> list[Mapping[str, Any]]:
return [row for row in obligations
if str(row.get("status") or "").lower() == "pending"
and not row.get("resolved_at") and not row.get("superseded_by_task_id")]
def _ref(row: Mapping[str, Any]) -> dict[str, Any]:
return {"source": "task", "id": str(row.get("id") or ""), "status": "pending",
"action_code": _code(row.get("action_code"))}
def decide_authoritative_operation(
projection: Mapping[str, Any] | None,
*, obligations: Iterable[Mapping[str, Any]] = (), now: datetime | None = None,
fiscal_complete: bool = True, reconciliation_blocking: bool = False,
hard_blocker: str | None = None,
) -> AuthoritativeOperationalDecision:
"""Return the sole operator decision, failing closed without a projection."""
now = now or datetime.now(timezone.utc)
if projection is None:
return AuthoritativeOperationalDecision(
"", "", "MISSING_PROJECTION", None, "REVIEW_REQUIRED", "review", True,
"MISSING_V2_PROJECTION", "A projeção Flow v2 está em falta; é necessária revisão, sem recorrer ao V1.",
confidence="low",
)
oid = str(projection.get("opportunity_id") or "")
canonical = str(projection.get("canonical_opportunity_id") or oid)
state = _code(projection.get("business_state"))
business_action = _code(projection.get("business_next_action")) or None
confidence = str(projection.get("confidence") or "low")
if projection.get("is_duplicate_representation"):
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, None, "not_current", False,
"DUPLICATE_SUPPRESSED", f"Representação duplicada do processo material canónico {canonical}.", confidence="high",
)
if hard_blocker:
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, hard_blocker, "blocked", True,
"HARD_FACTUAL_BLOCKER", "Um bloqueio factual ou de sistema impede a ação atual.", hard_blocker, confidence=confidence,
)
active = _active_obligations(obligations)
by_code: dict[str, list[Mapping[str, Any]]] = {}
for row in active:
by_code.setdefault(_code(row.get("action_code")), []).append(row)
# A formal-document prerequisite blocks only a transition which needs it.
if business_action in FORMAL_ACTIONS and not fiscal_complete:
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, "VALIDATE_FISCAL_CUSTOMER", "blocked", True,
"FISCAL_IDENTITY_REQUIRED", "Validar os dados fiscais antes de criar o documento oficial.",
"VALIDATE_FISCAL_CUSTOMER", confidence=confidence,
)
if business_action in FORMAL_ACTIONS and reconciliation_blocking:
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, "RECONCILE_DOCUMENTS", "blocked", True,
"DOCUMENT_RECONCILIATION_REQUIRED", "Confirmar a ligação do documento formal atual.",
"RECONCILE_DOCUMENTS", confidence=confidence,
)
def choose(codes: Iterable[str], queue: str, reason: str):
for code in codes:
rows = by_code.get(code, [])
if rows:
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, code, queue, True, reason,
"Existe uma obrigação operacional pendente e válida.",
obligation_source_refs=tuple(_ref(row) for row in rows), confidence=confidence,
)
return None
picked = choose(("REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW", "REVIEW_RECONSTRUCTED_PROCESS"), "review", "ACTIVE_REVIEW_OBLIGATION")
picked = picked or choose(("SUPPORT",), "do_now", "ACTIVE_SUPPORT_OBLIGATION")
picked = picked or choose(("SEND_INFO",), "do_now", "ACTIVE_RESPONSE_OBLIGATION")
picked = picked or choose(("CALL_CUSTOMER",), "do_now", "EXPLICIT_CALL_OBLIGATION")
if picked:
return picked
waiting_state = state in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}
followup_order = ("FOLLOW_UP_CUSTOMER_REVIEW",) if state == "AWAITING_CUSTOMER" else tuple(PAYMENT_FOLLOWUPS)
if waiting_state:
for code in followup_order:
rows = by_code.get(code, [])
if not rows:
continue
due_rows = [row for row in rows if _dt(row.get("due_at")) is None or _dt(row.get("due_at")) <= now]
if due_rows:
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, code, "do_now", True, "DUE_FOLLOW_UP",
"O follow-up ativo chegou à data e continua por satisfazer.",
obligation_source_refs=tuple(_ref(row) for row in due_rows), confidence=confidence,
)
return AuthoritativeOperationalDecision(
oid, canonical, state, business_action, None, "waiting", False, "FOLLOW_UP_NOT_DUE",
"O follow-up ativo ainda não chegou à data.",
obligation_source_refs=tuple(_ref(row) for row in rows), confidence=confidence,
)
if state in TERMINAL_STATES:
return AuthoritativeOperationalDecision(oid, canonical, state, business_action, None, "not_current", False,
"TERMINAL_BUSINESS_STATE", "O processo factual está concluído.", confidence=confidence)
if business_action:
queue = "review" if business_action in REVIEW_ACTIONS else "do_now"
return AuthoritativeOperationalDecision(oid, canonical, state, business_action, business_action, queue, True,
"BUSINESS_NEXT_ACTION", str(projection.get("reason_text") or "A ação decorre do estado factual Flow v2."), confidence=confidence)
if waiting_state:
return AuthoritativeOperationalDecision(oid, canonical, state, None, None, "waiting", False,
"WAITING_EXTERNAL_EVENT", "O processo aguarda um evento externo.", confidence=confidence)
return AuthoritativeOperationalDecision(oid, canonical, state, None, None, "backlog", False,
"NO_CURRENT_INTERNAL_ACTION", "Não existe ação interna atual.", confidence=confidence)

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

@@ -321,13 +321,14 @@ def _opportunity_item(
rows: list[Mapping[str, Any]],
decision: Mapping[str, Any],
) -> tuple[dict[str, Any] | None, list[dict[str, Any]]]:
authoritative = decision.get("authoritative_v2") is True
scheduled_call = next((
row for row in rows
if row.get("source") == "task"
and _s(row.get("status")).lower() == "pending"
and _canonical_code(row.get("action_code")) == "CALL_CUSTOMER"
), None)
if scheduled_call:
if scheduled_call and not authoritative:
decision = {
**dict(decision),
"action_code": "CALL_CUSTOMER",
@@ -347,26 +348,29 @@ def _opportunity_item(
evidence_rows.append(row)
if action_code in NO_WORK_ACTION_CODES:
if not _has_explicit_unresolved_review(evidence_rows):
if decision.get("suppress_current_card"):
return None, independent_exceptions
# The central decision normally wins, but an explicit unresolved review
# marker must not be hidden by NO_ACTION/NOT_FOUND.
action_code = "REVIEW"
decision = {
**dict(decision),
"action_code": action_code,
"label": "Rever evidência pendente",
"description": "A decisão central indica que não há ação, mas existe evidência explicitamente marcada para revisão.",
"reason": "Existe revisão humana não resolvida apesar da decisão central sem ação.",
"priority": "normal",
"target_url": f"/opportunities/{opportunity_id}",
}
keep_authoritative_noncurrent = authoritative and _s(decision.get("operational_queue")) in {"waiting", "backlog"}
if not keep_authoritative_noncurrent:
if not _has_explicit_unresolved_review(evidence_rows):
return None, independent_exceptions
# The central decision normally wins, but an explicit unresolved
# review marker must not be hidden by NO_ACTION/NOT_FOUND.
action_code = "REVIEW"
decision = {
**dict(decision), "action_code": action_code,
"label": "Rever evidência pendente",
"description": "A decisão central indica que não há ação, mas existe evidência explicitamente marcada para revisão.",
"reason": "Existe revisão humana não resolvida apesar da decisão central sem ação.",
"priority": "normal", "target_url": f"/opportunities/{opportunity_id}",
}
matching_task = scheduled_call or _matching_task(evidence_rows, action_code)
primary = matching_task or next((row for row in evidence_rows if row.get("source") == "task"), None)
primary = primary or (evidence_rows[0] if evidence_rows else {})
process_key = f"opportunity:{opportunity_id}"
waiting = _is_waiting(decision)
decision_queue = _s(decision.get("operational_queue")) if authoritative else ""
opportunity_row = next((row for row in evidence_rows if _s(row.get("fiscal_customer_name"))), None)
opportunity_row = opportunity_row or next((row for row in evidence_rows if _s(row.get("opportunity_title"))), None)
opportunity_row = opportunity_row or primary
@@ -383,7 +387,7 @@ def _opportunity_item(
"detail": _s(decision.get("description") or decision.get("reason")),
"why_human_required": _s(decision.get("reason") or decision.get("description")),
"priority": _s(decision.get("priority")) or _best_priority(evidence_rows),
"operational_queue": "waiting" if waiting else "do_now",
"operational_queue": decision_queue or ("waiting" if waiting else "do_now"),
"queue": _s((matching_task or primary).get("queue")) or _s(primary.get("queue")) or "rever",
"status": "waiting" if waiting else "pending",
"href": _s((matching_task or {}).get("href")) or _s(decision.get("target_url")) or f"/opportunities/{opportunity_id}",
@@ -400,6 +404,10 @@ def _opportunity_item(
],
"decision_version": decision.get("decision_version"),
"decision": dict(decision),
"business_state": decision.get("business_state"),
"business_next_action": decision.get("business_next_action"),
"effective_action": decision.get("effective_action"),
"canonical_opportunity_id": decision.get("canonical_opportunity_id") or opportunity_id,
})
return item, independent_exceptions

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

@@ -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

@@ -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

@@ -157,6 +157,23 @@ def apply_operational_eligibility(items: Iterable[Mapping[str, Any]], *, now: da
projected = []
for source in items:
item = dict(source)
decision = item.get("decision") if isinstance(item.get("decision"), Mapping) else {}
if decision.get("authoritative_v2") is True:
item["eligibility"] = {
"eligible": bool(decision.get("eligible")),
"queue": decision.get("operational_queue"),
"reason_code": decision.get("reason_code"),
"reason_text": decision.get("description"),
"blocking_action_code": decision.get("blocking_action_code"),
"obligation_source_refs": decision.get("obligation_source_refs") or [],
"confidence": decision.get("confidence"),
}
item["operational_queue"] = decision.get("operational_queue")
item["why_human_required"] = decision.get("description")
item["eligibility_reason_code"] = decision.get("reason_code")
item["eligible"] = bool(decision.get("eligible"))
projected.append(item)
continue
eligibility = evaluate_operational_eligibility(item, now=now)
item["eligibility"] = eligibility.to_dict()
item["operational_queue"] = eligibility.queue

View File

@@ -8,6 +8,7 @@ from __future__ import annotations
import os
import re
import subprocess
from contextlib import nullcontext
from typing import Any, Dict, List, Optional
from sqlalchemy import text
@@ -26,7 +27,7 @@ from app.canonical_operations import (
opportunity_ids_from_work_seeds,
partition_canonical_items,
)
from app.opportunity_next_action_service import get_opportunity_next_actions
from app.opportunity_next_action_service import get_opportunity_next_actions, get_opportunity_next_actions_for_mode
from app.document_reconciliation_service import active_document_link_exclusion_sql
@@ -104,7 +105,7 @@ def _norm_identity(value: Any) -> str:
def _identity_tokens(value: Any) -> set[str]:
return {
result = {
token
for token in _norm_identity(value).split()
if len(token) >= 3 and token not in _IDENTITY_STOPWORDS
@@ -302,7 +303,9 @@ def _attach_operation_urls(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
return _sanitize_operation_identities(cleaned)
def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
def get_operations_summary(
limit: int = 24, *, flow_mode: str | None = None, connection: Any = None,
) -> Dict[str, Any]:
"""Build the /operations work queue summary.
/operations is intentionally not a mini-dashboard. It returns a compact
@@ -315,7 +318,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
# after canonicalization. This prevents duplicate rows from consuming the
# display limit while still avoiding an unbounded history scan.
candidate_limit = max(200, min(display_limit * 20, 1000))
with engine.begin() as conn:
active_flow_mode = str(flow_mode if flow_mode is not None else settings.blif_flow_v2_mode or "off").strip().lower()
with (nullcontext(connection) if connection is not None else engine.begin()) as conn:
counts = conn.execute(text("""
SELECT
(SELECT COUNT(*) FROM opportunities WHERE status = 'open')::int AS open_opportunities,
@@ -672,12 +676,50 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
conversation_timeline = {str(row["conversation_id"]): dict(row) for row in timeline_rows}
seed_rows = [dict(r) for r in work_seed_rows]
# In authoritative mode the factual projection, not a legacy task/stage,
# selects material processes that have current business work. Synthetic
# rows are read-model seeds only and are never persisted.
if active_flow_mode == "authoritative":
seeded = {str(row.get("opportunity_id") or "") for row in seed_rows}
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
v2_seeds = conn.execute(text("""
SELECT p.opportunity_id::text, p.business_next_action, p.derived_at,
o.title, o.value_amount, o.currency, o.lifecycle_state,
c.name AS customer_name
FROM opportunity_flow_state_v2 p
JOIN opportunities o ON o.id = p.opportunity_id
LEFT JOIN customers c ON c.id = o.local_customer_id
WHERE p.is_duplicate_representation = false
AND (p.business_next_action IS NOT NULL
OR p.business_state IN ('AWAITING_CUSTOMER', 'AWAITING_PAYMENT'))
""")).mappings().all()
for projection_seed in v2_seeds:
oid = str(projection_seed["opportunity_id"])
if oid in seeded:
continue
seed_rows.append({
"source": "flow_v2", "id": oid, "opportunity_id": oid,
"created_at": projection_seed.get("derived_at"), "due_at": None,
"priority": "normal", "queue": "vendas",
"action_code": projection_seed.get("business_next_action"), "status": "projected",
"title": projection_seed.get("title") or "Oportunidade",
"opportunity_title": projection_seed.get("title") or "",
"customer_name": projection_seed.get("customer_name") or "",
"fiscal_customer_name": projection_seed.get("customer_name") or "",
"item_metadata": {}, "opportunity_metadata": {},
"opportunity_lifecycle_state": projection_seed.get("lifecycle_state") or "",
"opportunity_value_amount": projection_seed.get("value_amount") or 0,
"opportunity_currency": projection_seed.get("currency") or "EUR",
"href": f"/opportunities/{oid}", "action_label": "Abrir",
})
for row in seed_rows:
timeline = conversation_timeline.get(str(row.get("conversation_id") or ""), {})
row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")})
normalized_rows = _normalise_work_item_intent(seed_rows)
opportunity_ids = opportunity_ids_from_work_seeds(normalized_rows)
decisions = get_opportunity_next_actions(opportunity_ids) if opportunity_ids else {}
decisions = (get_opportunity_next_actions_for_mode(
opportunity_ids, flow_mode=active_flow_mode, connection=connection,
) if opportunity_ids else {})
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
canonical_items = _attach_operation_urls(list(projection["items"]))
partition = partition_canonical_items(canonical_items, display_limit=display_limit)
@@ -722,6 +764,11 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
**diagnostics,
"projection_metrics": diagnostics,
}
if flow_mode is not None:
# Complete, untruncated read-model population for offline simulation;
# normal API/GET callers do not receive this audit-only field.
result["simulation_all_items"] = canonical_items
return result
def get_system_health_summary() -> Dict[str, Any]:

View File

@@ -7,11 +7,15 @@ 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
from app.authoritative_operational_adapter import decide_authoritative_operation
# 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 +25,9 @@ from app.domain.opportunity_flow import (
)
logger = logging.getLogger(__name__)
@dataclass
class OpportunityNextAction:
action_code: str
@@ -37,6 +44,70 @@ 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 {}
# Lightweight repository fakes used by pure unit tests intentionally expose
# only ``begin``. Comparison observation is optional and must not alter the
# V1 contract or its query count when that read capability is absent.
if not hasattr(engine, "connect"):
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 != "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 +231,8 @@ def get_opportunity_next_action(
decision is now produced by the company workflow engine.
"""
if _flow_v2_mode() == "authoritative":
return get_opportunity_next_actions([opportunity_id]).get(opportunity_id) or decide_authoritative_operation(None).to_dict()
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 +245,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):
@@ -302,10 +376,24 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
if not ids:
return {}
return get_opportunity_next_actions_for_mode(ids, flow_mode=_flow_v2_mode())
def get_opportunity_next_actions_for_mode(
opportunity_ids: Iterable[str], *, flow_mode: str, connection: Any = None,
) -> Dict[str, Dict[str, Any]]:
"""Explicit-mode read boundary used by the read-only cutover simulation."""
ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip()))
if not ids:
return {}
if flow_mode == "authoritative":
return _get_authoritative_opportunity_decisions(ids, connection=connection)
# Read path only: schema creation belongs to startup/migrations. In
# particular, GET /operations must never perform DDL while calculating its
# canonical projection.
with engine.begin() as conn:
from contextlib import nullcontext
with (nullcontext(connection) if connection is not None else engine.begin()) as conn:
opportunities = _bulk_rows(conn, """
SELECT id::text, stage, status, title,
local_customer_id::text AS fiscal_customer_id,
@@ -376,4 +464,60 @@ 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()
if flow_mode == "compare":
_observe_flow_v2(decisions)
return decisions
def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any = None) -> Dict[str, Dict[str, Any]]:
"""Read-only set-oriented adapter input loader; never derives or persists V2."""
from contextlib import nullcontext
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
projections = _bulk_rows(conn, """
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
FROM opportunity_flow_state_v2 WHERE opportunity_id::text IN :opportunity_ids
""", ids)
tasks = _bulk_rows(conn, """
SELECT id::text, opportunity_id::text, action_code, status, due_at,
resolved_at, superseded_by_task_id::text, created_at, metadata
FROM tasks WHERE opportunity_id::text IN :opportunity_ids AND status = 'pending'
""", ids)
customers = _bulk_rows(conn, """
SELECT o.id::text AS opportunity_id,
(c.tax_id IS NOT NULL AND c.tax_id <> '' AND c.email IS NOT NULL AND c.email <> ''
AND c.street_name IS NOT NULL AND c.street_name <> ''
AND c.postal_zone IS NOT NULL AND c.postal_zone <> ''
AND c.city_name IS NOT NULL AND c.city_name <> '') AS fiscal_complete
FROM opportunities o LEFT JOIN customers c ON c.id = o.local_customer_id
WHERE o.id::text IN :opportunity_ids
""", ids)
reconciliation = _bulk_rows(conn, """
SELECT opportunity_id::text
FROM reconciliation_items
WHERE opportunity_id::text IN :opportunity_ids
AND status IN ('needs_review','conflict')
AND COALESCE((payload->>'reconciliation_blocks_current_action')::boolean, false)
""", ids)
projection_by_id = {str(row["opportunity_id"]): row for row in projections}
tasks_by_id = _group_by_opportunity(tasks)
fiscal_by_id = {str(row["opportunity_id"]): bool(row.get("fiscal_complete")) for row in customers}
reconciliation_ids = {str(row["opportunity_id"]) for row in reconciliation}
decisions = {}
for oid in ids:
projection = projection_by_id.get(oid)
decision = decide_authoritative_operation(
projection, obligations=tasks_by_id.get(oid, []),
fiscal_complete=fiscal_by_id.get(oid, False),
reconciliation_blocking=oid in reconciliation_ids,
)
value = decision.to_dict()
if projection is None:
value["opportunity_id"] = oid
value["target_url"] = f"/opportunities/{oid}"
else:
value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}"
decisions[oid] = value
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,226 @@
#!/usr/bin/env python3
"""Simulate authoritative BLIF Flow v2 without changing application mode."""
from __future__ import annotations
import argparse
import json
import os
from pathlib import Path
import sys
from typing import Any
JSON_PATH = Path("/tmp/blif_flow_v2_authoritative_simulation.json")
TEXT_PATH = Path("/tmp/blif_flow_v2_authoritative_operations.txt")
CURRENT = {"do_now", "review", "blocked", "exception"}
CLASSIFICATIONS = {
"SATISFIED", "SUPERSEDED", "DUPLICATE_SUPPRESSED", "PREMATURE_REMOVED",
"REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE",
}
NAMED = {
"Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "do_now"),
"Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
"ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "do_now"),
"CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "do_now"),
"X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, "not_current"),
"X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, "suppressed"),
"RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "review"),
"RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", None, "suppressed"),
}
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(
description="Read-only V1 versus simulated authoritative BLIF Flow v2 Operations audit."
)
parser.add_argument(
"--production-readonly-simulation", action="store_true",
help="Required opt-in when current_database() is clientflow; requires compare mode and a read-only transaction.",
)
return parser.parse_args(argv)
def _load_runtime():
"""Application imports are intentionally delayed until after argparse."""
root = Path(__file__).resolve().parents[1]
if str(root) not in sys.path:
sys.path.insert(0, str(root))
from app.db import engine
from app.operations_service import get_operations_summary
from app.opportunity_next_action_service import get_opportunity_next_actions_for_mode
return engine, get_operations_summary, get_opportunity_next_actions_for_mode
def _all(summary: dict[str, Any]) -> list[dict[str, Any]]:
if "simulation_all_items" in summary:
return list(summary.get("simulation_all_items") or [])
return sum((list(summary.get(key) or []) for key in
("work_items", "waiting_items", "backlog_items", "not_current_items")), [])
def _metrics(summary: dict[str, Any]) -> dict[str, int]:
queues = [str(item.get("operational_queue") or "") for item in _all(summary)]
return {"current": sum(q in CURRENT for q in queues), "do_now": queues.count("do_now"),
"review": queues.count("review"), "waiting": queues.count("waiting"),
"backlog": queues.count("backlog"), "blocked": queues.count("blocked"),
"not_current": queues.count("not_current")}
def _key(item: dict[str, Any]) -> str:
return str(item.get("work_item_key") or item.get("process_key") or item.get("opportunity_id") or item.get("id"))
def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[str, str]:
if new is not None:
return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item."
decision = old.get("decision") if isinstance(old.get("decision"), dict) else {}
code = str(decision.get("reason_code") or old.get("eligibility_reason_code") or "").upper()
if "DUPLICATE" in code:
return "DUPLICATE_SUPPRESSED", "The local opportunity is a duplicate material representation."
if "SATISF" in code or "ANSWERED" in code:
return "SATISFIED", "Later factual evidence satisfies the old obligation."
if "SUPERSE" in code or "TERMINAL" in code:
return "SUPERSEDED", "A later factual state supersedes the legacy card."
if "PREMATURE" in code:
return "PREMATURE_REMOVED", "The legacy action is premature for the factual state."
if old.get("source") in {"task", "opportunity"}:
return "LEGACY_ONLY", "The card is supported only by legacy operational representation."
return "UNSAFE_FALSE_NEGATIVE", "No safe factual explanation was found for removing this current card."
def build_report(
*, identity: dict[str, Any], configured_mode: str, v1: dict[str, Any], v2: dict[str, Any],
named_decisions: dict[str, dict[str, Any]], missing_projections: int,
) -> dict[str, Any]:
v1_current = {_key(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT}
v2_current = {_key(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT}
removed, changed = [], []
for key, old in v1_current.items():
new = v2_current.get(key)
if new is not None and new.get("action_code") == old.get("action_code"):
continue
classification, reason = _classification(old, new)
row = {"work_item_key": key, "opportunity_id": old.get("opportunity_id"),
"v1_action": old.get("action_code"), "simulated_v2_action": (new or {}).get("action_code"),
"classification": classification, "reason": reason,
"source_refs": old.get("source_refs") or []}
(changed if new is not None else removed).append(row)
added = [{"work_item_key": key, "opportunity_id": item.get("opportunity_id"),
"action": item.get("action_code"), "source_refs": item.get("source_refs") or []}
for key, item in v2_current.items() if key not in v1_current]
material_counts: dict[str, int] = {}
for item in v2_current.values():
material = str(item.get("canonical_opportunity_id") or item.get("opportunity_id") or item.get("process_key"))
material_counts[material] = material_counts.get(material, 0) + 1
duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1]
named_cases = {}
for name, (oid, expected_action, expected_queue) in NAMED.items():
decision = named_decisions.get(oid) or {}
suppressed = decision.get("suppress_current_card") is True
actual_queue = "suppressed" if suppressed else decision.get("operational_queue")
actual_action = decision.get("effective_action")
passed = actual_action == expected_action and actual_queue == expected_queue
named_cases[name] = {"opportunity_id": oid, "expected_action": expected_action,
"expected_queue": expected_queue, "actual_action": actual_action,
"actual_queue": actual_queue, "passed": passed}
fiscal = []
for item in v2_current.values():
decision = item.get("decision") if isinstance(item.get("decision"), dict) else {}
if item.get("action_code") == "VALIDATE_FISCAL_CUSTOMER":
fiscal.append({"opportunity_id": item.get("opportunity_id"),
"business_action": decision.get("business_next_action"),
"reason_code": decision.get("reason_code"),
"reason": decision.get("description")})
unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed)
return {
"identity": identity, "configured_mode": configured_mode,
"v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2),
"semantic_changes": len(removed) + len(changed) + len(added),
"removed_cards": removed, "added_cards": added, "changed_actions": changed,
"duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates),
"missing_projections": missing_projections, "unsafe_false_negatives": unsafe,
"named_cases": named_cases, "fiscal_validation_cases": fiscal,
}
def _write_reports(report: dict[str, Any]) -> None:
JSON_PATH.write_text(json.dumps(report, indent=2, ensure_ascii=False, default=str) + "\n")
lines = ["BLIF Flow v2 authoritative Operations simulation", "",
f"Identity: {report['identity']}", f"Configured mode: {report['configured_mode']}",
f"V1: {report['v1_metrics']}", f"Simulated authoritative V2: {report['authoritative_metrics']}",
f"Removed: {len(report['removed_cards'])}", f"Added: {len(report['added_cards'])}",
f"Changed actions: {len(report['changed_actions'])}",
f"Duplicate current cards: {report['duplicate_current_cards']}",
f"Missing projections: {report['missing_projections']}",
f"UNSAFE_FALSE_NEGATIVE: {report['unsafe_false_negatives']}", "", "Named cases:"]
lines.extend(f"- {name}: {case}" for name, case in report["named_cases"].items())
lines.extend(["", f"Fiscal validation cases: {len(report['fiscal_validation_cases'])}"])
TEXT_PATH.write_text("\n".join(lines) + "\n")
def exit_code(report: dict[str, Any]) -> int:
failed_named = any(not case.get("passed") for case in report.get("named_cases", {}).values())
unsafe = int(report.get("unsafe_false_negatives") or 0)
missing = int(report.get("missing_projections") or 0)
duplicates = int(report.get("duplicate_current_cards") or 0)
return 2 if unsafe or missing or duplicates or failed_named else 0
def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, Any]:
engine, get_operations, get_decisions = runtime_loader()
configured_mode = os.environ.get("BLIF_FLOW_V2_MODE")
conn = engine.connect()
try:
# One explicit transaction and one factual snapshot for identity, V1,
# V2, projection completeness, and named-case decisions.
conn.exec_driver_sql("BEGIN READ ONLY")
identity = dict(conn.exec_driver_sql("""
SELECT current_database() AS database, current_user AS user,
current_setting('transaction_read_only') AS transaction_read_only
""").mappings().one())
production = identity.get("database") == "clientflow"
if production and not args.production_readonly_simulation:
raise RuntimeError("production simulation requires --production-readonly-simulation")
if args.production_readonly_simulation:
if identity.get("database") != "clientflow":
raise RuntimeError("production simulation requires current_database() = clientflow")
if identity.get("transaction_read_only") != "on":
raise RuntimeError("production simulation requires transaction_read_only = on")
if configured_mode is None:
raise RuntimeError("production simulation requires explicit BLIF_FLOW_V2_MODE=compare")
if configured_mode.strip().lower() != "compare":
raise RuntimeError("production simulation requires BLIF_FLOW_V2_MODE=compare")
v1 = get_operations(limit=200, flow_mode="off", connection=conn)
v2 = get_operations(limit=200, flow_mode="authoritative", connection=conn)
projection_counts = conn.exec_driver_sql("""
SELECT (SELECT count(*) FROM opportunities) AS opportunities,
(SELECT count(*) FROM opportunity_flow_state_v2) AS projections
""").mappings().one()
named_ids = [value[0] for value in NAMED.values()]
named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn)
if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode:
raise RuntimeError("configured BLIF_FLOW_V2_MODE changed during simulation")
report = build_report(
identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2,
named_decisions=named_decisions,
missing_projections=max(0, int(projection_counts["opportunities"]) - int(projection_counts["projections"])),
)
_write_reports(report)
return report
finally:
conn.exec_driver_sql("ROLLBACK")
conn.close()
def main(argv: list[str] | None = None, *, runtime_loader=_load_runtime) -> int:
args = parse_args(argv) # --help exits before runtime_loader/application imports.
report = run(args, runtime_loader=runtime_loader)
print(json.dumps(report, indent=2, ensure_ascii=False, default=str))
return exit_code(report)
if __name__ == "__main__":
raise SystemExit(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

@@ -0,0 +1,99 @@
from datetime import datetime, timedelta, timezone
from app.authoritative_operational_adapter import decide_authoritative_operation
from app.canonical_operations import canonicalize_operations
NOW = datetime(2026, 8, 16, tzinfo=timezone.utc)
def projection(state="AWAITING_CUSTOMER", action=None, **extra):
return {"opportunity_id": "opp", "canonical_opportunity_id": "opp",
"business_state": state, "business_next_action": action,
"confidence": "high", **extra}
def task(code, due=NOW, **extra):
return {"id": code, "status": "pending", "action_code": code, "due_at": due, **extra}
def test_business_baselines_and_terminal():
for state, action in [("PROFORMA_REQUIRED", "CREATE_PROFORMA"),
("INVOICE_REQUIRED", "CREATE_INVOICE"),
("ODOO_ORDER_REQUIRED", "PREPARE_ORDER"),
("ODOO_ORDER_CREATED", "VALIDATE_ODOO_ORDER")]:
got = decide_authoritative_operation(projection(state, action), now=NOW)
assert (got.effective_action, got.queue) == (action, "do_now")
assert decide_authoritative_operation(projection("COMPLETED"), now=NOW).queue == "not_current"
def test_missing_projection_fails_closed_without_v1():
got = decide_authoritative_operation(None, now=NOW)
assert (got.effective_action, got.queue, got.reason_code) == ("REVIEW_REQUIRED", "review", "MISSING_V2_PROJECTION")
def test_due_and_future_customer_followup_semantics():
due = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW")], now=NOW)
assert (due.effective_action, due.queue) == ("FOLLOW_UP_CUSTOMER_REVIEW", "do_now")
future = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW", NOW + timedelta(days=1))], now=NOW)
assert (future.effective_action, future.queue) == (None, "waiting")
assert decide_authoritative_operation(projection(), now=NOW).queue == "waiting"
def test_only_active_non_superseded_tasks_override_and_precedence():
stale = task("SUPPORT", resolved_at=NOW)
assert decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), obligations=[stale], now=NOW).effective_action == "CREATE_PROFORMA"
obligations = [task("CALL_CUSTOMER"), task("SEND_INFO"), task("SUPPORT"), task("REVIEW_MANUALLY")]
assert decide_authoritative_operation(projection(), obligations=obligations, now=NOW).effective_action == "REVIEW_MANUALLY"
assert decide_authoritative_operation(projection(), obligations=[task("SUPPORT")], now=NOW).effective_action == "SUPPORT"
assert decide_authoritative_operation(projection(), obligations=[task("SEND_INFO")], now=NOW).effective_action == "SEND_INFO"
assert decide_authoritative_operation(projection(), obligations=[task("CALL_CUSTOMER")], now=NOW).effective_action == "CALL_CUSTOMER"
def test_fiscal_and_reconciliation_only_block_formal_action():
inquiry = decide_authoritative_operation(projection("INQUIRY", "SEND_INFO"), fiscal_complete=False, reconciliation_blocking=True, now=NOW)
assert inquiry.effective_action == "SEND_INFO"
formal = decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), fiscal_complete=False, now=NOW)
assert (formal.effective_action, formal.blocking_action_code) == ("VALIDATE_FISCAL_CUSTOMER", "VALIDATE_FISCAL_CUSTOMER")
def test_duplicate_is_suppressed_even_with_pending_legacy_task():
decision = decide_authoritative_operation(projection("REVIEW_REQUIRED", "REVIEW_REQUIRED",
is_duplicate_representation=True,
canonical_opportunity_id="canonical"),
obligations=[task("REVIEW_MANUALLY")], now=NOW).to_dict()
source = {"source": "task", "id": "t", "status": "pending", "action_code": "REVIEW_MANUALLY",
"opportunity_id": "opp", "created_at": NOW}
assert canonicalize_operations([source], {"opp": decision})["items"] == []
def test_authoritative_queue_and_action_survive_canonical_boundary():
decision = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW")], now=NOW).to_dict()
source = {"source": "task", "id": "old", "status": "pending", "action_code": "SEND_INVOICE",
"opportunity_id": "opp", "created_at": NOW}
item = canonicalize_operations([source], {"opp": decision})["items"][0]
assert (item["action_code"], item["operational_queue"]) == ("FOLLOW_UP_CUSTOMER_REVIEW", "do_now")
def test_named_case_acceptance_decisions():
cases = {
"5c33db95-fab8-477a-bddd-0b9cc8f91302": ("PROFORMA_REQUIRED", "CREATE_PROFORMA", [], "CREATE_PROFORMA", "do_now"),
"fd79b9a1-07e6-4f61-95e8-09eab89c155e": ("AWAITING_CUSTOMER", None, [task("FOLLOW_UP_CUSTOMER_REVIEW")], "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
"61f1c955-a372-4ea7-b9b0-b8528d74a141": ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER", [], "PREPARE_ORDER", "do_now"),
"e3b23ac5-84db-4763-8a31-a684e873032c": ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER", [], "PREPARE_ORDER", "do_now"),
"dc89a466-db24-401b-bfe9-d47644b2d0c8": ("COMPLETED", None, [], None, "not_current"),
"fd221608-e007-4043-a23d-07e0c119a345": ("REVIEW_REQUIRED", "REVIEW_REQUIRED", [], "REVIEW_REQUIRED", "review"),
}
for oid, (state, action, obligations, effective, queue) in cases.items():
got = decide_authoritative_operation({**projection(state, action), "opportunity_id": oid,
"canonical_opportunity_id": oid},
obligations=obligations, now=NOW)
assert (got.effective_action, got.queue) == (effective, queue)
for duplicate, canonical in [
("1816a06e-9a69-4a9b-9279-1263156892d3", "dc89a466-db24-401b-bfe9-d47644b2d0c8"),
("434124fb-ac19-4d78-909a-55761d7e8daa", "fd221608-e007-4043-a23d-07e0c119a345"),
]:
got = decide_authoritative_operation({**projection("REVIEW_REQUIRED", "REVIEW_REQUIRED"),
"opportunity_id": duplicate, "canonical_opportunity_id": canonical,
"is_duplicate_representation": True}, now=NOW)
assert (got.effective_action, got.queue, got.canonical_opportunity_id) == (None, "not_current", canonical)

View File

@@ -0,0 +1,138 @@
from argparse import Namespace
import importlib.util
from pathlib import Path
import pytest
PATH = Path("scripts/simulate_blif_flow_v2_authoritative.py")
SPEC = importlib.util.spec_from_file_location("authoritative_simulation", PATH)
sim = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(sim)
class Result:
def __init__(self, row): self.row = row
def mappings(self): return self
def one(self): return self.row
class Connection:
def __init__(self, database="clientflow", readonly="on"):
self.database, self.readonly, self.sql = database, readonly, []
def exec_driver_sql(self, sql):
self.sql.append(sql.strip())
if "current_database" in sql:
return Result({"database": self.database, "user": "audit", "transaction_read_only": self.readonly})
if "count(*) FROM opportunities" in sql:
return Result({"opportunities": 328, "projections": 328})
return Result({})
def close(self): pass
class Engine:
def __init__(self, connection): self.connection = connection
def connect(self): return self.connection
def named_decisions(fail=False):
result = {}
for name, (oid, action, queue) in sim.NAMED.items():
result[oid] = {"effective_action": action,
"operational_queue": "not_current" if queue == "suppressed" else queue,
"suppress_current_card": queue == "suppressed"}
if fail:
result[sim.NAMED["Instalbeira"][0]]["effective_action"] = "SEND_INVOICE"
return result
def runtime(connection, *, fail_named=False, modes=None):
empty = {"work_items": [], "waiting_items": [], "backlog_items": [], "not_current_items": []}
def operations(**kwargs):
assert kwargs["connection"] is connection
if modes is not None:
modes.append(kwargs["flow_mode"])
return empty
def decisions(ids, **kwargs):
assert kwargs == {"flow_mode": "authoritative", "connection": connection}
return named_decisions(fail_named)
return lambda: (Engine(connection), operations, decisions)
def test_help_performs_zero_db_work():
called = False
def loader():
nonlocal called
called = True
raise AssertionError("DB/application runtime must not load")
with pytest.raises(SystemExit) as exc:
sim.main(["--help"], runtime_loader=loader)
assert exc.value.code == 0 and called is False
def test_production_requires_explicit_flag(monkeypatch):
monkeypatch.setenv("BLIF_FLOW_V2_MODE", "compare")
with pytest.raises(RuntimeError, match="requires --production"):
sim.run(Namespace(production_readonly_simulation=False), runtime_loader=runtime(Connection()))
@pytest.mark.parametrize("database,readonly,mode,message", [
("test_db", "on", "compare", "current_database"),
("clientflow", "off", "compare", "transaction_read_only"),
("clientflow", "on", None, "explicit BLIF_FLOW"),
("clientflow", "on", "off", "BLIF_FLOW_V2_MODE=compare"),
("clientflow", "on", "shadow", "BLIF_FLOW_V2_MODE=compare"),
("clientflow", "on", "authoritative", "BLIF_FLOW_V2_MODE=compare"),
])
def test_production_guards(monkeypatch, database, readonly, mode, message):
if mode is None:
monkeypatch.delenv("BLIF_FLOW_V2_MODE", raising=False)
else:
monkeypatch.setenv("BLIF_FLOW_V2_MODE", mode)
with pytest.raises(RuntimeError, match=message):
sim.run(Namespace(production_readonly_simulation=True),
runtime_loader=runtime(Connection(database, readonly)))
def test_safe_simulation_uses_explicit_modes_without_changing_config_and_has_no_db_writes(monkeypatch, tmp_path):
monkeypatch.setenv("BLIF_FLOW_V2_MODE", "compare")
monkeypatch.setattr(sim, "JSON_PATH", tmp_path / "report.json")
monkeypatch.setattr(sim, "TEXT_PATH", tmp_path / "report.txt")
connection = Connection()
modes = []
report = sim.run(Namespace(production_readonly_simulation=True), runtime_loader=runtime(connection, modes=modes))
assert report["configured_mode"] == "compare"
assert sim.exit_code(report) == 0
assert [sql for sql in connection.sql if sql.split(None, 1)[0].upper() in {"INSERT", "UPDATE", "DELETE", "MERGE"}] == []
assert connection.sql[0] == "BEGIN READ ONLY" and connection.sql[-1] == "ROLLBACK"
assert modes == ["off", "authoritative"]
def safe_report():
return sim.build_report(identity={}, configured_mode="compare",
v1={"work_items": []}, v2={"work_items": []},
named_decisions=named_decisions(), missing_projections=0)
def test_unsafe_false_negative_is_nonzero():
report = safe_report()
report["unsafe_false_negatives"] = 1
assert sim.exit_code(report) != 0
def test_missing_projection_is_nonzero():
report = safe_report()
report["missing_projections"] = 1
assert sim.exit_code(report) != 0
def test_duplicate_current_card_is_nonzero():
report = safe_report()
report["duplicate_current_cards"] = 1
assert sim.exit_code(report) != 0
def test_named_case_failure_is_nonzero():
report = sim.build_report(identity={}, configured_mode="compare", v1={"work_items": []},
v2={"work_items": []}, named_decisions=named_decisions(True), missing_projections=0)
assert sim.exit_code(report) != 0

View File

@@ -0,0 +1,88 @@
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_empty_bulk_is_safe(monkeypatch):
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "authoritative")
assert 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

@@ -8,6 +8,12 @@ import app.opportunity_next_action_service as service
from app.domain.opportunity_flow import build_opportunity_evidence
@pytest.fixture(autouse=True)
def _v1_contract_mode(monkeypatch):
"""These tests measure the legacy bulk contract, independently of env mode."""
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "off")
class _Result:
def __init__(self, rows):
self._rows = rows