From fc092a96b693dedbf21e8ec87fa4d54f82ca7720 Mon Sep 17 00:00:00 2001 From: plx Date: Sun, 16 Aug 2026 01:15:26 +0000 Subject: [PATCH] feat: add BLIF Flow v2 authoritative operational adapter --- app/admin_ui/labels.py | 8 + app/authoritative_operational_adapter.py | 182 ++++++++++++++++++ app/canonical_operations.py | 38 ++-- app/operational_eligibility.py | 17 ++ app/operations_service.py | 36 ++++ app/opportunity_next_action_service.py | 68 ++++++- .../simulate_blif_flow_v2_authoritative.py | 95 +++++++++ .../test_authoritative_operational_adapter.py | 99 ++++++++++ tests/test_blif_flow_v2_cutover.py | 5 +- tests/test_opportunity_next_action_bulk.py | 6 + 10 files changed, 531 insertions(+), 23 deletions(-) create mode 100644 app/authoritative_operational_adapter.py create mode 100644 scripts/simulate_blif_flow_v2_authoritative.py create mode 100644 tests/test_authoritative_operational_adapter.py diff --git a/app/admin_ui/labels.py b/app/admin_ui/labels.py index 1b208dd..c634c48 100644 --- a/app/admin_ui/labels.py +++ b/app/admin_ui/labels.py @@ -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", diff --git a/app/authoritative_operational_adapter.py b/app/authoritative_operational_adapter.py new file mode 100644 index 0000000..a0b2b18 --- /dev/null +++ b/app/authoritative_operational_adapter.py @@ -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) diff --git a/app/canonical_operations.py b/app/canonical_operations.py index 7bfa540..c09f655 100644 --- a/app/canonical_operations.py +++ b/app/canonical_operations.py @@ -321,13 +321,14 @@ def _opportunity_item( rows: list[Mapping[str, Any]], decision: Mapping[str, Any], ) -> tuple[dict[str, Any] | None, list[dict[str, Any]]]: + authoritative = decision.get("authoritative_v2") is True scheduled_call = next(( row for row in rows if row.get("source") == "task" and _s(row.get("status")).lower() == "pending" and _canonical_code(row.get("action_code")) == "CALL_CUSTOMER" ), None) - if scheduled_call: + if scheduled_call and not authoritative: decision = { **dict(decision), "action_code": "CALL_CUSTOMER", @@ -347,26 +348,29 @@ def _opportunity_item( evidence_rows.append(row) if action_code in NO_WORK_ACTION_CODES: - if not _has_explicit_unresolved_review(evidence_rows): + if decision.get("suppress_current_card"): return None, independent_exceptions - # The central decision normally wins, but an explicit unresolved review - # marker must not be hidden by NO_ACTION/NOT_FOUND. - action_code = "REVIEW" - decision = { - **dict(decision), - "action_code": action_code, - "label": "Rever evidência pendente", - "description": "A decisão central indica que não há ação, mas existe evidência explicitamente marcada para revisão.", - "reason": "Existe revisão humana não resolvida apesar da decisão central sem ação.", - "priority": "normal", - "target_url": f"/opportunities/{opportunity_id}", - } + keep_authoritative_noncurrent = authoritative and _s(decision.get("operational_queue")) in {"waiting", "backlog"} + if not keep_authoritative_noncurrent: + if not _has_explicit_unresolved_review(evidence_rows): + return None, independent_exceptions + # The central decision normally wins, but an explicit unresolved + # review marker must not be hidden by NO_ACTION/NOT_FOUND. + action_code = "REVIEW" + decision = { + **dict(decision), "action_code": action_code, + "label": "Rever evidência pendente", + "description": "A decisão central indica que não há ação, mas existe evidência explicitamente marcada para revisão.", + "reason": "Existe revisão humana não resolvida apesar da decisão central sem ação.", + "priority": "normal", "target_url": f"/opportunities/{opportunity_id}", + } matching_task = scheduled_call or _matching_task(evidence_rows, action_code) primary = matching_task or next((row for row in evidence_rows if row.get("source") == "task"), None) primary = primary or (evidence_rows[0] if evidence_rows else {}) process_key = f"opportunity:{opportunity_id}" waiting = _is_waiting(decision) + decision_queue = _s(decision.get("operational_queue")) if authoritative else "" opportunity_row = next((row for row in evidence_rows if _s(row.get("fiscal_customer_name"))), None) opportunity_row = opportunity_row or next((row for row in evidence_rows if _s(row.get("opportunity_title"))), None) opportunity_row = opportunity_row or primary @@ -383,7 +387,7 @@ def _opportunity_item( "detail": _s(decision.get("description") or decision.get("reason")), "why_human_required": _s(decision.get("reason") or decision.get("description")), "priority": _s(decision.get("priority")) or _best_priority(evidence_rows), - "operational_queue": "waiting" if waiting else "do_now", + "operational_queue": decision_queue or ("waiting" if waiting else "do_now"), "queue": _s((matching_task or primary).get("queue")) or _s(primary.get("queue")) or "rever", "status": "waiting" if waiting else "pending", "href": _s((matching_task or {}).get("href")) or _s(decision.get("target_url")) or f"/opportunities/{opportunity_id}", @@ -400,6 +404,10 @@ def _opportunity_item( ], "decision_version": decision.get("decision_version"), "decision": dict(decision), + "business_state": decision.get("business_state"), + "business_next_action": decision.get("business_next_action"), + "effective_action": decision.get("effective_action"), + "canonical_opportunity_id": decision.get("canonical_opportunity_id") or opportunity_id, }) return item, independent_exceptions diff --git a/app/operational_eligibility.py b/app/operational_eligibility.py index 5798fb7..57a84e3 100644 --- a/app/operational_eligibility.py +++ b/app/operational_eligibility.py @@ -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 diff --git a/app/operations_service.py b/app/operations_service.py index 8434826..6a109a6 100644 --- a/app/operations_service.py +++ b/app/operations_service.py @@ -672,6 +672,42 @@ 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 str(settings.blif_flow_v2_mode or "off").strip().lower() == "authoritative": + seeded = {str(row.get("opportunity_id") or "") for row in seed_rows} + with engine.connect() as conn: + v2_seeds = conn.execute(text(""" + SELECT p.opportunity_id::text, p.business_next_action, p.derived_at, + o.title, o.value_amount, o.currency, o.lifecycle_state, + c.name AS customer_name + FROM opportunity_flow_state_v2 p + JOIN opportunities o ON o.id = p.opportunity_id + LEFT JOIN customers c ON c.id = o.local_customer_id + WHERE p.is_duplicate_representation = false + AND (p.business_next_action IS NOT NULL + OR p.business_state IN ('AWAITING_CUSTOMER', 'AWAITING_PAYMENT')) + """)).mappings().all() + for projection_seed in v2_seeds: + oid = str(projection_seed["opportunity_id"]) + if oid in seeded: + continue + seed_rows.append({ + "source": "flow_v2", "id": oid, "opportunity_id": oid, + "created_at": projection_seed.get("derived_at"), "due_at": None, + "priority": "normal", "queue": "vendas", + "action_code": projection_seed.get("business_next_action"), "status": "projected", + "title": projection_seed.get("title") or "Oportunidade", + "opportunity_title": projection_seed.get("title") or "", + "customer_name": projection_seed.get("customer_name") or "", + "fiscal_customer_name": projection_seed.get("customer_name") or "", + "item_metadata": {}, "opportunity_metadata": {}, + "opportunity_lifecycle_state": projection_seed.get("lifecycle_state") or "", + "opportunity_value_amount": projection_seed.get("value_amount") or 0, + "opportunity_currency": projection_seed.get("currency") or "EUR", + "href": f"/opportunities/{oid}", "action_label": "Abrir", + }) for row in seed_rows: timeline = conversation_timeline.get(str(row.get("conversation_id") or ""), {}) row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")}) diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index 80fa3cf..c041c55 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -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,12 +372,13 @@ 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 {} + if _flow_v2_mode() == "authoritative": + return _get_authoritative_opportunity_decisions(ids) + # Read path only: schema creation belongs to startup/migrations. In # particular, GET /operations must never perform DDL while calculating its # canonical projection. @@ -450,3 +455,56 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict() _observe_flow_v2(decisions) return decisions + + +def _get_authoritative_opportunity_decisions(ids: list[str]) -> Dict[str, Dict[str, Any]]: + """Read-only set-oriented adapter input loader; never derives or persists V2.""" + with 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 diff --git a/scripts/simulate_blif_flow_v2_authoritative.py b/scripts/simulate_blif_flow_v2_authoritative.py new file mode 100644 index 0000000..79102b5 --- /dev/null +++ b/scripts/simulate_blif_flow_v2_authoritative.py @@ -0,0 +1,95 @@ +#!/usr/bin/env python3 +"""Read-only audit of the exact Operations boundary in V1 and authoritative V2.""" +from __future__ import annotations + +import json +from pathlib import Path +import sys + +ROOT = Path(__file__).resolve().parents[1] +if str(ROOT) not in sys.path: + sys.path.insert(0, str(ROOT)) + +from app.config import settings +from app.operations_service import get_operations_summary + +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"} + + +def _all(summary): + return list(summary.get("work_items") or []) + list(summary.get("waiting_items") or []) + list(summary.get("backlog_items") or []) + list(summary.get("not_current_items") or []) + + +def _metrics(summary): + items = _all(summary) + queues = [str(item.get("operational_queue") or "") for item in items] + return { + "current": sum(queue in CURRENT for queue in queues), + "do_now": queues.count("do_now"), "review": queues.count("review"), + "waiting": queues.count("waiting"), "backlog": queues.count("backlog"), + "not_current": queues.count("not_current"), + } + + +def _identity(item): + return str(item.get("canonical_opportunity_id") or item.get("opportunity_id") or item.get("process_key") or item.get("work_item_key")) + + +def _classification(old, new): + if new is not None: + return "REPLACED_BY_CORRECT_ACTION" + decision = old.get("decision") or {} + code = str(decision.get("reason_code") or old.get("eligibility_reason_code") or "").upper() + if "DUPLICATE" in code: + return "DUPLICATE_SUPPRESSED" + if "SATISF" in code or "ANSWERED" in code: + return "SATISFIED" + if "SUPERSE" in code or "TERMINAL" in code: + return "SUPERSEDED" + if "PREMATURE" in code: + return "PREMATURE_REMOVED" + if old.get("source") in {"task", "opportunity"}: + return "LEGACY_ONLY" + return "UNSAFE_FALSE_NEGATIVE" + + +def simulate(): + original = settings.blif_flow_v2_mode + try: + settings.blif_flow_v2_mode = "off" + v1 = get_operations_summary(limit=200) + settings.blif_flow_v2_mode = "authoritative" + v2 = get_operations_summary(limit=200) + finally: + settings.blif_flow_v2_mode = original + v1_current = {_identity(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT} + v2_current = {_identity(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT} + removed = [] + for key, old in v1_current.items(): + new = v2_current.get(key) + if new is None or new.get("action_code") != old.get("action_code"): + removed.append({"identity": key, "v1_action": old.get("action_code"), + "v2_action": (new or {}).get("action_code"), + "classification": _classification(old, new)}) + added = [{"identity": key, "action": item.get("action_code")} + for key, item in v2_current.items() + if key not in v1_current or v1_current[key].get("action_code") != item.get("action_code")] + result = {"v1": _metrics(v1), "authoritative_v2": _metrics(v2), + "cards_removed": removed, "cards_added_or_replaced": added, + "unsafe_false_negatives": sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed)} + JSON_PATH.write_text(json.dumps(result, indent=2, ensure_ascii=False, default=str) + "\n") + lines = ["BLIF Flow v2 authoritative Operations simulation", "", f"V1: {result['v1']}", + f"Authoritative V2: {result['authoritative_v2']}", + f"Cards removed: {len(removed)}", f"Cards added/replaced: {len(added)}", + f"UNSAFE_FALSE_NEGATIVE: {result['unsafe_false_negatives']}"] + TEXT_PATH.write_text("\n".join(lines) + "\n") + return result + + +if __name__ == "__main__": + audit = simulate() + print(json.dumps(audit, indent=2, ensure_ascii=False, default=str)) + if audit["unsafe_false_negatives"]: + raise SystemExit(2) diff --git a/tests/test_authoritative_operational_adapter.py b/tests/test_authoritative_operational_adapter.py new file mode 100644 index 0000000..310080c --- /dev/null +++ b/tests/test_authoritative_operational_adapter.py @@ -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) diff --git a/tests/test_blif_flow_v2_cutover.py b/tests/test_blif_flow_v2_cutover.py index 1f5752a..7a1430c 100644 --- a/tests/test_blif_flow_v2_cutover.py +++ b/tests/test_blif_flow_v2_cutover.py @@ -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): diff --git a/tests/test_opportunity_next_action_bulk.py b/tests/test_opportunity_next_action_bulk.py index babf0bc..269caf2 100644 --- a/tests/test_opportunity_next_action_bulk.py +++ b/tests/test_opportunity_next_action_bulk.py @@ -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