Compare commits
4 Commits
fix/blif-f
...
fix/blif-f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f96b1b3d6e | ||
|
|
213ba6b41b | ||
|
|
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",
|
||||
|
||||
196
app/authoritative_operational_adapter.py
Normal file
196
app/authoritative_operational_adapter.py
Normal 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)
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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]:
|
||||
|
||||
@@ -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
|
||||
|
||||
300
scripts/simulate_blif_flow_v2_authoritative.py
Normal file
300
scripts/simulate_blif_flow_v2_authoritative.py
Normal 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())
|
||||
146
tests/test_authoritative_operational_adapter.py
Normal file
146
tests/test_authoritative_operational_adapter.py
Normal 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)
|
||||
168
tests/test_blif_flow_v2_authoritative_simulation.py
Normal file
168
tests/test_blif_flow_v2_authoritative_simulation.py
Normal 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"}
|
||||
@@ -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