Compare commits
2 Commits
fix/blif-f
...
ac385342f5
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac385342f5 | ||
|
|
fc092a96b6 |
@@ -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",
|
||||
|
||||
182
app/authoritative_operational_adapter.py
Normal file
182
app/authoritative_operational_adapter.py
Normal file
@@ -0,0 +1,182 @@
|
||||
"""Authoritative BLIF Flow v2 -> operator decision adapter.
|
||||
|
||||
The persisted projection owns factual business state. Pending tasks are only
|
||||
operational obligations and may override that state when this policy proves
|
||||
that they are current. This module is deliberately pure and performs no I/O.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import asdict, dataclass
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Iterable, Mapping
|
||||
|
||||
from app.admin_ui.labels import action_label
|
||||
|
||||
|
||||
FORMAL_ACTIONS = {"CREATE_PROFORMA", "CREATE_INVOICE"}
|
||||
REVIEW_ACTIONS = {"REVIEW", "REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS"}
|
||||
PAYMENT_FOLLOWUPS = {"FOLLOW_UP_PAYMENT", "FOLLOW_UP_PROFORMA"}
|
||||
VALID_OVERRIDES = REVIEW_ACTIONS | {"SUPPORT", "SEND_INFO", "CALL_CUSTOMER", "FOLLOW_UP_CUSTOMER_REVIEW"} | PAYMENT_FOLLOWUPS
|
||||
TERMINAL_STATES = {"COMPLETED", "LOST", "NO_INTEREST"}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class AuthoritativeOperationalDecision:
|
||||
opportunity_id: str
|
||||
canonical_opportunity_id: str
|
||||
business_state: str
|
||||
business_next_action: str | None
|
||||
effective_action: str | None
|
||||
queue: str
|
||||
eligible: bool
|
||||
reason_code: str
|
||||
reason_text: str
|
||||
blocking_action_code: str | None = None
|
||||
obligation_source_refs: tuple[dict[str, Any], ...] = ()
|
||||
confidence: str = "high"
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
result = asdict(self)
|
||||
result["obligation_source_refs"] = list(self.obligation_source_refs)
|
||||
# Shared legacy presentation contract consumed by Operations and detail.
|
||||
result.update({
|
||||
"action_code": self.effective_action or "NO_ACTION",
|
||||
"label": action_label(self.effective_action, "Sem ação"),
|
||||
"description": self.reason_text,
|
||||
"priority": "alta" if self.queue in {"review", "blocked", "exception"} else "normal",
|
||||
"can_execute": self.eligible and self.effective_action is not None,
|
||||
"reason_if_blocked": self.reason_code if not self.eligible else None,
|
||||
"operational_queue": self.queue,
|
||||
"authoritative_v2": True,
|
||||
"suppress_current_card": self.queue == "not_current",
|
||||
})
|
||||
return result
|
||||
|
||||
|
||||
def _code(value: Any) -> str:
|
||||
return str(value or "").strip().upper()
|
||||
|
||||
|
||||
def _dt(value: Any) -> datetime | None:
|
||||
if not value:
|
||||
return None
|
||||
if isinstance(value, datetime):
|
||||
parsed = value
|
||||
else:
|
||||
try:
|
||||
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
|
||||
except ValueError:
|
||||
return None
|
||||
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def _active_obligations(obligations: Iterable[Mapping[str, Any]]) -> list[Mapping[str, Any]]:
|
||||
return [row for row in obligations
|
||||
if str(row.get("status") or "").lower() == "pending"
|
||||
and not row.get("resolved_at") and not row.get("superseded_by_task_id")]
|
||||
|
||||
|
||||
def _ref(row: Mapping[str, Any]) -> dict[str, Any]:
|
||||
return {"source": "task", "id": str(row.get("id") or ""), "status": "pending",
|
||||
"action_code": _code(row.get("action_code"))}
|
||||
|
||||
|
||||
def decide_authoritative_operation(
|
||||
projection: Mapping[str, Any] | None,
|
||||
*, obligations: Iterable[Mapping[str, Any]] = (), now: datetime | None = None,
|
||||
fiscal_complete: bool = True, reconciliation_blocking: bool = False,
|
||||
hard_blocker: str | None = None,
|
||||
) -> AuthoritativeOperationalDecision:
|
||||
"""Return the sole operator decision, failing closed without a projection."""
|
||||
now = now or datetime.now(timezone.utc)
|
||||
if projection is None:
|
||||
return AuthoritativeOperationalDecision(
|
||||
"", "", "MISSING_PROJECTION", None, "REVIEW_REQUIRED", "review", True,
|
||||
"MISSING_V2_PROJECTION", "A projeção Flow v2 está em falta; é necessária revisão, sem recorrer ao V1.",
|
||||
confidence="low",
|
||||
)
|
||||
oid = str(projection.get("opportunity_id") or "")
|
||||
canonical = str(projection.get("canonical_opportunity_id") or oid)
|
||||
state = _code(projection.get("business_state"))
|
||||
business_action = _code(projection.get("business_next_action")) or None
|
||||
confidence = str(projection.get("confidence") or "low")
|
||||
if projection.get("is_duplicate_representation"):
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, None, "not_current", False,
|
||||
"DUPLICATE_SUPPRESSED", f"Representação duplicada do processo material canónico {canonical}.", confidence="high",
|
||||
)
|
||||
if hard_blocker:
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, hard_blocker, "blocked", True,
|
||||
"HARD_FACTUAL_BLOCKER", "Um bloqueio factual ou de sistema impede a ação atual.", hard_blocker, confidence=confidence,
|
||||
)
|
||||
|
||||
active = _active_obligations(obligations)
|
||||
by_code: dict[str, list[Mapping[str, Any]]] = {}
|
||||
for row in active:
|
||||
by_code.setdefault(_code(row.get("action_code")), []).append(row)
|
||||
|
||||
# A formal-document prerequisite blocks only a transition which needs it.
|
||||
if business_action in FORMAL_ACTIONS and not fiscal_complete:
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, "VALIDATE_FISCAL_CUSTOMER", "blocked", True,
|
||||
"FISCAL_IDENTITY_REQUIRED", "Validar os dados fiscais antes de criar o documento oficial.",
|
||||
"VALIDATE_FISCAL_CUSTOMER", confidence=confidence,
|
||||
)
|
||||
if business_action in FORMAL_ACTIONS and reconciliation_blocking:
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, "RECONCILE_DOCUMENTS", "blocked", True,
|
||||
"DOCUMENT_RECONCILIATION_REQUIRED", "Confirmar a ligação do documento formal atual.",
|
||||
"RECONCILE_DOCUMENTS", confidence=confidence,
|
||||
)
|
||||
|
||||
def choose(codes: Iterable[str], queue: str, reason: str):
|
||||
for code in codes:
|
||||
rows = by_code.get(code, [])
|
||||
if rows:
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, code, queue, True, reason,
|
||||
"Existe uma obrigação operacional pendente e válida.",
|
||||
obligation_source_refs=tuple(_ref(row) for row in rows), confidence=confidence,
|
||||
)
|
||||
return None
|
||||
|
||||
picked = choose(("REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW", "REVIEW_RECONSTRUCTED_PROCESS"), "review", "ACTIVE_REVIEW_OBLIGATION")
|
||||
picked = picked or choose(("SUPPORT",), "do_now", "ACTIVE_SUPPORT_OBLIGATION")
|
||||
picked = picked or choose(("SEND_INFO",), "do_now", "ACTIVE_RESPONSE_OBLIGATION")
|
||||
picked = picked or choose(("CALL_CUSTOMER",), "do_now", "EXPLICIT_CALL_OBLIGATION")
|
||||
if picked:
|
||||
return picked
|
||||
|
||||
waiting_state = state in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}
|
||||
followup_order = ("FOLLOW_UP_CUSTOMER_REVIEW",) if state == "AWAITING_CUSTOMER" else tuple(PAYMENT_FOLLOWUPS)
|
||||
if waiting_state:
|
||||
for code in followup_order:
|
||||
rows = by_code.get(code, [])
|
||||
if not rows:
|
||||
continue
|
||||
due_rows = [row for row in rows if _dt(row.get("due_at")) is None or _dt(row.get("due_at")) <= now]
|
||||
if due_rows:
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, code, "do_now", True, "DUE_FOLLOW_UP",
|
||||
"O follow-up ativo chegou à data e continua por satisfazer.",
|
||||
obligation_source_refs=tuple(_ref(row) for row in due_rows), confidence=confidence,
|
||||
)
|
||||
return AuthoritativeOperationalDecision(
|
||||
oid, canonical, state, business_action, None, "waiting", False, "FOLLOW_UP_NOT_DUE",
|
||||
"O follow-up ativo ainda não chegou à data.",
|
||||
obligation_source_refs=tuple(_ref(row) for row in rows), confidence=confidence,
|
||||
)
|
||||
|
||||
if state in TERMINAL_STATES:
|
||||
return AuthoritativeOperationalDecision(oid, canonical, state, business_action, None, "not_current", False,
|
||||
"TERMINAL_BUSINESS_STATE", "O processo factual está concluído.", confidence=confidence)
|
||||
if business_action:
|
||||
queue = "review" if business_action in REVIEW_ACTIONS else "do_now"
|
||||
return AuthoritativeOperationalDecision(oid, canonical, state, business_action, business_action, queue, True,
|
||||
"BUSINESS_NEXT_ACTION", str(projection.get("reason_text") or "A ação decorre do estado factual Flow v2."), confidence=confidence)
|
||||
if waiting_state:
|
||||
return AuthoritativeOperationalDecision(oid, canonical, state, None, None, "waiting", False,
|
||||
"WAITING_EXTERNAL_EVENT", "O processo aguarda um evento externo.", confidence=confidence)
|
||||
return AuthoritativeOperationalDecision(oid, canonical, state, None, None, "backlog", False,
|
||||
"NO_CURRENT_INTERNAL_ACTION", "Não existe ação interna atual.", confidence=confidence)
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -8,6 +8,7 @@ from __future__ import annotations
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
from contextlib import nullcontext
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from sqlalchemy import text
|
||||
@@ -26,7 +27,7 @@ from app.canonical_operations import (
|
||||
opportunity_ids_from_work_seeds,
|
||||
partition_canonical_items,
|
||||
)
|
||||
from app.opportunity_next_action_service import get_opportunity_next_actions
|
||||
from app.opportunity_next_action_service import get_opportunity_next_actions, get_opportunity_next_actions_for_mode
|
||||
from app.document_reconciliation_service import active_document_link_exclusion_sql
|
||||
|
||||
|
||||
@@ -104,7 +105,7 @@ def _norm_identity(value: Any) -> str:
|
||||
|
||||
|
||||
def _identity_tokens(value: Any) -> set[str]:
|
||||
return {
|
||||
result = {
|
||||
token
|
||||
for token in _norm_identity(value).split()
|
||||
if len(token) >= 3 and token not in _IDENTITY_STOPWORDS
|
||||
@@ -302,7 +303,9 @@ def _attach_operation_urls(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
||||
return _sanitize_operation_identities(cleaned)
|
||||
|
||||
|
||||
def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
def get_operations_summary(
|
||||
limit: int = 24, *, flow_mode: str | None = None, connection: Any = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Build the /operations work queue summary.
|
||||
|
||||
/operations is intentionally not a mini-dashboard. It returns a compact
|
||||
@@ -315,7 +318,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
# after canonicalization. This prevents duplicate rows from consuming the
|
||||
# display limit while still avoiding an unbounded history scan.
|
||||
candidate_limit = max(200, min(display_limit * 20, 1000))
|
||||
with engine.begin() as conn:
|
||||
active_flow_mode = str(flow_mode if flow_mode is not None else settings.blif_flow_v2_mode or "off").strip().lower()
|
||||
with (nullcontext(connection) if connection is not None else engine.begin()) as conn:
|
||||
counts = conn.execute(text("""
|
||||
SELECT
|
||||
(SELECT COUNT(*) FROM opportunities WHERE status = 'open')::int AS open_opportunities,
|
||||
@@ -672,12 +676,50 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
conversation_timeline = {str(row["conversation_id"]): dict(row) for row in timeline_rows}
|
||||
|
||||
seed_rows = [dict(r) for r in work_seed_rows]
|
||||
# In authoritative mode the factual projection, not a legacy task/stage,
|
||||
# selects material processes that have current business work. Synthetic
|
||||
# rows are read-model seeds only and are never persisted.
|
||||
if active_flow_mode == "authoritative":
|
||||
seeded = {str(row.get("opportunity_id") or "") for row in seed_rows}
|
||||
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
|
||||
v2_seeds = conn.execute(text("""
|
||||
SELECT p.opportunity_id::text, p.business_next_action, p.derived_at,
|
||||
o.title, o.value_amount, o.currency, o.lifecycle_state,
|
||||
c.name AS customer_name
|
||||
FROM opportunity_flow_state_v2 p
|
||||
JOIN opportunities o ON o.id = p.opportunity_id
|
||||
LEFT JOIN customers c ON c.id = o.local_customer_id
|
||||
WHERE p.is_duplicate_representation = false
|
||||
AND (p.business_next_action IS NOT NULL
|
||||
OR p.business_state IN ('AWAITING_CUSTOMER', 'AWAITING_PAYMENT'))
|
||||
""")).mappings().all()
|
||||
for projection_seed in v2_seeds:
|
||||
oid = str(projection_seed["opportunity_id"])
|
||||
if oid in seeded:
|
||||
continue
|
||||
seed_rows.append({
|
||||
"source": "flow_v2", "id": oid, "opportunity_id": oid,
|
||||
"created_at": projection_seed.get("derived_at"), "due_at": None,
|
||||
"priority": "normal", "queue": "vendas",
|
||||
"action_code": projection_seed.get("business_next_action"), "status": "projected",
|
||||
"title": projection_seed.get("title") or "Oportunidade",
|
||||
"opportunity_title": projection_seed.get("title") or "",
|
||||
"customer_name": projection_seed.get("customer_name") or "",
|
||||
"fiscal_customer_name": projection_seed.get("customer_name") or "",
|
||||
"item_metadata": {}, "opportunity_metadata": {},
|
||||
"opportunity_lifecycle_state": projection_seed.get("lifecycle_state") or "",
|
||||
"opportunity_value_amount": projection_seed.get("value_amount") or 0,
|
||||
"opportunity_currency": projection_seed.get("currency") or "EUR",
|
||||
"href": f"/opportunities/{oid}", "action_label": "Abrir",
|
||||
})
|
||||
for row in seed_rows:
|
||||
timeline = conversation_timeline.get(str(row.get("conversation_id") or ""), {})
|
||||
row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")})
|
||||
normalized_rows = _normalise_work_item_intent(seed_rows)
|
||||
opportunity_ids = opportunity_ids_from_work_seeds(normalized_rows)
|
||||
decisions = get_opportunity_next_actions(opportunity_ids) if opportunity_ids else {}
|
||||
decisions = (get_opportunity_next_actions_for_mode(
|
||||
opportunity_ids, flow_mode=active_flow_mode, connection=connection,
|
||||
) if opportunity_ids else {})
|
||||
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
|
||||
canonical_items = _attach_operation_urls(list(projection["items"]))
|
||||
partition = partition_canonical_items(canonical_items, display_limit=display_limit)
|
||||
@@ -722,6 +764,11 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
**diagnostics,
|
||||
"projection_metrics": diagnostics,
|
||||
}
|
||||
if flow_mode is not None:
|
||||
# Complete, untruncated read-model population for offline simulation;
|
||||
# normal API/GET callers do not receive this audit-only field.
|
||||
result["simulation_all_items"] = canonical_items
|
||||
return result
|
||||
|
||||
|
||||
def get_system_health_summary() -> Dict[str, Any]:
|
||||
|
||||
@@ -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,60 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
|
||||
company_profile="blif",
|
||||
)
|
||||
decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict()
|
||||
_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 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
|
||||
|
||||
226
scripts/simulate_blif_flow_v2_authoritative.py
Normal file
226
scripts/simulate_blif_flow_v2_authoritative.py
Normal file
@@ -0,0 +1,226 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Simulate authoritative BLIF Flow v2 without changing application mode."""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import sys
|
||||
from typing import Any
|
||||
|
||||
JSON_PATH = Path("/tmp/blif_flow_v2_authoritative_simulation.json")
|
||||
TEXT_PATH = Path("/tmp/blif_flow_v2_authoritative_operations.txt")
|
||||
CURRENT = {"do_now", "review", "blocked", "exception"}
|
||||
CLASSIFICATIONS = {
|
||||
"SATISFIED", "SUPERSEDED", "DUPLICATE_SUPPRESSED", "PREMATURE_REMOVED",
|
||||
"REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE",
|
||||
}
|
||||
NAMED = {
|
||||
"Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "do_now"),
|
||||
"Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
|
||||
"ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "do_now"),
|
||||
"CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "do_now"),
|
||||
"X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, "not_current"),
|
||||
"X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, "suppressed"),
|
||||
"RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "review"),
|
||||
"RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", None, "suppressed"),
|
||||
}
|
||||
|
||||
|
||||
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Read-only V1 versus simulated authoritative BLIF Flow v2 Operations audit."
|
||||
)
|
||||
parser.add_argument(
|
||||
"--production-readonly-simulation", action="store_true",
|
||||
help="Required opt-in when current_database() is clientflow; requires compare mode and a read-only transaction.",
|
||||
)
|
||||
return parser.parse_args(argv)
|
||||
|
||||
|
||||
def _load_runtime():
|
||||
"""Application imports are intentionally delayed until after argparse."""
|
||||
root = Path(__file__).resolve().parents[1]
|
||||
if str(root) not in sys.path:
|
||||
sys.path.insert(0, str(root))
|
||||
from app.db import engine
|
||||
from app.operations_service import get_operations_summary
|
||||
from app.opportunity_next_action_service import get_opportunity_next_actions_for_mode
|
||||
return engine, get_operations_summary, get_opportunity_next_actions_for_mode
|
||||
|
||||
|
||||
def _all(summary: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
if "simulation_all_items" in summary:
|
||||
return list(summary.get("simulation_all_items") or [])
|
||||
return sum((list(summary.get(key) or []) for key in
|
||||
("work_items", "waiting_items", "backlog_items", "not_current_items")), [])
|
||||
|
||||
|
||||
def _metrics(summary: dict[str, Any]) -> dict[str, int]:
|
||||
queues = [str(item.get("operational_queue") or "") for item in _all(summary)]
|
||||
return {"current": sum(q in CURRENT for q in queues), "do_now": queues.count("do_now"),
|
||||
"review": queues.count("review"), "waiting": queues.count("waiting"),
|
||||
"backlog": queues.count("backlog"), "blocked": queues.count("blocked"),
|
||||
"not_current": queues.count("not_current")}
|
||||
|
||||
|
||||
def _key(item: dict[str, Any]) -> str:
|
||||
return str(item.get("work_item_key") or item.get("process_key") or item.get("opportunity_id") or item.get("id"))
|
||||
|
||||
|
||||
def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[str, str]:
|
||||
if new is not None:
|
||||
return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item."
|
||||
decision = old.get("decision") if isinstance(old.get("decision"), dict) else {}
|
||||
code = str(decision.get("reason_code") or old.get("eligibility_reason_code") or "").upper()
|
||||
if "DUPLICATE" in code:
|
||||
return "DUPLICATE_SUPPRESSED", "The local opportunity is a duplicate material representation."
|
||||
if "SATISF" in code or "ANSWERED" in code:
|
||||
return "SATISFIED", "Later factual evidence satisfies the old obligation."
|
||||
if "SUPERSE" in code or "TERMINAL" in code:
|
||||
return "SUPERSEDED", "A later factual state supersedes the legacy card."
|
||||
if "PREMATURE" in code:
|
||||
return "PREMATURE_REMOVED", "The legacy action is premature for the factual state."
|
||||
if old.get("source") in {"task", "opportunity"}:
|
||||
return "LEGACY_ONLY", "The card is supported only by legacy operational representation."
|
||||
return "UNSAFE_FALSE_NEGATIVE", "No safe factual explanation was found for removing this current card."
|
||||
|
||||
|
||||
def build_report(
|
||||
*, identity: dict[str, Any], configured_mode: str, v1: dict[str, Any], v2: dict[str, Any],
|
||||
named_decisions: dict[str, dict[str, Any]], missing_projections: int,
|
||||
) -> dict[str, Any]:
|
||||
v1_current = {_key(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT}
|
||||
v2_current = {_key(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT}
|
||||
removed, changed = [], []
|
||||
for key, old in v1_current.items():
|
||||
new = v2_current.get(key)
|
||||
if new is not None and new.get("action_code") == old.get("action_code"):
|
||||
continue
|
||||
classification, reason = _classification(old, new)
|
||||
row = {"work_item_key": key, "opportunity_id": old.get("opportunity_id"),
|
||||
"v1_action": old.get("action_code"), "simulated_v2_action": (new or {}).get("action_code"),
|
||||
"classification": classification, "reason": reason,
|
||||
"source_refs": old.get("source_refs") or []}
|
||||
(changed if new is not None else removed).append(row)
|
||||
added = [{"work_item_key": key, "opportunity_id": item.get("opportunity_id"),
|
||||
"action": item.get("action_code"), "source_refs": item.get("source_refs") or []}
|
||||
for key, item in v2_current.items() if key not in v1_current]
|
||||
|
||||
material_counts: dict[str, int] = {}
|
||||
for item in v2_current.values():
|
||||
material = str(item.get("canonical_opportunity_id") or item.get("opportunity_id") or item.get("process_key"))
|
||||
material_counts[material] = material_counts.get(material, 0) + 1
|
||||
duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1]
|
||||
|
||||
named_cases = {}
|
||||
for name, (oid, expected_action, expected_queue) in NAMED.items():
|
||||
decision = named_decisions.get(oid) or {}
|
||||
suppressed = decision.get("suppress_current_card") is True
|
||||
actual_queue = "suppressed" if suppressed else decision.get("operational_queue")
|
||||
actual_action = decision.get("effective_action")
|
||||
passed = actual_action == expected_action and actual_queue == expected_queue
|
||||
named_cases[name] = {"opportunity_id": oid, "expected_action": expected_action,
|
||||
"expected_queue": expected_queue, "actual_action": actual_action,
|
||||
"actual_queue": actual_queue, "passed": passed}
|
||||
|
||||
fiscal = []
|
||||
for item in v2_current.values():
|
||||
decision = item.get("decision") if isinstance(item.get("decision"), dict) else {}
|
||||
if item.get("action_code") == "VALIDATE_FISCAL_CUSTOMER":
|
||||
fiscal.append({"opportunity_id": item.get("opportunity_id"),
|
||||
"business_action": decision.get("business_next_action"),
|
||||
"reason_code": decision.get("reason_code"),
|
||||
"reason": decision.get("description")})
|
||||
unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed)
|
||||
return {
|
||||
"identity": identity, "configured_mode": configured_mode,
|
||||
"v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2),
|
||||
"semantic_changes": len(removed) + len(changed) + len(added),
|
||||
"removed_cards": removed, "added_cards": added, "changed_actions": changed,
|
||||
"duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates),
|
||||
"missing_projections": missing_projections, "unsafe_false_negatives": unsafe,
|
||||
"named_cases": named_cases, "fiscal_validation_cases": fiscal,
|
||||
}
|
||||
|
||||
|
||||
def _write_reports(report: dict[str, Any]) -> None:
|
||||
JSON_PATH.write_text(json.dumps(report, indent=2, ensure_ascii=False, default=str) + "\n")
|
||||
lines = ["BLIF Flow v2 authoritative Operations simulation", "",
|
||||
f"Identity: {report['identity']}", f"Configured mode: {report['configured_mode']}",
|
||||
f"V1: {report['v1_metrics']}", f"Simulated authoritative V2: {report['authoritative_metrics']}",
|
||||
f"Removed: {len(report['removed_cards'])}", f"Added: {len(report['added_cards'])}",
|
||||
f"Changed actions: {len(report['changed_actions'])}",
|
||||
f"Duplicate current cards: {report['duplicate_current_cards']}",
|
||||
f"Missing projections: {report['missing_projections']}",
|
||||
f"UNSAFE_FALSE_NEGATIVE: {report['unsafe_false_negatives']}", "", "Named cases:"]
|
||||
lines.extend(f"- {name}: {case}" for name, case in report["named_cases"].items())
|
||||
lines.extend(["", f"Fiscal validation cases: {len(report['fiscal_validation_cases'])}"])
|
||||
TEXT_PATH.write_text("\n".join(lines) + "\n")
|
||||
|
||||
|
||||
def exit_code(report: dict[str, Any]) -> int:
|
||||
failed_named = any(not case.get("passed") for case in report.get("named_cases", {}).values())
|
||||
unsafe = int(report.get("unsafe_false_negatives") or 0)
|
||||
missing = int(report.get("missing_projections") or 0)
|
||||
duplicates = int(report.get("duplicate_current_cards") or 0)
|
||||
return 2 if unsafe or missing or duplicates or failed_named else 0
|
||||
|
||||
|
||||
def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, Any]:
|
||||
engine, get_operations, get_decisions = runtime_loader()
|
||||
configured_mode = os.environ.get("BLIF_FLOW_V2_MODE")
|
||||
conn = engine.connect()
|
||||
try:
|
||||
# One explicit transaction and one factual snapshot for identity, V1,
|
||||
# V2, projection completeness, and named-case decisions.
|
||||
conn.exec_driver_sql("BEGIN READ ONLY")
|
||||
identity = dict(conn.exec_driver_sql("""
|
||||
SELECT current_database() AS database, current_user AS user,
|
||||
current_setting('transaction_read_only') AS transaction_read_only
|
||||
""").mappings().one())
|
||||
production = identity.get("database") == "clientflow"
|
||||
if production and not args.production_readonly_simulation:
|
||||
raise RuntimeError("production simulation requires --production-readonly-simulation")
|
||||
if args.production_readonly_simulation:
|
||||
if identity.get("database") != "clientflow":
|
||||
raise RuntimeError("production simulation requires current_database() = clientflow")
|
||||
if identity.get("transaction_read_only") != "on":
|
||||
raise RuntimeError("production simulation requires transaction_read_only = on")
|
||||
if configured_mode is None:
|
||||
raise RuntimeError("production simulation requires explicit BLIF_FLOW_V2_MODE=compare")
|
||||
if configured_mode.strip().lower() != "compare":
|
||||
raise RuntimeError("production simulation requires BLIF_FLOW_V2_MODE=compare")
|
||||
|
||||
v1 = get_operations(limit=200, flow_mode="off", connection=conn)
|
||||
v2 = get_operations(limit=200, flow_mode="authoritative", connection=conn)
|
||||
projection_counts = conn.exec_driver_sql("""
|
||||
SELECT (SELECT count(*) FROM opportunities) AS opportunities,
|
||||
(SELECT count(*) FROM opportunity_flow_state_v2) AS projections
|
||||
""").mappings().one()
|
||||
named_ids = [value[0] for value in NAMED.values()]
|
||||
named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn)
|
||||
if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode:
|
||||
raise RuntimeError("configured BLIF_FLOW_V2_MODE changed during simulation")
|
||||
report = build_report(
|
||||
identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2,
|
||||
named_decisions=named_decisions,
|
||||
missing_projections=max(0, int(projection_counts["opportunities"]) - int(projection_counts["projections"])),
|
||||
)
|
||||
_write_reports(report)
|
||||
return report
|
||||
finally:
|
||||
conn.exec_driver_sql("ROLLBACK")
|
||||
conn.close()
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None, *, runtime_loader=_load_runtime) -> int:
|
||||
args = parse_args(argv) # --help exits before runtime_loader/application imports.
|
||||
report = run(args, runtime_loader=runtime_loader)
|
||||
print(json.dumps(report, indent=2, ensure_ascii=False, default=str))
|
||||
return exit_code(report)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
99
tests/test_authoritative_operational_adapter.py
Normal file
99
tests/test_authoritative_operational_adapter.py
Normal file
@@ -0,0 +1,99 @@
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from app.authoritative_operational_adapter import decide_authoritative_operation
|
||||
from app.canonical_operations import canonicalize_operations
|
||||
|
||||
|
||||
NOW = datetime(2026, 8, 16, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def projection(state="AWAITING_CUSTOMER", action=None, **extra):
|
||||
return {"opportunity_id": "opp", "canonical_opportunity_id": "opp",
|
||||
"business_state": state, "business_next_action": action,
|
||||
"confidence": "high", **extra}
|
||||
|
||||
|
||||
def task(code, due=NOW, **extra):
|
||||
return {"id": code, "status": "pending", "action_code": code, "due_at": due, **extra}
|
||||
|
||||
|
||||
def test_business_baselines_and_terminal():
|
||||
for state, action in [("PROFORMA_REQUIRED", "CREATE_PROFORMA"),
|
||||
("INVOICE_REQUIRED", "CREATE_INVOICE"),
|
||||
("ODOO_ORDER_REQUIRED", "PREPARE_ORDER"),
|
||||
("ODOO_ORDER_CREATED", "VALIDATE_ODOO_ORDER")]:
|
||||
got = decide_authoritative_operation(projection(state, action), now=NOW)
|
||||
assert (got.effective_action, got.queue) == (action, "do_now")
|
||||
assert decide_authoritative_operation(projection("COMPLETED"), now=NOW).queue == "not_current"
|
||||
|
||||
|
||||
def test_missing_projection_fails_closed_without_v1():
|
||||
got = decide_authoritative_operation(None, now=NOW)
|
||||
assert (got.effective_action, got.queue, got.reason_code) == ("REVIEW_REQUIRED", "review", "MISSING_V2_PROJECTION")
|
||||
|
||||
|
||||
def test_due_and_future_customer_followup_semantics():
|
||||
due = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW")], now=NOW)
|
||||
assert (due.effective_action, due.queue) == ("FOLLOW_UP_CUSTOMER_REVIEW", "do_now")
|
||||
future = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW", NOW + timedelta(days=1))], now=NOW)
|
||||
assert (future.effective_action, future.queue) == (None, "waiting")
|
||||
assert decide_authoritative_operation(projection(), now=NOW).queue == "waiting"
|
||||
|
||||
|
||||
def test_only_active_non_superseded_tasks_override_and_precedence():
|
||||
stale = task("SUPPORT", resolved_at=NOW)
|
||||
assert decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), obligations=[stale], now=NOW).effective_action == "CREATE_PROFORMA"
|
||||
obligations = [task("CALL_CUSTOMER"), task("SEND_INFO"), task("SUPPORT"), task("REVIEW_MANUALLY")]
|
||||
assert decide_authoritative_operation(projection(), obligations=obligations, now=NOW).effective_action == "REVIEW_MANUALLY"
|
||||
assert decide_authoritative_operation(projection(), obligations=[task("SUPPORT")], now=NOW).effective_action == "SUPPORT"
|
||||
assert decide_authoritative_operation(projection(), obligations=[task("SEND_INFO")], now=NOW).effective_action == "SEND_INFO"
|
||||
assert decide_authoritative_operation(projection(), obligations=[task("CALL_CUSTOMER")], now=NOW).effective_action == "CALL_CUSTOMER"
|
||||
|
||||
|
||||
def test_fiscal_and_reconciliation_only_block_formal_action():
|
||||
inquiry = decide_authoritative_operation(projection("INQUIRY", "SEND_INFO"), fiscal_complete=False, reconciliation_blocking=True, now=NOW)
|
||||
assert inquiry.effective_action == "SEND_INFO"
|
||||
formal = decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), fiscal_complete=False, now=NOW)
|
||||
assert (formal.effective_action, formal.blocking_action_code) == ("VALIDATE_FISCAL_CUSTOMER", "VALIDATE_FISCAL_CUSTOMER")
|
||||
|
||||
|
||||
def test_duplicate_is_suppressed_even_with_pending_legacy_task():
|
||||
decision = decide_authoritative_operation(projection("REVIEW_REQUIRED", "REVIEW_REQUIRED",
|
||||
is_duplicate_representation=True,
|
||||
canonical_opportunity_id="canonical"),
|
||||
obligations=[task("REVIEW_MANUALLY")], now=NOW).to_dict()
|
||||
source = {"source": "task", "id": "t", "status": "pending", "action_code": "REVIEW_MANUALLY",
|
||||
"opportunity_id": "opp", "created_at": NOW}
|
||||
assert canonicalize_operations([source], {"opp": decision})["items"] == []
|
||||
|
||||
|
||||
def test_authoritative_queue_and_action_survive_canonical_boundary():
|
||||
decision = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW")], now=NOW).to_dict()
|
||||
source = {"source": "task", "id": "old", "status": "pending", "action_code": "SEND_INVOICE",
|
||||
"opportunity_id": "opp", "created_at": NOW}
|
||||
item = canonicalize_operations([source], {"opp": decision})["items"][0]
|
||||
assert (item["action_code"], item["operational_queue"]) == ("FOLLOW_UP_CUSTOMER_REVIEW", "do_now")
|
||||
|
||||
|
||||
def test_named_case_acceptance_decisions():
|
||||
cases = {
|
||||
"5c33db95-fab8-477a-bddd-0b9cc8f91302": ("PROFORMA_REQUIRED", "CREATE_PROFORMA", [], "CREATE_PROFORMA", "do_now"),
|
||||
"fd79b9a1-07e6-4f61-95e8-09eab89c155e": ("AWAITING_CUSTOMER", None, [task("FOLLOW_UP_CUSTOMER_REVIEW")], "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
|
||||
"61f1c955-a372-4ea7-b9b0-b8528d74a141": ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER", [], "PREPARE_ORDER", "do_now"),
|
||||
"e3b23ac5-84db-4763-8a31-a684e873032c": ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER", [], "PREPARE_ORDER", "do_now"),
|
||||
"dc89a466-db24-401b-bfe9-d47644b2d0c8": ("COMPLETED", None, [], None, "not_current"),
|
||||
"fd221608-e007-4043-a23d-07e0c119a345": ("REVIEW_REQUIRED", "REVIEW_REQUIRED", [], "REVIEW_REQUIRED", "review"),
|
||||
}
|
||||
for oid, (state, action, obligations, effective, queue) in cases.items():
|
||||
got = decide_authoritative_operation({**projection(state, action), "opportunity_id": oid,
|
||||
"canonical_opportunity_id": oid},
|
||||
obligations=obligations, now=NOW)
|
||||
assert (got.effective_action, got.queue) == (effective, queue)
|
||||
for duplicate, canonical in [
|
||||
("1816a06e-9a69-4a9b-9279-1263156892d3", "dc89a466-db24-401b-bfe9-d47644b2d0c8"),
|
||||
("434124fb-ac19-4d78-909a-55761d7e8daa", "fd221608-e007-4043-a23d-07e0c119a345"),
|
||||
]:
|
||||
got = decide_authoritative_operation({**projection("REVIEW_REQUIRED", "REVIEW_REQUIRED"),
|
||||
"opportunity_id": duplicate, "canonical_opportunity_id": canonical,
|
||||
"is_duplicate_representation": True}, now=NOW)
|
||||
assert (got.effective_action, got.queue, got.canonical_opportunity_id) == (None, "not_current", canonical)
|
||||
138
tests/test_blif_flow_v2_authoritative_simulation.py
Normal file
138
tests/test_blif_flow_v2_authoritative_simulation.py
Normal file
@@ -0,0 +1,138 @@
|
||||
from argparse import Namespace
|
||||
import importlib.util
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
PATH = Path("scripts/simulate_blif_flow_v2_authoritative.py")
|
||||
SPEC = importlib.util.spec_from_file_location("authoritative_simulation", PATH)
|
||||
sim = importlib.util.module_from_spec(SPEC)
|
||||
SPEC.loader.exec_module(sim)
|
||||
|
||||
|
||||
class Result:
|
||||
def __init__(self, row): self.row = row
|
||||
def mappings(self): return self
|
||||
def one(self): return self.row
|
||||
|
||||
|
||||
class Connection:
|
||||
def __init__(self, database="clientflow", readonly="on"):
|
||||
self.database, self.readonly, self.sql = database, readonly, []
|
||||
def exec_driver_sql(self, sql):
|
||||
self.sql.append(sql.strip())
|
||||
if "current_database" in sql:
|
||||
return Result({"database": self.database, "user": "audit", "transaction_read_only": self.readonly})
|
||||
if "count(*) FROM opportunities" in sql:
|
||||
return Result({"opportunities": 328, "projections": 328})
|
||||
return Result({})
|
||||
def close(self): pass
|
||||
|
||||
|
||||
class Engine:
|
||||
def __init__(self, connection): self.connection = connection
|
||||
def connect(self): return self.connection
|
||||
|
||||
|
||||
def named_decisions(fail=False):
|
||||
result = {}
|
||||
for name, (oid, action, queue) in sim.NAMED.items():
|
||||
result[oid] = {"effective_action": action,
|
||||
"operational_queue": "not_current" if queue == "suppressed" else queue,
|
||||
"suppress_current_card": queue == "suppressed"}
|
||||
if fail:
|
||||
result[sim.NAMED["Instalbeira"][0]]["effective_action"] = "SEND_INVOICE"
|
||||
return result
|
||||
|
||||
|
||||
def runtime(connection, *, fail_named=False, modes=None):
|
||||
empty = {"work_items": [], "waiting_items": [], "backlog_items": [], "not_current_items": []}
|
||||
def operations(**kwargs):
|
||||
assert kwargs["connection"] is connection
|
||||
if modes is not None:
|
||||
modes.append(kwargs["flow_mode"])
|
||||
return empty
|
||||
def decisions(ids, **kwargs):
|
||||
assert kwargs == {"flow_mode": "authoritative", "connection": connection}
|
||||
return named_decisions(fail_named)
|
||||
return lambda: (Engine(connection), operations, decisions)
|
||||
|
||||
|
||||
def test_help_performs_zero_db_work():
|
||||
called = False
|
||||
def loader():
|
||||
nonlocal called
|
||||
called = True
|
||||
raise AssertionError("DB/application runtime must not load")
|
||||
with pytest.raises(SystemExit) as exc:
|
||||
sim.main(["--help"], runtime_loader=loader)
|
||||
assert exc.value.code == 0 and called is False
|
||||
|
||||
|
||||
def test_production_requires_explicit_flag(monkeypatch):
|
||||
monkeypatch.setenv("BLIF_FLOW_V2_MODE", "compare")
|
||||
with pytest.raises(RuntimeError, match="requires --production"):
|
||||
sim.run(Namespace(production_readonly_simulation=False), runtime_loader=runtime(Connection()))
|
||||
|
||||
|
||||
@pytest.mark.parametrize("database,readonly,mode,message", [
|
||||
("test_db", "on", "compare", "current_database"),
|
||||
("clientflow", "off", "compare", "transaction_read_only"),
|
||||
("clientflow", "on", None, "explicit BLIF_FLOW"),
|
||||
("clientflow", "on", "off", "BLIF_FLOW_V2_MODE=compare"),
|
||||
("clientflow", "on", "shadow", "BLIF_FLOW_V2_MODE=compare"),
|
||||
("clientflow", "on", "authoritative", "BLIF_FLOW_V2_MODE=compare"),
|
||||
])
|
||||
def test_production_guards(monkeypatch, database, readonly, mode, message):
|
||||
if mode is None:
|
||||
monkeypatch.delenv("BLIF_FLOW_V2_MODE", raising=False)
|
||||
else:
|
||||
monkeypatch.setenv("BLIF_FLOW_V2_MODE", mode)
|
||||
with pytest.raises(RuntimeError, match=message):
|
||||
sim.run(Namespace(production_readonly_simulation=True),
|
||||
runtime_loader=runtime(Connection(database, readonly)))
|
||||
|
||||
|
||||
def test_safe_simulation_uses_explicit_modes_without_changing_config_and_has_no_db_writes(monkeypatch, tmp_path):
|
||||
monkeypatch.setenv("BLIF_FLOW_V2_MODE", "compare")
|
||||
monkeypatch.setattr(sim, "JSON_PATH", tmp_path / "report.json")
|
||||
monkeypatch.setattr(sim, "TEXT_PATH", tmp_path / "report.txt")
|
||||
connection = Connection()
|
||||
modes = []
|
||||
report = sim.run(Namespace(production_readonly_simulation=True), runtime_loader=runtime(connection, modes=modes))
|
||||
assert report["configured_mode"] == "compare"
|
||||
assert sim.exit_code(report) == 0
|
||||
assert [sql for sql in connection.sql if sql.split(None, 1)[0].upper() in {"INSERT", "UPDATE", "DELETE", "MERGE"}] == []
|
||||
assert connection.sql[0] == "BEGIN READ ONLY" and connection.sql[-1] == "ROLLBACK"
|
||||
assert modes == ["off", "authoritative"]
|
||||
|
||||
|
||||
def safe_report():
|
||||
return sim.build_report(identity={}, configured_mode="compare",
|
||||
v1={"work_items": []}, v2={"work_items": []},
|
||||
named_decisions=named_decisions(), missing_projections=0)
|
||||
|
||||
|
||||
def test_unsafe_false_negative_is_nonzero():
|
||||
report = safe_report()
|
||||
report["unsafe_false_negatives"] = 1
|
||||
assert sim.exit_code(report) != 0
|
||||
|
||||
|
||||
def test_missing_projection_is_nonzero():
|
||||
report = safe_report()
|
||||
report["missing_projections"] = 1
|
||||
assert sim.exit_code(report) != 0
|
||||
|
||||
|
||||
def test_duplicate_current_card_is_nonzero():
|
||||
report = safe_report()
|
||||
report["duplicate_current_cards"] = 1
|
||||
assert sim.exit_code(report) != 0
|
||||
|
||||
|
||||
def test_named_case_failure_is_nonzero():
|
||||
report = sim.build_report(identity={}, configured_mode="compare", v1={"work_items": []},
|
||||
v2={"work_items": []}, named_decisions=named_decisions(True), missing_projections=0)
|
||||
assert sim.exit_code(report) != 0
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user