6 Commits

15 changed files with 1627 additions and 230 deletions

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

@@ -0,0 +1,196 @@
"""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,
"is_duplicate_representation": self.reason_code == "DUPLICATE_SUPPRESSED",
"suppress_current_card": self.reason_code == "DUPLICATE_SUPPRESSED",
})
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 _unanswered_inbound(row: Mapping[str, Any]) -> bool:
inbound = _dt(row.get("latest_public_inbound"))
outbound = _dt(row.get("latest_public_outbound"))
return bool(inbound and (outbound is None or outbound <= inbound))
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")),
"source_system": str(row.get("source_system") or "")}
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,
)
# Terminal factual state cannot be reopened by a leftover operational task.
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)
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")
valid_send_info = [row for row in by_code.get("SEND_INFO", []) if (
state != "AWAITING_CUSTOMER" or _unanswered_inbound(row)
)]
if valid_send_info and (state in {"INQUIRY", "AWAITING_CUSTOMER"} or business_action == "SEND_INFO"):
by_code["SEND_INFO"] = valid_send_info
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 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

@@ -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,11 @@ 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,
"material_process_key": decision.get("material_process_key"),
})
return item, independent_exceptions

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
@@ -37,6 +38,53 @@ def _int(value: Any) -> int:
return 0
def _authoritative_material_identity(item: Dict[str, Any]) -> str:
return str(item.get("material_process_key") or item.get("canonical_opportunity_id")
or item.get("process_key") or item.get("work_item_key") or item.get("id"))
def _collapse_authoritative_material_processes(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
"""Enforce one current card per material process without deleting evidence."""
current_queues = {"do_now", "review", "blocked", "exception"}
action_rank = {
"REVIEW_MANUALLY": 10, "REVIEW_RECONSTRUCTED_PROCESS": 11, "REVIEW_REQUIRED": 12,
"SUPPORT": 20, "SEND_INFO": 30, "CALL_CUSTOMER": 40,
"FOLLOW_UP_CUSTOMER_REVIEW": 50, "FOLLOW_UP_PAYMENT": 51,
}
groups: Dict[str, List[Dict[str, Any]]] = {}
noncurrent: List[Dict[str, Any]] = []
for item in items:
if item.get("operational_queue") not in current_queues:
noncurrent.append(item)
continue
groups.setdefault(_authoritative_material_identity(item), []).append(item)
selected: List[Dict[str, Any]] = []
for identity, group in groups.items():
winner = min(group, key=lambda row: (
0 if row.get("operational_queue") in {"blocked", "exception"} else 1,
action_rank.get(str(row.get("action_code") or ""), 100),
str(row.get("work_item_key") or ""),
))
merged = dict(winner)
if len(group) > 1:
merged["suppressed_competing_actions"] = [
{"work_item_key": row.get("work_item_key"), "action_code": row.get("action_code"),
"source_refs": row.get("source_refs") or []}
for row in group if row is not winner
]
refs, seen = list(merged.get("source_refs") or []), set()
for ref in refs:
seen.add((str(ref.get("source")), str(ref.get("id"))))
for row in group:
for ref in row.get("source_refs") or []:
key = (str(ref.get("source")), str(ref.get("id")))
if key not in seen:
refs.append(ref); seen.add(key)
merged["source_refs"] = refs
selected.append(merged)
return selected + noncurrent
def _is_manually_resolved_error(value: Any) -> bool:
text_value = str(value or "").strip().casefold()
return "limpo manualmente" in text_value or "resolvido manualmente" in text_value
@@ -104,7 +152,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 +350,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 +365,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,14 +723,57 @@ 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_state, 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")
or ("WAIT_PAYMENT" if projection_seed.get("business_state") == "AWAITING_PAYMENT"
else "WAIT_CUSTOMER")),
"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"]))
if active_flow_mode == "authoritative":
canonical_items = _collapse_authoritative_material_processes(canonical_items)
partition = partition_canonical_items(canonical_items, display_limit=display_limit)
actionable_items = partition["actionable_items"]
all_waiting_items = partition["waiting_items"]
@@ -708,7 +802,7 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"backlog_total": partition["backlog_total"],
"not_current_total": partition["not_current_total"],
}
return {
result = {
"counts": cleaned_counts,
"recent_outbox": [dict(r) for r in recent_outbox],
"recent_documents": [dict(r) for r in recent_documents],
@@ -722,6 +816,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

@@ -15,6 +15,7 @@ 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,
@@ -51,6 +52,11 @@ def _load_v2_comparison_rows(opportunity_ids: list[str]) -> dict[str, dict[str,
"""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,
@@ -95,8 +101,6 @@ def compare_v1_v2_decisions(
def _observe_flow_v2(v1_decisions: Dict[str, Dict[str, Any]]) -> None:
mode = _flow_v2_mode()
if mode == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
if mode != "compare" or not v1_decisions:
return
rows = _load_v2_comparison_rows(list(v1_decisions))
@@ -228,7 +232,7 @@ def get_opportunity_next_action(
"""
if _flow_v2_mode() == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
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:
@@ -368,16 +372,28 @@ def _bulk_operation_snapshots(
def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Dict[str, Any]]:
"""Return the same decisions as the single-item API with a fixed query count."""
if _flow_v2_mode() == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip()))
if not ids:
return {}
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,
@@ -448,5 +464,88 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
company_profile="blif",
)
decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict()
_observe_flow_v2(decisions)
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 t.id::text, t.opportunity_id::text, t.action_code, t.status, t.due_at,
t.source_system, t.conversation_id, t.resolved_at, t.resolution_code,
t.superseded_by_task_id::text, t.created_at, t.metadata,
(SELECT max(m.created_at) FROM messages m
WHERE m.source_system = 'chatwoot' AND m.conversation_id = t.conversation_id
AND m.direction = 'inbound') AS latest_public_inbound,
(SELECT max(m.created_at) FROM messages m
WHERE m.source_system = 'chatwoot' AND m.conversation_id = t.conversation_id
AND m.direction = 'outbound') AS latest_public_outbound
FROM tasks t WHERE t.opportunity_id::text IN :opportunity_ids AND t.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()
winning_ids = {str(ref.get("id") or "") for ref in value.get("obligation_source_refs") or []}
obligation_audit = []
for task in tasks_by_id.get(oid, []):
task_id = str(task.get("id") or "")
if task_id in winning_ids:
classification = "VALID_ACTIVE_OBLIGATION"
elif task.get("superseded_by_task_id"):
classification = "SUPERSEDED_OBLIGATION"
elif task.get("resolved_at") or task.get("resolution_code"):
classification = "SATISFIED_OBLIGATION"
elif ((task_code := str(task.get("action_code") or "").upper()).startswith("SEND_")
and task_code != "SEND_INFO") or task_code == "CREATE_JASMIN_QUOTE":
classification = "PREMATURE_OR_STALE_OBLIGATION"
else:
classification = "STALE_LEGACY_OBLIGATION"
obligation_audit.append({
"id": task_id, "action_code": task.get("action_code"), "status": task.get("status"),
"source_system": task.get("source_system"), "classification": classification,
})
value["obligation_audit_refs"] = obligation_audit
if projection is None:
value["opportunity_id"] = oid
value["target_url"] = f"/opportunities/{oid}"
else:
value["material_process_key"] = projection.get("material_process_key")
value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}"
decisions[oid] = value
return decisions

View File

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

View File

@@ -99,15 +99,29 @@ def _load(
) -> dict[str, Any]:
with engine.connect() as conn:
conn = conn.execution_options(isolation_level="AUTOCOMMIT")
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if identity[0] != expected_database or (expected_user and identity[1] != expected_user):
raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
if require_read_only and identity[2] != "on":
raise RuntimeError(f"read-only simulation requires transaction_read_only=on: {identity!r}")
conn.execute(text("BEGIN READ ONLY"))
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,
@@ -157,7 +171,8 @@ def _load(
WHERE status IN ('open','needs_review','conflict') ORDER BY created_at
""")).mappings()]
finally:
conn.execute(text("ROLLBACK"))
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),

