3 Commits

13 changed files with 980 additions and 39 deletions

View File

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

View File

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

View File

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

View File

@@ -157,6 +157,23 @@ def apply_operational_eligibility(items: Iterable[Mapping[str, Any]], *, now: da
projected = [] projected = []
for source in items: for source in items:
item = dict(source) 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) eligibility = evaluate_operational_eligibility(item, now=now)
item["eligibility"] = eligibility.to_dict() item["eligibility"] = eligibility.to_dict()
item["operational_queue"] = eligibility.queue item["operational_queue"] = eligibility.queue

View File

@@ -8,6 +8,7 @@ from __future__ import annotations
import os import os
import re import re
import subprocess import subprocess
from contextlib import nullcontext
from typing import Any, Dict, List, Optional from typing import Any, Dict, List, Optional
from sqlalchemy import text from sqlalchemy import text
@@ -26,7 +27,7 @@ from app.canonical_operations import (
opportunity_ids_from_work_seeds, opportunity_ids_from_work_seeds,
partition_canonical_items, 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 from app.document_reconciliation_service import active_document_link_exclusion_sql
@@ -104,7 +105,7 @@ def _norm_identity(value: Any) -> str:
def _identity_tokens(value: Any) -> set[str]: def _identity_tokens(value: Any) -> set[str]:
return { result = {
token token
for token in _norm_identity(value).split() for token in _norm_identity(value).split()
if len(token) >= 3 and token not in _IDENTITY_STOPWORDS if len(token) >= 3 and token not in _IDENTITY_STOPWORDS
@@ -302,7 +303,9 @@ def _attach_operation_urls(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
return _sanitize_operation_identities(cleaned) 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. """Build the /operations work queue summary.
/operations is intentionally not a mini-dashboard. It returns a compact /operations is intentionally not a mini-dashboard. It returns a compact
@@ -315,7 +318,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
# after canonicalization. This prevents duplicate rows from consuming the # after canonicalization. This prevents duplicate rows from consuming the
# display limit while still avoiding an unbounded history scan. # display limit while still avoiding an unbounded history scan.
candidate_limit = max(200, min(display_limit * 20, 1000)) 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(""" counts = conn.execute(text("""
SELECT SELECT
(SELECT COUNT(*) FROM opportunities WHERE status = 'open')::int AS open_opportunities, (SELECT COUNT(*) FROM opportunities WHERE status = 'open')::int AS open_opportunities,
@@ -672,12 +676,50 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
conversation_timeline = {str(row["conversation_id"]): dict(row) for row in timeline_rows} conversation_timeline = {str(row["conversation_id"]): dict(row) for row in timeline_rows}
seed_rows = [dict(r) for r in work_seed_rows] seed_rows = [dict(r) for r in work_seed_rows]
# In authoritative mode the factual projection, not a legacy task/stage,
# selects material processes that have current business work. Synthetic
# rows are read-model seeds only and are never persisted.
if active_flow_mode == "authoritative":
seeded = {str(row.get("opportunity_id") or "") for row in seed_rows}
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
v2_seeds = conn.execute(text("""
SELECT p.opportunity_id::text, p.business_next_action, p.derived_at,
o.title, o.value_amount, o.currency, o.lifecycle_state,
c.name AS customer_name
FROM opportunity_flow_state_v2 p
JOIN opportunities o ON o.id = p.opportunity_id
LEFT JOIN customers c ON c.id = o.local_customer_id
WHERE p.is_duplicate_representation = false
AND (p.business_next_action IS NOT NULL
OR p.business_state IN ('AWAITING_CUSTOMER', 'AWAITING_PAYMENT'))
""")).mappings().all()
for projection_seed in v2_seeds:
oid = str(projection_seed["opportunity_id"])
if oid in seeded:
continue
seed_rows.append({
"source": "flow_v2", "id": oid, "opportunity_id": oid,
"created_at": projection_seed.get("derived_at"), "due_at": None,
"priority": "normal", "queue": "vendas",
"action_code": projection_seed.get("business_next_action"), "status": "projected",
"title": projection_seed.get("title") or "Oportunidade",
"opportunity_title": projection_seed.get("title") or "",
"customer_name": projection_seed.get("customer_name") or "",
"fiscal_customer_name": projection_seed.get("customer_name") or "",
"item_metadata": {}, "opportunity_metadata": {},
"opportunity_lifecycle_state": projection_seed.get("lifecycle_state") or "",
"opportunity_value_amount": projection_seed.get("value_amount") or 0,
"opportunity_currency": projection_seed.get("currency") or "EUR",
"href": f"/opportunities/{oid}", "action_label": "Abrir",
})
for row in seed_rows: for row in seed_rows:
timeline = conversation_timeline.get(str(row.get("conversation_id") or ""), {}) 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")}) row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")})
normalized_rows = _normalise_work_item_intent(seed_rows) normalized_rows = _normalise_work_item_intent(seed_rows)
opportunity_ids = opportunity_ids_from_work_seeds(normalized_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]) projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
canonical_items = _attach_operation_urls(list(projection["items"])) canonical_items = _attach_operation_urls(list(projection["items"]))
partition = partition_canonical_items(canonical_items, display_limit=display_limit) partition = partition_canonical_items(canonical_items, display_limit=display_limit)
@@ -722,6 +764,11 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
**diagnostics, **diagnostics,
"projection_metrics": 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]: 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.db import engine
from app.config import settings 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. # Backward-compatible static anchors from v1.5.59: quotation_doc, invoice_doc, confirmar pagamento antes de emitir fatura, Pagamento confirmado com base em, Criar/enviar fatura.
from app.domain.opportunity_flow import ( from app.domain.opportunity_flow import (
OpportunityEvidence, OpportunityEvidence,
@@ -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.""" """Read the persisted projection only; comparison must never derive writes."""
if not opportunity_ids: if not opportunity_ids:
return {} 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: with engine.connect() as conn:
rows = conn.execute(text(""" rows = conn.execute(text("""
SELECT opportunity_id::text, material_process_key, 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: def _observe_flow_v2(v1_decisions: Dict[str, Dict[str, Any]]) -> None:
mode = _flow_v2_mode() 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: if mode != "compare" or not v1_decisions:
return return
rows = _load_v2_comparison_rows(list(v1_decisions)) rows = _load_v2_comparison_rows(list(v1_decisions))
@@ -228,7 +232,7 @@ def get_opportunity_next_action(
""" """
if _flow_v2_mode() == "authoritative": 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) evidence = (_build_db_evidence(opportunity_id, preloaded=preloaded)
if preloaded is not None else _build_db_evidence(opportunity_id)) if preloaded is not None else _build_db_evidence(opportunity_id))
if evidence is None: if evidence is None:
@@ -368,16 +372,28 @@ def _bulk_operation_snapshots(
def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Dict[str, Any]]: def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Dict[str, Any]]:
"""Return the same decisions as the single-item API with a fixed query count.""" """Return the same decisions as the single-item API with a fixed query count."""
if _flow_v2_mode() == "authoritative":
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip())) ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip()))
if not ids: if not ids:
return {} return {}
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 # Read path only: schema creation belongs to startup/migrations. In
# particular, GET /operations must never perform DDL while calculating its # particular, GET /operations must never perform DDL while calculating its
# canonical projection. # 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, """ opportunities = _bulk_rows(conn, """
SELECT id::text, stage, status, title, SELECT id::text, stage, status, title,
local_customer_id::text AS fiscal_customer_id, local_customer_id::text AS fiscal_customer_id,
@@ -448,5 +464,60 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
company_profile="blif", company_profile="blif",
) )
decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict() decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict()
if flow_mode == "compare":
_observe_flow_v2(decisions) _observe_flow_v2(decisions)
return decisions return decisions
def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any = None) -> Dict[str, Dict[str, Any]]:
"""Read-only set-oriented adapter input loader; never derives or persists V2."""
from contextlib import nullcontext
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
projections = _bulk_rows(conn, """
SELECT opportunity_id::text, material_process_key,
canonical_opportunity_id::text, is_duplicate_representation,
business_state, business_next_action, diagnostic_status,
confidence, reason_code, reason_text
FROM opportunity_flow_state_v2 WHERE opportunity_id::text IN :opportunity_ids
""", ids)
tasks = _bulk_rows(conn, """
SELECT id::text, opportunity_id::text, action_code, status, due_at,
resolved_at, superseded_by_task_id::text, created_at, metadata
FROM tasks WHERE opportunity_id::text IN :opportunity_ids AND status = 'pending'
""", ids)
customers = _bulk_rows(conn, """
SELECT o.id::text AS opportunity_id,
(c.tax_id IS NOT NULL AND c.tax_id <> '' AND c.email IS NOT NULL AND c.email <> ''
AND c.street_name IS NOT NULL AND c.street_name <> ''
AND c.postal_zone IS NOT NULL AND c.postal_zone <> ''
AND c.city_name IS NOT NULL AND c.city_name <> '') AS fiscal_complete
FROM opportunities o LEFT JOIN customers c ON c.id = o.local_customer_id
WHERE o.id::text IN :opportunity_ids
""", ids)
reconciliation = _bulk_rows(conn, """
SELECT opportunity_id::text
FROM reconciliation_items
WHERE opportunity_id::text IN :opportunity_ids
AND status IN ('needs_review','conflict')
AND COALESCE((payload->>'reconciliation_blocks_current_action')::boolean, false)
""", ids)
projection_by_id = {str(row["opportunity_id"]): row for row in projections}
tasks_by_id = _group_by_opportunity(tasks)
fiscal_by_id = {str(row["opportunity_id"]): bool(row.get("fiscal_complete")) for row in customers}
reconciliation_ids = {str(row["opportunity_id"]) for row in reconciliation}
decisions = {}
for oid in ids:
projection = projection_by_id.get(oid)
decision = decide_authoritative_operation(
projection, obligations=tasks_by_id.get(oid, []),
fiscal_complete=fiscal_by_id.get(oid, False),
reconciliation_blocking=oid in reconciliation_ids,
)
value = decision.to_dict()
if projection is None:
value["opportunity_id"] = oid
value["target_url"] = f"/opportunities/{oid}"
else:
value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}"
decisions[oid] = value
return decisions

