4 Commits

11 changed files with 1078 additions and 31 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

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

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