View File

@@ -0,0 +1,300 @@
#!/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", "CREATE_PROFORMA", "do_now"),
"Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", None, "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
"ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "PREPARE_ORDER", "do_now"),
"CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "PREPARE_ORDER", "do_now"),
"X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, None, "not_current"),
"X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, None, "suppressed"),
"RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS", "review"),
"RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", "REVIEW_REQUIRED", None, "suppressed"),
}
DUPLICATE_CANONICAL = {
"1816a06e-9a69-4a9b-9279-1263156892d3": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
"434124fb-ac19-4d78-909a-55761d7e8daa": "fd221608-e007-4043-a23d-07e0c119a345",
}
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 _material_identity(item: dict[str, Any]) -> str:
return str(item.get("material_process_key") or item.get("canonical_opportunity_id")
or item.get("opportunity_id") or item.get("process_key") or _key(item))
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,
projection_state_counts: dict[str, int] | None = None,
) -> dict[str, Any]:
v1_current = {_material_identity(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT}
v2_current = {_material_identity(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 = {"material_process_identity": key, "work_item_key": _key(old), "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 = []
for key, item in v2_current.items():
if key in v1_current:
continue
decision = item.get("decision") if isinstance(item.get("decision"), dict) else {}
winning = list(decision.get("obligation_source_refs") or [])
reason_code = decision.get("reason_code")
obligation_kind = ("VALID_ACTIVE_OBLIGATION" if str(reason_code).startswith(("ACTIVE_", "DUE_", "EXPLICIT_"))
else "FACTUAL_V2_BUSINESS_ACTION")
added.append({
"material_process_identity": key, "work_item_key": _key(item),
"opportunity_id": item.get("opportunity_id"),
"canonical_opportunity_id": item.get("canonical_opportunity_id"),
"material_process_key": item.get("material_process_key"),
"business_state": decision.get("business_state"),
"business_next_action": decision.get("business_next_action"),
"effective_action": decision.get("effective_action") or item.get("action_code"),
"queue": item.get("operational_queue"), "reason_code": reason_code,
"winning_obligation_task_id": (winning[0].get("id") if winning else None),
"task_status": (winning[0].get("status") if winning else None),
"source_system": (winning[0].get("source_system") if winning else None),
"authoritative_basis": obligation_kind,
"why_absent_from_v1": "NO_V1_CURRENT_CARD_FOR_MATERIAL_PROCESS",
"source_refs": item.get("source_refs") or [],
"non_winning_obligations": [ref for ref in decision.get("obligation_audit_refs") or []
if ref.get("classification") != "VALID_ACTIVE_OBLIGATION"],
})
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_business, 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")
actual_business = decision.get("business_next_action")
duplicate_case = expected_queue == "suppressed"
expected_canonical = DUPLICATE_CANONICAL.get(oid)
canonical_ok = not duplicate_case or decision.get("canonical_opportunity_id") == expected_canonical
passed = ((duplicate_case or actual_business == expected_business)
and actual_action == expected_action and actual_queue == expected_queue and canonical_ok)
named_cases[name] = {"opportunity_id": oid, "expected_business_next_action": expected_business,
"actual_business_next_action": actual_business, "expected_action": expected_action,
"expected_queue": expected_queue, "actual_action": actual_action,
"actual_queue": actual_queue, "passed": passed}
if duplicate_case:
named_cases[name].update({"expected_canonical_opportunity_id": expected_canonical,
"actual_canonical_opportunity_id": decision.get("canonical_opportunity_id")})
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)
invalid_obligation_cards = []
for item in v2_current.values():
decision = item.get("decision") if isinstance(item.get("decision"), dict) else {}
refs = decision.get("obligation_source_refs") or []
if refs and any(str(ref.get("status") or "").lower() != "pending" for ref in refs):
invalid_obligation_cards.append({"opportunity_id": item.get("opportunity_id"),
"effective_action": item.get("action_code"), "refs": refs})
grouped_added = {code: [row for row in added if row["effective_action"] == code]
for code in ("SEND_INFO", "SEND_QUOTE", "SEND_PROFORMA", "REVIEW_REQUIRED")}
state_counts = projection_state_counts or {}
waiting_projection_count = sum(state_counts.get(state, 0) for state in ("AWAITING_CUSTOMER", "AWAITING_PAYMENT"))
rendered_waiting = _metrics(v2)["waiting"]
represented_waiting_states = sum(
1 for item in _all(v2)
if (item.get("decision") or {}).get("business_state") in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}
)
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,
"added_card_diagnostics": added, "grouped_added_diagnostics": grouped_added,
"waiting_state_diagnostics": {
"projection_count": waiting_projection_count, "represented_processes": represented_waiting_states,
"rendered_waiting_cards": rendered_waiting,
"promoted_by_due_obligation": max(0, represented_waiting_states - rendered_waiting),
"rendering_policy": "Waiting states carry no invented internal action; they render in waiting only when selected as Operations processes.",
"lost": max(0, waiting_projection_count - represented_waiting_states),
},
"invalid_obligation_current_cards": invalid_obligation_cards,
}
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)
invalid_obligations = len(report.get("invalid_obligation_current_cards") or [])
return 2 if unsafe or missing or duplicates or failed_named or invalid_obligations 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()
state_rows = conn.exec_driver_sql("""
SELECT business_state, count(*) AS count
FROM opportunity_flow_state_v2 GROUP BY business_state
""").mappings().all()
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"])),
projection_state_counts={str(row["business_state"]): int(row["count"]) for row in state_rows},
)
_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,146 @@
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_terminal_state_ignores_leftover_pending_obligation():
got = decide_authoritative_operation(projection("COMPLETED"),
obligations=[task("REVIEW_RECONSTRUCTED_PROCESS")], now=NOW)
assert (got.effective_action, got.queue, got.reason_code) == (None, "not_current", "TERMINAL_BUSINESS_STATE")
def test_specific_reconstructed_review_overrides_generic_business_review():
got = decide_authoritative_operation(projection("REVIEW_REQUIRED", "REVIEW_REQUIRED"),
obligations=[task("REVIEW_RECONSTRUCTED_PROCESS")], now=NOW)
assert (got.business_next_action, got.effective_action, got.queue) == (
"REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS", "review")
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")
future_item = canonicalize_operations([{
"source": "flow_v2", "id": "opp", "opportunity_id": "opp",
"status": "projected", "action_code": "WAIT_CUSTOMER", "created_at": NOW,
}], {"opp": future.to_dict()})["items"]
assert len(future_item) == 1 and future_item[0]["operational_queue"] == "waiting"
assert decide_authoritative_operation(projection(), now=NOW).queue == "waiting"
def test_unanswered_send_info_remains_valid_on_awaiting_customer():
inbound = NOW - timedelta(hours=2)
obligation = task("SEND_INFO", latest_public_inbound=inbound, latest_public_outbound=None)
got = decide_authoritative_operation(projection(), obligations=[obligation], now=NOW)
assert (got.effective_action, got.queue, got.reason_code) == (
"SEND_INFO", "do_now", "ACTIVE_RESPONSE_OBLIGATION")
def test_send_info_satisfied_by_later_outbound_is_not_current():
inbound = NOW - timedelta(hours=2)
obligation = task("SEND_INFO", latest_public_inbound=inbound,
latest_public_outbound=inbound + timedelta(minutes=10))
got = decide_authoritative_operation(projection(), obligations=[obligation], now=NOW)
assert (got.effective_action, got.queue) == (None, "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("INQUIRY", "SEND_INFO"), obligations=[task("SEND_INFO")], now=NOW).effective_action == "SEND_INFO"
assert decide_authoritative_operation(projection(), obligations=[task("SEND_INFO")], now=NOW).effective_action is None
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_xmat_duplicate_suppresses_diagnostic_review_business_action():
got = decide_authoritative_operation({
**projection("REVIEW_REQUIRED", "REVIEW_REQUIRED"),
"opportunity_id": "1816a06e-9a69-4a9b-9279-1263156892d3",
"canonical_opportunity_id": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
"is_duplicate_representation": True,
}, now=NOW).to_dict()
assert got["business_next_action"] == "REVIEW_REQUIRED"
assert got["effective_action"] is None and got["suppress_current_card"] is True
assert got["canonical_opportunity_id"] == "dc89a466-db24-401b-bfe9-d47644b2d0c8"
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,168 @@
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
def all(self): return self.row if isinstance(self.row, list) else [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})
if "GROUP BY business_state" in sql:
return Result([])
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, business_action, action, queue) in sim.NAMED.items():
result[oid] = {"business_next_action": business_action, "effective_action": action,
"operational_queue": "not_current" if queue == "suppressed" else queue,
"suppress_current_card": queue == "suppressed",
"canonical_opportunity_id": sim.DUPLICATE_CANONICAL.get(oid, oid)}
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
def test_action_changes_correlate_by_material_process_not_action_key():
v1_item = {"work_item_key": "opportunity:o:action:OLD", "opportunity_id": "o",
"action_code": "OLD", "operational_queue": "do_now"}
v2_item = {"work_item_key": "opportunity:o:action:NEW", "opportunity_id": "o",
"canonical_opportunity_id": "o", "action_code": "NEW", "operational_queue": "do_now"}
report = sim.build_report(identity={}, configured_mode="compare",
v1={"work_items": [v1_item]}, v2={"work_items": [v2_item]},
named_decisions=named_decisions(), missing_projections=0)
assert len(report["changed_actions"]) == 1
assert report["removed_cards"] == [] and report["added_cards"] == []
def test_authoritative_material_collapse_selects_one_card_and_retains_evidence():
from app.operations_service import _collapse_authoritative_material_processes
rows = [
{"process_key": "conversation:chatwoot:1863", "work_item_key": "x:SEND_INFO",
"action_code": "SEND_INFO", "operational_queue": "do_now", "source_refs": [{"source": "task", "id": "1"}]},
{"process_key": "conversation:chatwoot:1863", "work_item_key": "x:REVIEW_MANUALLY",
"action_code": "REVIEW_MANUALLY", "operational_queue": "review", "source_refs": [{"source": "communication", "id": "2"}]},
]
result = _collapse_authoritative_material_processes(rows)
assert len(result) == 1
assert result[0]["action_code"] == "REVIEW_MANUALLY"
assert {ref["id"] for ref in result[0]["source_refs"]} == {"1", "2"}

View File

@@ -76,10 +76,9 @@ def test_duplicate_representation_stays_suppressed():
assert suppressed.effective_operational_queue == "not_current"
def test_authoritative_mode_is_fail_closed(monkeypatch):
def test_authoritative_empty_bulk_is_safe(monkeypatch):
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "authoritative")
with pytest.raises(RuntimeError, match="disabled"):
service.get_opportunity_next_actions([])
assert service.get_opportunity_next_actions([]) == {}
def test_shadow_mode_has_no_decision_side_effect(monkeypatch):

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