View File

@@ -99,6 +99,15 @@ def _load(
) -> dict[str, Any]: ) -> dict[str, Any]:
with engine.connect() as conn: with engine.connect() as conn:
conn = conn.execution_options(isolation_level="AUTOCOMMIT") conn = conn.execution_options(isolation_level="AUTOCOMMIT")
transaction_started = False
try:
# A fresh connection commonly reports transaction_read_only=off.
# Establish the protected transaction on the same connection used
# for every factual read before validating a strict audit.
if require_read_only:
conn.execute(text("BEGIN READ ONLY"))
transaction_started = True
identity = conn.execute(text( identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')" "SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one() )).one()
@@ -106,8 +115,13 @@ def _load(
raise RuntimeError(f"refusing unexpected database identity: {identity!r}") raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
if require_read_only and identity[2] != "on": if require_read_only and identity[2] != "on":
raise RuntimeError(f"read-only simulation requires transaction_read_only=on: {identity!r}") 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")) conn.execute(text("BEGIN READ ONLY"))
try: transaction_started = True
opportunities = [dict(row) for row in conn.execute(text(""" opportunities = [dict(row) for row in conn.execute(text("""
SELECT o.*, o.id::text AS id, o.local_customer_id::text, SELECT o.*, o.id::text AS id, o.local_customer_id::text,
c.name AS linked_customer_name, c.tax_id, c.email AS fiscal_email, c.name AS linked_customer_name, c.tax_id, c.email AS fiscal_email,
@@ -157,6 +171,7 @@ def _load(
WHERE status IN ('open','needs_review','conflict') ORDER BY created_at WHERE status IN ('open','needs_review','conflict') ORDER BY created_at
""")).mappings()] """)).mappings()]
finally: finally:
if transaction_started:
conn.execute(text("ROLLBACK")) conn.execute(text("ROLLBACK"))
return { return {
"identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]}, "identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},

View File

@@ -0,0 +1,226 @@
#!/usr/bin/env python3
"""Simulate authoritative BLIF Flow v2 without changing application mode."""
from __future__ import annotations
import argparse
import json
import os
from pathlib import Path
import sys
from typing import Any
JSON_PATH = Path("/tmp/blif_flow_v2_authoritative_simulation.json")
TEXT_PATH = Path("/tmp/blif_flow_v2_authoritative_operations.txt")
CURRENT = {"do_now", "review", "blocked", "exception"}
CLASSIFICATIONS = {
"SATISFIED", "SUPERSEDED", "DUPLICATE_SUPPRESSED", "PREMATURE_REMOVED",
"REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE",
}
NAMED = {
"Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "do_now"),
"Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
"ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "do_now"),
"CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "do_now"),
"X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, "not_current"),
"X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, "suppressed"),
"RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "review"),
"RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", None, "suppressed"),
}
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
parser = argparse.ArgumentParser(
description="Read-only V1 versus simulated authoritative BLIF Flow v2 Operations audit."
)
parser.add_argument(
"--production-readonly-simulation", action="store_true",
help="Required opt-in when current_database() is clientflow; requires compare mode and a read-only transaction.",
)
return parser.parse_args(argv)
def _load_runtime():
"""Application imports are intentionally delayed until after argparse."""
root = Path(__file__).resolve().parents[1]
if str(root) not in sys.path:
sys.path.insert(0, str(root))
from app.db import engine
from app.operations_service import get_operations_summary
from app.opportunity_next_action_service import get_opportunity_next_actions_for_mode
return engine, get_operations_summary, get_opportunity_next_actions_for_mode
def _all(summary: dict[str, Any]) -> list[dict[str, Any]]:
if "simulation_all_items" in summary:
return list(summary.get("simulation_all_items") or [])
return sum((list(summary.get(key) or []) for key in
("work_items", "waiting_items", "backlog_items", "not_current_items")), [])
def _metrics(summary: dict[str, Any]) -> dict[str, int]:
queues = [str(item.get("operational_queue") or "") for item in _all(summary)]
return {"current": sum(q in CURRENT for q in queues), "do_now": queues.count("do_now"),
"review": queues.count("review"), "waiting": queues.count("waiting"),
"backlog": queues.count("backlog"), "blocked": queues.count("blocked"),
"not_current": queues.count("not_current")}
def _key(item: dict[str, Any]) -> str:
return str(item.get("work_item_key") or item.get("process_key") or item.get("opportunity_id") or item.get("id"))
def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[str, str]:
if new is not None:
return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item."
decision = old.get("decision") if isinstance(old.get("decision"), dict) else {}
code = str(decision.get("reason_code") or old.get("eligibility_reason_code") or "").upper()
if "DUPLICATE" in code:
return "DUPLICATE_SUPPRESSED", "The local opportunity is a duplicate material representation."
if "SATISF" in code or "ANSWERED" in code:
return "SATISFIED", "Later factual evidence satisfies the old obligation."
if "SUPERSE" in code or "TERMINAL" in code:
return "SUPERSEDED", "A later factual state supersedes the legacy card."
if "PREMATURE" in code:
return "PREMATURE_REMOVED", "The legacy action is premature for the factual state."
if old.get("source") in {"task", "opportunity"}:
return "LEGACY_ONLY", "The card is supported only by legacy operational representation."
return "UNSAFE_FALSE_NEGATIVE", "No safe factual explanation was found for removing this current card."
def build_report(
*, identity: dict[str, Any], configured_mode: str, v1: dict[str, Any], v2: dict[str, Any],
named_decisions: dict[str, dict[str, Any]], missing_projections: int,
) -> dict[str, Any]:
v1_current = {_key(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT}
v2_current = {_key(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT}
removed, changed = [], []
for key, old in v1_current.items():
new = v2_current.get(key)
if new is not None and new.get("action_code") == old.get("action_code"):
continue
classification, reason = _classification(old, new)
row = {"work_item_key": key, "opportunity_id": old.get("opportunity_id"),
"v1_action": old.get("action_code"), "simulated_v2_action": (new or {}).get("action_code"),
"classification": classification, "reason": reason,
"source_refs": old.get("source_refs") or []}
(changed if new is not None else removed).append(row)
added = [{"work_item_key": key, "opportunity_id": item.get("opportunity_id"),
"action": item.get("action_code"), "source_refs": item.get("source_refs") or []}
for key, item in v2_current.items() if key not in v1_current]
material_counts: dict[str, int] = {}
for item in v2_current.values():
material = str(item.get("canonical_opportunity_id") or item.get("opportunity_id") or item.get("process_key"))
material_counts[material] = material_counts.get(material, 0) + 1
duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1]
named_cases = {}
for name, (oid, expected_action, expected_queue) in NAMED.items():
decision = named_decisions.get(oid) or {}
suppressed = decision.get("suppress_current_card") is True
actual_queue = "suppressed" if suppressed else decision.get("operational_queue")
actual_action = decision.get("effective_action")
passed = actual_action == expected_action and actual_queue == expected_queue
named_cases[name] = {"opportunity_id": oid, "expected_action": expected_action,
"expected_queue": expected_queue, "actual_action": actual_action,
"actual_queue": actual_queue, "passed": passed}
fiscal = []
for item in v2_current.values():
decision = item.get("decision") if isinstance(item.get("decision"), dict) else {}
if item.get("action_code") == "VALIDATE_FISCAL_CUSTOMER":
fiscal.append({"opportunity_id": item.get("opportunity_id"),
"business_action": decision.get("business_next_action"),
"reason_code": decision.get("reason_code"),
"reason": decision.get("description")})
unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed)
return {
"identity": identity, "configured_mode": configured_mode,
"v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2),
"semantic_changes": len(removed) + len(changed) + len(added),
"removed_cards": removed, "added_cards": added, "changed_actions": changed,
"duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates),
"missing_projections": missing_projections, "unsafe_false_negatives": unsafe,
"named_cases": named_cases, "fiscal_validation_cases": fiscal,
}
def _write_reports(report: dict[str, Any]) -> None:
JSON_PATH.write_text(json.dumps(report, indent=2, ensure_ascii=False, default=str) + "\n")
lines = ["BLIF Flow v2 authoritative Operations simulation", "",
f"Identity: {report['identity']}", f"Configured mode: {report['configured_mode']}",
f"V1: {report['v1_metrics']}", f"Simulated authoritative V2: {report['authoritative_metrics']}",
f"Removed: {len(report['removed_cards'])}", f"Added: {len(report['added_cards'])}",
f"Changed actions: {len(report['changed_actions'])}",
f"Duplicate current cards: {report['duplicate_current_cards']}",
f"Missing projections: {report['missing_projections']}",
f"UNSAFE_FALSE_NEGATIVE: {report['unsafe_false_negatives']}", "", "Named cases:"]
lines.extend(f"- {name}: {case}" for name, case in report["named_cases"].items())
lines.extend(["", f"Fiscal validation cases: {len(report['fiscal_validation_cases'])}"])
TEXT_PATH.write_text("\n".join(lines) + "\n")
def exit_code(report: dict[str, Any]) -> int:
failed_named = any(not case.get("passed") for case in report.get("named_cases", {}).values())
unsafe = int(report.get("unsafe_false_negatives") or 0)
missing = int(report.get("missing_projections") or 0)
duplicates = int(report.get("duplicate_current_cards") or 0)
return 2 if unsafe or missing or duplicates or failed_named else 0
def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, Any]:
engine, get_operations, get_decisions = runtime_loader()
configured_mode = os.environ.get("BLIF_FLOW_V2_MODE")
conn = engine.connect()
try:
# One explicit transaction and one factual snapshot for identity, V1,
# V2, projection completeness, and named-case decisions.
conn.exec_driver_sql("BEGIN READ ONLY")
identity = dict(conn.exec_driver_sql("""
SELECT current_database() AS database, current_user AS user,
current_setting('transaction_read_only') AS transaction_read_only
""").mappings().one())
production = identity.get("database") == "clientflow"
if production and not args.production_readonly_simulation:
raise RuntimeError("production simulation requires --production-readonly-simulation")
if args.production_readonly_simulation:
if identity.get("database") != "clientflow":
raise RuntimeError("production simulation requires current_database() = clientflow")
if identity.get("transaction_read_only") != "on":
raise RuntimeError("production simulation requires transaction_read_only = on")
if configured_mode is None:
raise RuntimeError("production simulation requires explicit BLIF_FLOW_V2_MODE=compare")
if configured_mode.strip().lower() != "compare":
raise RuntimeError("production simulation requires BLIF_FLOW_V2_MODE=compare")
v1 = get_operations(limit=200, flow_mode="off", connection=conn)
v2 = get_operations(limit=200, flow_mode="authoritative", connection=conn)
projection_counts = conn.exec_driver_sql("""
SELECT (SELECT count(*) FROM opportunities) AS opportunities,
(SELECT count(*) FROM opportunity_flow_state_v2) AS projections
""").mappings().one()
named_ids = [value[0] for value in NAMED.values()]
named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn)
if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode:
raise RuntimeError("configured BLIF_FLOW_V2_MODE changed during simulation")
report = build_report(
identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2,
named_decisions=named_decisions,
missing_projections=max(0, int(projection_counts["opportunities"]) - int(projection_counts["projections"])),
)
_write_reports(report)
return report
finally:
conn.exec_driver_sql("ROLLBACK")
conn.close()
def main(argv: list[str] | None = None, *, runtime_loader=_load_runtime) -> int:
args = parse_args(argv) # --help exits before runtime_loader/application imports.
report = run(args, runtime_loader=runtime_loader)
print(json.dumps(report, indent=2, ensure_ascii=False, default=str))
return exit_code(report)
if __name__ == "__main__":
raise SystemExit(main())

View File

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

View File

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

View File

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

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 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: class _Result:
def __init__(self, rows): def __init__(self, rows):
self._rows = rows self._rows = rows