Compare commits
5 Commits
feat/blif-
...
ac385342f5
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac385342f5 | ||
|
|
fc092a96b6 | ||
|
|
7cc9fecbdc | ||
|
|
e456cbc0e4 | ||
|
|
e3b8750ebb |
@@ -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)
|
||||
@@ -19,6 +19,7 @@ from app.config import settings
|
||||
|
||||
FLOW_VERSION = "blif-flow-v2-shadow-3"
|
||||
WRITE_DATABASE_ALLOWLIST = frozenset({"clientflow_codex_test"})
|
||||
PROJECTION_WRITE_TABLES = frozenset({"opportunity_flow_state_v2", "opportunity_flow_transitions"})
|
||||
|
||||
|
||||
def _jsonable(value: Any) -> Any:
|
||||
@@ -78,16 +79,24 @@ def _projection_value(row: dict[str, Any], derived_at: datetime) -> dict[str, An
|
||||
return stable | {"source_fingerprint": fingerprint, "derived_at": derived_at}
|
||||
|
||||
|
||||
def _derive_all() -> list[dict[str, Any]]:
|
||||
def _derive_all(
|
||||
*,
|
||||
expected_database: str = "clientflow_codex_test",
|
||||
expected_user: str | None = "clientflow_codex_test",
|
||||
expected_opportunity_count: int | None = 328,
|
||||
require_opportunities: bool = False,
|
||||
) -> list[dict[str, Any]]:
|
||||
# Reuse the validated shadow evidence adapter without making it authoritative.
|
||||
from scripts.simulate_blif_flow_v2 import collect
|
||||
|
||||
report = collect(
|
||||
expected_database="clientflow_codex_test",
|
||||
expected_user="clientflow_codex_test",
|
||||
expected_database=expected_database,
|
||||
expected_user=expected_user,
|
||||
# collect() still opens its factual read phase with BEGIN READ ONLY;
|
||||
# the session default may be read-write in the isolated test database.
|
||||
require_read_only=False,
|
||||
expected_opportunity_count=expected_opportunity_count,
|
||||
require_opportunities=require_opportunities,
|
||||
)
|
||||
return list(report["opportunities"])
|
||||
|
||||
@@ -99,6 +108,10 @@ def rebuild_blif_flow_v2_projection(
|
||||
allowed_databases: frozenset[str] = WRITE_DATABASE_ALLOWLIST,
|
||||
target_schema: str = "public",
|
||||
connection: Any | None = None,
|
||||
derive_expected_database: str = "clientflow_codex_test",
|
||||
derive_expected_user: str | None = "clientflow_codex_test",
|
||||
expected_opportunity_count: int | None = 328,
|
||||
require_opportunities: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
"""Idempotently rebuild projection rows and state-change transitions.
|
||||
|
||||
@@ -114,7 +127,14 @@ def rebuild_blif_flow_v2_projection(
|
||||
if not re.fullmatch(r"[a-z_][a-z0-9_]*", target_schema):
|
||||
raise ValueError("invalid target_schema")
|
||||
|
||||
rows = list(derived_rows) if derived_rows is not None else _derive_all()
|
||||
rows = list(derived_rows) if derived_rows is not None else _derive_all(
|
||||
expected_database=derive_expected_database,
|
||||
expected_user=derive_expected_user,
|
||||
expected_opportunity_count=expected_opportunity_count,
|
||||
require_opportunities=require_opportunities,
|
||||
)
|
||||
if require_opportunities and not rows:
|
||||
raise RuntimeError("Flow v2 projection requires at least one opportunity")
|
||||
derived_at = datetime.now(timezone.utc)
|
||||
values = [_projection_value(row, derived_at) for row in rows]
|
||||
if len({value["opportunity_id"] for value in values}) != len(values):
|
||||
|
||||
@@ -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,19 +348,21 @@ def _opportunity_item(
|
||||
evidence_rows.append(row)
|
||||
|
||||
if action_code in NO_WORK_ACTION_CODES:
|
||||
if decision.get("suppress_current_card"):
|
||||
return None, independent_exceptions
|
||||
keep_authoritative_noncurrent = authoritative and _s(decision.get("operational_queue")) in {"waiting", "backlog"}
|
||||
if not keep_authoritative_noncurrent:
|
||||
if not _has_explicit_unresolved_review(evidence_rows):
|
||||
return None, independent_exceptions
|
||||
# The central decision normally wins, but an explicit unresolved review
|
||||
# marker must not be hidden by NO_ACTION/NOT_FOUND.
|
||||
# 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,
|
||||
**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}",
|
||||
"priority": "normal", "target_url": f"/opportunities/{opportunity_id}",
|
||||
}
|
||||
|
||||
matching_task = scheduled_call or _matching_task(evidence_rows, action_code)
|
||||
@@ -367,6 +370,7 @@ def _opportunity_item(
|
||||
primary = primary or (evidence_rows[0] if evidence_rows else {})
|
||||
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()
|
||||
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
|
||||
|
||||
@@ -1,228 +1,328 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Read-only compatibility/cutover audit for BLIF Flow v2."""
|
||||
"""Run a strictly read-only BLIF Flow v2 shadow/cutover audit."""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
from collections import Counter
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from typing import Any, Sequence
|
||||
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
os.chdir(ROOT)
|
||||
|
||||
from sqlalchemy import text
|
||||
PRODUCTION_DATABASE = "clientflow"
|
||||
TEST_DATABASE = "clientflow_codex_test"
|
||||
AUDIT = Path("/tmp/blif_flow_v2_production_shadow_audit.json")
|
||||
SEMANTIC = Path("/tmp/blif_flow_v2_production_semantic_compare.json")
|
||||
SUMMARY = Path("/tmp/blif_flow_v2_production_shadow_summary.txt")
|
||||
CLASSIFICATIONS = {
|
||||
"SEMANTICALLY_EQUIVALENT", "V1_OPERATIONAL_OVERRIDE", "V2_CORRECTS_V1",
|
||||
"LEGACY_ONLY", "REAL_CONFLICT", "MISSING_PROJECTION",
|
||||
}
|
||||
CURRENT_QUEUES = {"do_now", "review", "exception"}
|
||||
OVERRIDE_PRECEDENCE = {
|
||||
"scheduled_call", "due_followup", "future_followup", "integration_exception",
|
||||
"document_prerequisite", "fiscal_prerequisite", "safe_preserve_v1",
|
||||
}
|
||||
LEGACY_ACTIONS = {
|
||||
"CREATE_JASMIN_QUOTE", "NO_ACTION", "WAIT_CUSTOMER", "WAIT_PAYMENT",
|
||||
"WAIT_PRODUCTION", "WAIT_LOGISTICS", "WAIT_SUPPLIER", "WAIT_SCHEDULED_DATE",
|
||||
}
|
||||
REVIEW_ACTIONS = {
|
||||
"REVIEW", "REVIEW_REQUIRED", "REVIEW_MANUALLY", "REVIEW_RECONSTRUCTED_PROCESS",
|
||||
"RECONCILE_DOCUMENTS", "VALIDATE_FISCAL_CUSTOMER", "REVIEW_EXCEPTION",
|
||||
}
|
||||
NAMED_IDS = {
|
||||
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
|
||||
"PANORAMIC": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
|
||||
"ENGEXICON": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
|
||||
"CONSTRURECUP": "e3b23ac5-84db-4763-8a31-a684e873032c",
|
||||
"X_MAT_CANONICAL": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
|
||||
"X_MAT_DUPLICATE": "1816a06e-9a69-4a9b-9279-1263156892d3",
|
||||
"RZSOLAR_CANONICAL": "fd221608-e007-4043-a23d-07e0c119a345",
|
||||
"RZSOLAR_DUPLICATE": "434124fb-ac19-4d78-909a-55761d7e8daa",
|
||||
}
|
||||
|
||||
from app.db import engine
|
||||
from app.opportunity_next_action_service import (
|
||||
_load_v2_comparison_rows, compare_v1_v2_decisions, get_opportunity_next_actions,
|
||||
|
||||
def build_parser() -> argparse.ArgumentParser:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument(
|
||||
"--production-readonly-audit", action="store_true",
|
||||
help="explicitly authorize a read-only audit of database clientflow",
|
||||
)
|
||||
from scripts.simulate_blif_flow_v2 import collect
|
||||
return parser
|
||||
|
||||
|
||||
EXPECTED_DATABASE = "clientflow_codex_test"
|
||||
CUTOVER = Path("/tmp/blif_flow_v2_cutover_report.json")
|
||||
INVENTORY = Path("/tmp/blif_flow_v2_legacy_field_inventory.json")
|
||||
COMPARE = Path("/tmp/blif_flow_v2_compare_report.json")
|
||||
def validate_execution(
|
||||
*, production_readonly_audit: bool, database: str, transaction_read_only: str,
|
||||
mode: str,
|
||||
) -> None:
|
||||
normalized_mode = str(mode or "").strip().lower()
|
||||
if production_readonly_audit:
|
||||
if database != PRODUCTION_DATABASE:
|
||||
raise RuntimeError(
|
||||
f"--production-readonly-audit requires database {PRODUCTION_DATABASE!r}, found {database!r}"
|
||||
)
|
||||
if transaction_read_only != "on":
|
||||
raise RuntimeError("production audit requires transaction_read_only=on")
|
||||
if normalized_mode not in {"shadow", "compare"}:
|
||||
raise RuntimeError("production audit requires BLIF_FLOW_V2_MODE=shadow or compare")
|
||||
return
|
||||
if database == PRODUCTION_DATABASE:
|
||||
raise RuntimeError("production database requires explicit --production-readonly-audit opt-in")
|
||||
if database != TEST_DATABASE:
|
||||
raise RuntimeError(f"default audit requires database {TEST_DATABASE!r}, found {database!r}")
|
||||
if transaction_read_only != "on":
|
||||
raise RuntimeError("cutover audit requires an explicit READ ONLY transaction")
|
||||
|
||||
|
||||
FIELD_INVENTORY = [
|
||||
{
|
||||
"field": "stage", "kind": "derived_compatibility", "factual": False,
|
||||
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
|
||||
"app/reconciliation_service.py", "app/odoo_service.py", "app/external_reconciliation_sync.py"],
|
||||
"readers": ["app/domain/opportunity_flow/evidence.py (V1 only)", "app/workflow_guard.py",
|
||||
"app/admin_ui/pages/opportunities.py", "app/admin_ui/pages/orders.py",
|
||||
"app/admin_dashboard.py", "app/revenue_forecast_service.py"],
|
||||
"flow_v2_replacement": "opportunity_flow_state_v2.business_state",
|
||||
"compatibility_required": True, "one_time_migration_required": False,
|
||||
"eventually_deprecatable": True,
|
||||
"notes": "Keep stored value initially for V1 boards, reports, forecasts and action guards; never use as factual V2 input. Consumers need code migration, not historical mass rewrite.",
|
||||
},
|
||||
{
|
||||
"field": "lifecycle_state", "kind": "operational_compatibility", "factual": False,
|
||||
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
|
||||
"readers": ["app/operational_eligibility.py", "app/operations_service.py",
|
||||
"app/admin_ui/pages/opportunities.py"],
|
||||
"flow_v2_replacement": "business state plus active operational obligations",
|
||||
"compatibility_required": True, "one_time_migration_required": False,
|
||||
"eventually_deprecatable": True,
|
||||
"notes": "Awaiting/recovery/nurture presentation remains V1 compatibility. It cannot create a scheduled obligation without an active task.",
|
||||
},
|
||||
{
|
||||
"field": "next_follow_up_at", "kind": "denormalized_followup_compatibility", "factual": False,
|
||||
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
|
||||
"readers": ["app/operations_service.py", "app/admin_ui/pages/opportunities.py"],
|
||||
"flow_v2_replacement": "pending follow-up task action_code/due_at",
|
||||
"compatibility_required": True, "one_time_migration_required": False,
|
||||
"eventually_deprecatable": True,
|
||||
"notes": "Historical timestamp is inert without a pending follow-up task; retain all 44 values for now.",
|
||||
},
|
||||
{
|
||||
"field": "nurture_until", "kind": "denormalized_schedule_compatibility", "factual": False,
|
||||
"writers": ["app/opportunity_service.py"],
|
||||
"readers": ["app/admin_ui/pages/opportunities.py"],
|
||||
"flow_v2_replacement": "explicit pending REVIEW_NURTURE/follow-up task due_at",
|
||||
"compatibility_required": True, "one_time_migration_required": False,
|
||||
"eventually_deprecatable": True,
|
||||
},
|
||||
{
|
||||
"field": "follow_up_attempts", "kind": "historical_counter", "factual": False,
|
||||
"writers": ["app/opportunity_service.py", "app/followup_service.py"],
|
||||
"readers": ["app/admin_ui/pages/opportunities.py"],
|
||||
"flow_v2_replacement": "event/task history aggregation", "compatibility_required": True,
|
||||
"one_time_migration_required": False, "eventually_deprecatable": True,
|
||||
},
|
||||
{
|
||||
"field": "last_action_code", "kind": "last-known-action_compatibility", "factual": False,
|
||||
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
|
||||
"app/reconciliation_service.py", "app/odoo_service.py"],
|
||||
"readers": ["app/workflow_guard.py", "app/admin_dashboard.py",
|
||||
"app/admin_ui/pages/opportunities.py", "app/action_prompt.py"],
|
||||
"flow_v2_replacement": "opportunity_flow_state_v2.business_next_action plus OperationalEligibility",
|
||||
"compatibility_required": True, "one_time_migration_required": False,
|
||||
"eventually_deprecatable": True,
|
||||
},
|
||||
{
|
||||
"field": "last_task_id", "kind": "historical_pointer", "factual": False,
|
||||
"writers": ["app/opportunity_service.py", "app/reconciliation_service.py",
|
||||
"app/company_opportunity_linking.py"],
|
||||
"readers": ["app/workflow_guard.py"],
|
||||
"flow_v2_replacement": "active canonical task lookup", "compatibility_required": True,
|
||||
"one_time_migration_required": False, "eventually_deprecatable": True,
|
||||
},
|
||||
{
|
||||
"field": "pending_primary_* / pending_follow_up_*", "kind": "runtime_read_model", "factual": False,
|
||||
"writers": ["none (SQL projections in app/opportunity_service.py)"],
|
||||
"readers": ["app/admin_ui/pages/opportunities.py"],
|
||||
"flow_v2_replacement": "active task lookup remains an explicit operational overlay",
|
||||
"compatibility_required": True, "one_time_migration_required": False,
|
||||
"eventually_deprecatable": False,
|
||||
"notes": "These virtual fields correctly derive active obligations and are not stored opportunity state.",
|
||||
},
|
||||
]
|
||||
def _code(value: Any) -> str:
|
||||
return str(value or "").strip().upper()
|
||||
|
||||
SOURCE_OF_TRUTH = [
|
||||
["business process state", "opportunity_flow_state_v2.business_state"],
|
||||
["business next action", "opportunity_flow_state_v2.business_next_action"],
|
||||
["time-sensitive operational queue", "OperationalEligibility"],
|
||||
["scheduled follow-up", "pending task action_code + due_at"],
|
||||
["formal document state", "commercial_documents + active opportunity_document_links"],
|
||||
["payment", "confirmed factual payment evidence/operation link"],
|
||||
["Odoo execution", "operation_links / factual Odoo evidence"],
|
||||
["messages/customer response", "messages and communications chronology"],
|
||||
["legacy stage", "compatibility only"],
|
||||
["tasks", "operator obligations/history; never business fact proof"],
|
||||
]
|
||||
|
||||
def classify_semantic_difference(
|
||||
*, v1_state: str | None, v1_action: str | None,
|
||||
v2_state: str | None, v2_action: str | None,
|
||||
v2_projection_present: bool = True,
|
||||
is_duplicate_representation: bool = False,
|
||||
operational_action: str | None = None,
|
||||
operational_precedence: str | None = None,
|
||||
operational_queue: str | None = None,
|
||||
v2_confidence: str | None = None,
|
||||
v2_diagnostic_status: str | None = None,
|
||||
) -> tuple[str, str]:
|
||||
"""Conservatively compare meanings rather than raw action vocabulary."""
|
||||
if not v2_projection_present:
|
||||
return "MISSING_PROJECTION", "No persisted Flow v2 projection exists for this opportunity."
|
||||
v1, v2, effective = _code(v1_action), _code(v2_action), _code(operational_action)
|
||||
precedence = str(operational_precedence or "").strip().lower()
|
||||
queue = str(operational_queue or "").strip().lower()
|
||||
if is_duplicate_representation:
|
||||
return "V2_CORRECTS_V1", "Material identity suppresses a duplicate representation without deleting evidence."
|
||||
if precedence in OVERRIDE_PRECEDENCE and effective and (v1 == effective or queue in CURRENT_QUEUES | {"waiting"}):
|
||||
return "V1_OPERATIONAL_OVERRIDE", f"Explicit operational precedence {precedence} validly overlays the V2 business transition."
|
||||
if v1 == v2 and v1:
|
||||
return "SEMANTICALLY_EQUIVALENT", "V1 and V2 select the same action."
|
||||
if not v1 and not v2:
|
||||
return "SEMANTICALLY_EQUIVALENT", "Neither model has a current business action."
|
||||
if v1 in REVIEW_ACTIONS and (v2 in REVIEW_ACTIONS or _code(v2_state) in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED"}):
|
||||
return "SEMANTICALLY_EQUIVALENT", "Both decisions require review/blocker handling."
|
||||
if v1 in {"NO_ACTION", ""} and not v2 and _code(v2_state) in {"COMPLETED", "LOST", "NO_INTEREST"}:
|
||||
return "SEMANTICALLY_EQUIVALENT", "Both decisions represent a terminal/non-current process."
|
||||
if v1.startswith("FOLLOW_UP_") and _code(v2_state) in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}:
|
||||
return "V1_OPERATIONAL_OVERRIDE", "A current follow-up obligation overlays a waiting V2 business state."
|
||||
if v1 == "CONFIRM_PAYMENT" and _code(v2_state) == "AWAITING_PAYMENT" and not v2:
|
||||
return "SEMANTICALLY_EQUIVALENT", "Both decisions mean payment remains outstanding; V1 names the compatibility action."
|
||||
if v1 in {"VALIDATE_FISCAL_CUSTOMER", "RECONCILE_DOCUMENTS"} and not effective and (
|
||||
queue in {"not_current", "backlog"} or _code(v2_state) in {"INQUIRY", "AWAITING_CUSTOMER", "AWAITING_PAYMENT"}
|
||||
):
|
||||
return "LEGACY_ONLY", "Legacy data-hygiene/blocker vocabulary is not a current factual V2 obligation."
|
||||
if precedence in {"safe_diagnostic_only", "safe_ambiguous_review"}:
|
||||
return "V2_CORRECTS_V1", "SAFE V2 prevents ambiguous historical compatibility state from creating current work."
|
||||
if _code(v2_state) in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED"}:
|
||||
return "V2_CORRECTS_V1", "V2 converts conflicting factual history into an explicit protected review state."
|
||||
if effective and effective == v2 and v1 != v2:
|
||||
return "V2_CORRECTS_V1", "The safe operational action follows the factual V2 transition rather than the legacy action."
|
||||
if v1 in LEGACY_ACTIONS:
|
||||
if v2 and _code(v2_state) not in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}:
|
||||
return "V2_CORRECTS_V1", "V2 replaces a legacy compatibility action with a factual business transition."
|
||||
return "LEGACY_ONLY", "V1 action is compatibility vocabulary with no native V2 business transition."
|
||||
if v2 and str(v2_confidence or "").lower() == "high" and str(v2_diagnostic_status or "").lower() == "clear":
|
||||
return "V2_CORRECTS_V1", "High-confidence factual V2 transition corrects a different legacy action."
|
||||
if effective and v1 == effective:
|
||||
return "V1_OPERATIONAL_OVERRIDE", "V1 matches the safe operational overlay rather than the business action."
|
||||
return "REAL_CONFLICT", "The available evidence does not establish equivalence, a valid override, or a safe V2 correction."
|
||||
|
||||
|
||||
def _write(path: Path, value: Any) -> None:
|
||||
path.write_text(json.dumps(value, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
|
||||
|
||||
|
||||
def _identity_and_ids() -> tuple[dict[str, str], list[str], int]:
|
||||
def _queue_totals(rows: list[dict[str, Any]], key: str) -> dict[str, int]:
|
||||
counts = Counter(str(row[key].get("operational_queue") or row[key].get("effective_operational_queue") or "not_current") for row in rows)
|
||||
return {
|
||||
"current": sum(counts[name] for name in CURRENT_QUEUES),
|
||||
"do_now": counts["do_now"], "review": counts["review"],
|
||||
"waiting": counts["waiting"], "backlog": counts["backlog"],
|
||||
"not_current": counts["not_current"],
|
||||
}
|
||||
|
||||
|
||||
def _operation_totals(source: dict[str, Any]) -> dict[str, int]:
|
||||
return {"current": int(source.get("current_work") or 0),
|
||||
"do_now": int(source.get("do_now") or 0),
|
||||
"review": int(source.get("review") or 0),
|
||||
"waiting": int(source.get("waiting") or 0),
|
||||
"backlog": int(source.get("backlog") or 0),
|
||||
"not_current": int(source.get("not_current") or 0)}
|
||||
|
||||
|
||||
def _database_snapshot(engine: Any) -> tuple[dict[str, str], list[str], dict[str, dict[str, Any]]]:
|
||||
from sqlalchemy import text
|
||||
with engine.connect() as conn:
|
||||
conn.exec_driver_sql("BEGIN READ ONLY")
|
||||
try:
|
||||
identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one()
|
||||
if identity[0] != EXPECTED_DATABASE or identity[2] != "on":
|
||||
raise RuntimeError(f"refusing unexpected/non-read-only database: {identity!r}")
|
||||
ids = [str(value) for value in conn.execute(text("SELECT opportunity_id FROM opportunity_flow_state_v2 ORDER BY opportunity_id")).scalars()]
|
||||
historical = conn.execute(text("""
|
||||
SELECT count(*) FROM opportunities o
|
||||
WHERE o.next_follow_up_at IS NOT NULL
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM tasks t WHERE t.opportunity_id=o.id AND t.status='pending'
|
||||
AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN
|
||||
('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
|
||||
)
|
||||
""")).scalar_one()
|
||||
identity = conn.execute(text(
|
||||
"SELECT current_database(),current_user,current_setting('transaction_read_only')"
|
||||
)).one()
|
||||
opportunity_ids = [str(value) for value in conn.execute(text(
|
||||
"SELECT id FROM opportunities ORDER BY id"
|
||||
)).scalars()]
|
||||
rows = conn.execute(text("""
|
||||
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, evidence_refs, flow_version
|
||||
FROM opportunity_flow_state_v2 ORDER BY opportunity_id
|
||||
""")).mappings().all()
|
||||
finally:
|
||||
conn.rollback()
|
||||
return {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]}, ids, int(historical)
|
||||
return (
|
||||
{"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
|
||||
opportunity_ids, {str(row["opportunity_id"]): dict(row) for row in rows},
|
||||
)
|
||||
|
||||
|
||||
def main() -> None:
|
||||
identity, ids, historical_timestamps = _identity_and_ids()
|
||||
if len(ids) != 328:
|
||||
raise RuntimeError(f"expected 328 persisted Flow v2 opportunities, found {len(ids)}")
|
||||
v1 = get_opportunity_next_actions(ids)
|
||||
v2 = _load_v2_comparison_rows(ids)
|
||||
comparisons = compare_v1_v2_decisions(v1, v2)
|
||||
counts = Counter("agree" if row["action_agrees"] else "different" for row in comparisons)
|
||||
counts["missing_projection"] = sum(not row["projection_present"] for row in comparisons)
|
||||
compare_report = {
|
||||
"generated_at": datetime.now(timezone.utc), "database": identity,
|
||||
"mode_semantics": "V1 returned; persisted V2 observed; no business mutation",
|
||||
"opportunity_count": len(comparisons), "counts": dict(counts), "comparisons": comparisons,
|
||||
def run_audit(*, production_readonly_audit: bool) -> dict[str, Any]:
|
||||
# Imports occur only after argparse, so --help cannot initialize DB code.
|
||||
from app.config import settings
|
||||
from app.db import engine
|
||||
from scripts.simulate_blif_flow_v2 import collect
|
||||
|
||||
identity, opportunity_ids, persisted = _database_snapshot(engine)
|
||||
mode = str(settings.blif_flow_v2_mode or "").strip().lower()
|
||||
validate_execution(
|
||||
production_readonly_audit=production_readonly_audit,
|
||||
database=identity["database"], transaction_read_only=identity["transaction_read_only"],
|
||||
mode=mode,
|
||||
)
|
||||
preflight = {**identity, "BLIF_FLOW_V2_MODE": mode,
|
||||
"production_readonly_opt_in": production_readonly_audit}
|
||||
print(json.dumps(preflight, sort_keys=True), flush=True)
|
||||
|
||||
expected_count = None if production_readonly_audit else 328
|
||||
projection = collect(
|
||||
expected_database=identity["database"], expected_user=identity["user"],
|
||||
require_read_only=production_readonly_audit,
|
||||
expected_opportunity_count=expected_count, require_opportunities=True,
|
||||
)
|
||||
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
|
||||
if set(records) != set(opportunity_ids):
|
||||
raise RuntimeError("factual collector did not return the complete opportunity universe")
|
||||
|
||||
comparisons = []
|
||||
for opportunity_id in opportunity_ids:
|
||||
record = records[opportunity_id]
|
||||
v1, operational = record["v1"], record["safe_v2"]
|
||||
v2 = persisted.get(opportunity_id)
|
||||
classification, reason = classify_semantic_difference(
|
||||
v1_state=v1.get("commercial_stage"), v1_action=v1.get("current_action"),
|
||||
v2_state=(v2 or {}).get("business_state"), v2_action=(v2 or {}).get("business_next_action"),
|
||||
v2_projection_present=v2 is not None,
|
||||
is_duplicate_representation=bool((v2 or {}).get("is_duplicate_representation")),
|
||||
operational_action=operational.get("effective_operational_action"),
|
||||
operational_precedence=operational.get("precedence"),
|
||||
operational_queue=operational.get("effective_operational_queue"),
|
||||
v2_confidence=(v2 or {}).get("confidence"),
|
||||
v2_diagnostic_status=(v2 or {}).get("diagnostic_status"),
|
||||
)
|
||||
comparisons.append({
|
||||
"opportunity_id": opportunity_id, "title": record.get("title"),
|
||||
"customer_name": record.get("customer"), "classification": classification,
|
||||
"classification_reason": reason,
|
||||
"v1": {"state": v1.get("commercial_stage"), "action": v1.get("current_action"),
|
||||
"queue": v1.get("operational_queue"), "reason_code": v1.get("reason")},
|
||||
"v2": {"business_state": (v2 or {}).get("business_state"),
|
||||
"business_next_action": (v2 or {}).get("business_next_action"),
|
||||
"reason_code": (v2 or {}).get("reason_code"),
|
||||
"diagnostic_status": (v2 or {}).get("diagnostic_status"),
|
||||
"confidence": (v2 or {}).get("confidence")},
|
||||
"operational_override": {"action": operational.get("effective_operational_action"),
|
||||
"queue": operational.get("effective_operational_queue"),
|
||||
"precedence": operational.get("precedence")},
|
||||
"material_process_key": (v2 or {}).get("material_process_key"),
|
||||
"canonical_opportunity_id": (v2 or {}).get("canonical_opportunity_id"),
|
||||
"is_duplicate_representation": bool((v2 or {}).get("is_duplicate_representation")),
|
||||
"evidence_refs": (v2 or {}).get("evidence_refs", []),
|
||||
})
|
||||
counts = Counter(row["classification"] for row in comparisons)
|
||||
for name in CLASSIFICATIONS:
|
||||
counts.setdefault(name, 0)
|
||||
real_conflicts = [row for row in comparisons if row["classification"] == "REAL_CONFLICT"]
|
||||
missing = [row for row in comparisons if row["classification"] == "MISSING_PROJECTION"]
|
||||
semantic_report = {
|
||||
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
|
||||
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
|
||||
"classification_counts": dict(sorted(counts.items())),
|
||||
"real_conflicts": real_conflicts, "missing_projections": missing,
|
||||
"comparisons": comparisons,
|
||||
}
|
||||
_write(COMPARE, compare_report)
|
||||
_write(SEMANTIC, semantic_report)
|
||||
|
||||
inventory = {
|
||||
"generated_at": datetime.now(timezone.utc), "database": identity,
|
||||
"fields": FIELD_INVENTORY,
|
||||
"summary": {
|
||||
"fields_requiring_data_migration": [],
|
||||
"compatibility_only_fields": [row["field"] for row in FIELD_INVENTORY if row["compatibility_required"]],
|
||||
"safe_to_deprecate_after_consumer_migration": [row["field"] for row in FIELD_INVENTORY if row["eventually_deprecatable"]],
|
||||
"historical_followup_timestamps_preserved": historical_timestamps,
|
||||
},
|
||||
}
|
||||
_write(INVENTORY, inventory)
|
||||
|
||||
projection = collect(expected_database=EXPECTED_DATABASE, expected_user=identity["user"], require_read_only=False)
|
||||
named = {}
|
||||
for name in ("INSTALBEIRA", "PANORAMIC SUCCESS", "ENGEXICON", "CONSTRURECUP", "X MAT", "RZSOLAR"):
|
||||
matches = [row for row in projection["opportunities"]
|
||||
if name.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()]
|
||||
named[name] = [{"opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"],
|
||||
"canonical_process_id": row["canonical_process_id"], "business_state": row["safe_v2"]["business_state"],
|
||||
"effective_action": row["safe_v2"]["effective_operational_action"],
|
||||
"queue": row["safe_v2"]["effective_operational_queue"],
|
||||
"duplicate_suppressed": row["safe_v2"]["precedence"] == "duplicate_representation"}
|
||||
for row in matches]
|
||||
cutover = {
|
||||
"generated_at": datetime.now(timezone.utc), "database": identity,
|
||||
"source_of_truth": [{"concern": concern, "source": source} for concern, source in SOURCE_OF_TRUTH],
|
||||
"data_migration": {"opportunities_requiring_mutation_before_shadow": 0,
|
||||
"opportunities_requiring_no_mutation_before_shadow": len(ids),
|
||||
"broad_repair_required": False,
|
||||
"stage_write_migration_required": False,
|
||||
"lifecycle_write_migration_required": False,
|
||||
"followup_timestamp_cleanup_required": False},
|
||||
"mode_contract": {
|
||||
"off": "No Flow v2 derivation or projection writes are triggered by runtime reads.",
|
||||
"shadow": "V1 returned; explicit projection rebuild is additive/idempotent; no tasks, stage, or UI behavior changed.",
|
||||
"compare": "V1 returned; V2 projection read and structured comparison logged; no disagreement writes.",
|
||||
"authoritative": "Disabled and fail-closed. No activation performed.",
|
||||
by_id = {row["opportunity_id"]: row for row in comparisons}
|
||||
for name, opportunity_id in NAMED_IDS.items():
|
||||
row = by_id[opportunity_id]
|
||||
named[name] = row
|
||||
# Use the complete canonical Operations candidate universe, including
|
||||
# preserved standalone obligations, rather than opportunity cards alone.
|
||||
v1_metrics = _operation_totals(projection["v1_totals"])
|
||||
safe_metrics = _operation_totals(projection["safe_v2_totals"])
|
||||
audit = {
|
||||
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
|
||||
"read_only": True, "business_writes": 0, "projection_writes": 0,
|
||||
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
|
||||
"canonical_count": sum(not row.get("is_duplicate_representation") for row in persisted.values()),
|
||||
"duplicate_representation_count": sum(bool(row.get("is_duplicate_representation")) for row in persisted.values()),
|
||||
"operations_metrics": {"v1": v1_metrics, "safe_v2": safe_metrics},
|
||||
"semantic_classification_counts": dict(sorted(counts.items())),
|
||||
"authoritative_cutover_blockers": {
|
||||
"real_conflicts": len(real_conflicts), "missing_projections": len(missing),
|
||||
"blocked": bool(real_conflicts or missing),
|
||||
},
|
||||
"production_shadow_blockers": [
|
||||
"Production migration 011 presence was not and must not be checked from this environment.",
|
||||
"Deployment must provide an explicit projection rebuild cadence and comparison-log monitoring/retention.",
|
||||
],
|
||||
"mode_off_deployment_blockers": [],
|
||||
"shadow_compare_activation_prerequisites": [
|
||||
"Install migration 011 through the normal controlled production migration process.",
|
||||
"Run mode=off first, then explicitly rebuild projection in shadow.",
|
||||
"Monitor structured disagreement rates and missing projections before compare enablement.",
|
||||
],
|
||||
"authoritative_activation_blockers": [
|
||||
"Authoritative switch intentionally raises and has no enabled code path.",
|
||||
"Operations/UI must consume V2 business state/action followed by OperationalEligibility.",
|
||||
"Legacy stage consumers in orders, forecasts, dashboard, workflow guards and opportunity columns require migration or explicit compatibility adapters.",
|
||||
"Explicit operational override precedence must be implemented for CALL_CUSTOMER, due follow-up, SUPPORT, SEND_INFO, manual review and blockers.",
|
||||
"Production shadow/compare observation, rollback criteria and zero-false-negative acceptance must be completed.",
|
||||
],
|
||||
"legacy_findings_are": "field-level compatibility findings, not required row mutations",
|
||||
"historical_followup_timestamps_preserved": historical_timestamps,
|
||||
"compare_summary": compare_report["counts"], "named_cases": named,
|
||||
"real_conflicts": real_conflicts, "missing_projections": missing,
|
||||
"named_cases": named,
|
||||
}
|
||||
_write(CUTOVER, cutover)
|
||||
print(json.dumps({"database": identity, "opportunities": len(ids),
|
||||
"compare": compare_report["counts"], "historical_timestamps": historical_timestamps,
|
||||
"mutation_required": 0, "outputs": [str(CUTOVER), str(INVENTORY), str(COMPARE)]}, indent=2))
|
||||
_write(AUDIT, audit)
|
||||
lines = [
|
||||
"BLIF FLOW V2 PRODUCTION SHADOW READ-ONLY AUDIT", "",
|
||||
f"database: {identity['database']}", f"user: {identity['user']}",
|
||||
f"transaction_read_only: {identity['transaction_read_only']}",
|
||||
f"BLIF_FLOW_V2_MODE: {mode}",
|
||||
f"production_readonly_opt_in: {production_readonly_audit}", "",
|
||||
f"opportunities: {len(opportunity_ids)}", f"projections: {len(persisted)}",
|
||||
f"canonical: {audit['canonical_count']}",
|
||||
f"duplicate representations: {audit['duplicate_representation_count']}", "",
|
||||
f"V1 metrics: {json.dumps(v1_metrics, sort_keys=True)}",
|
||||
f"SAFE V2 metrics: {json.dumps(safe_metrics, sort_keys=True)}", "",
|
||||
f"semantic classifications: {json.dumps(dict(sorted(counts.items())), sort_keys=True)}",
|
||||
f"REAL_CONFLICT blockers: {len(real_conflicts)}",
|
||||
f"MISSING_PROJECTION blockers: {len(missing)}",
|
||||
"business writes: 0", "projection writes: 0",
|
||||
]
|
||||
SUMMARY.write_text("\n".join(lines) + "\n", encoding="utf-8")
|
||||
return audit
|
||||
|
||||
|
||||
def main(argv: Sequence[str] | None = None) -> int:
|
||||
# --help exits here before application/database imports or connections.
|
||||
args = build_parser().parse_args(argv)
|
||||
result = run_audit(production_readonly_audit=args.production_readonly_audit)
|
||||
print(json.dumps({
|
||||
"database": result["preflight"]["database"],
|
||||
"opportunities": result["opportunity_count"], "projections": result["projection_count"],
|
||||
"semantic_classifications": result["semantic_classification_counts"],
|
||||
"authoritative_cutover_blockers": result["authoritative_cutover_blockers"],
|
||||
"outputs": [str(AUDIT), str(SEMANTIC), str(SUMMARY)],
|
||||
}, indent=2, sort_keys=True))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
raise SystemExit(main())
|
||||
|
||||
@@ -1,18 +1,84 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Rebuild additive BLIF Flow v2 projection tables in safe shadow mode."""
|
||||
"""Rebuild additive BLIF Flow v2 projection tables with explicit safeguards."""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Sequence
|
||||
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
os.chdir(ROOT)
|
||||
|
||||
|
||||
def build_parser() -> argparse.ArgumentParser:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument(
|
||||
"--production-shadow", action="store_true",
|
||||
help="explicitly authorize projection-only shadow writes to database clientflow",
|
||||
)
|
||||
return parser
|
||||
|
||||
|
||||
def validate_execution(*, production_shadow: bool, mode: str, database: str) -> None:
|
||||
"""Validate CLI intent independently of DATABASE_URL inference."""
|
||||
normalized_mode = str(mode or "").strip().lower()
|
||||
if production_shadow:
|
||||
if normalized_mode != "shadow":
|
||||
raise RuntimeError("--production-shadow requires BLIF_FLOW_V2_MODE=shadow")
|
||||
if database != "clientflow":
|
||||
raise RuntimeError(
|
||||
f"--production-shadow requires database 'clientflow', found {database!r}"
|
||||
)
|
||||
return
|
||||
if database == "clientflow":
|
||||
raise RuntimeError("production database requires explicit --production-shadow opt-in")
|
||||
|
||||
|
||||
def main(argv: Sequence[str] | None = None) -> int:
|
||||
# argparse handles --help and exits before any application/DB import below.
|
||||
args = build_parser().parse_args(argv)
|
||||
|
||||
from sqlalchemy import text
|
||||
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
|
||||
from app.config import settings
|
||||
from app.db import engine
|
||||
|
||||
mode = str(settings.blif_flow_v2_mode or "").strip().lower()
|
||||
with engine.connect() as conn:
|
||||
identity = conn.execute(text(
|
||||
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
||||
)).one()
|
||||
database, user, transaction_read_only = identity
|
||||
validate_execution(
|
||||
production_shadow=args.production_shadow, mode=mode, database=database,
|
||||
)
|
||||
print(json.dumps({
|
||||
"database": database, "user": user, "blif_flow_v2_mode": mode,
|
||||
"transaction_read_only": transaction_read_only,
|
||||
"production_shadow_opt_in": args.production_shadow,
|
||||
}, sort_keys=True), flush=True)
|
||||
|
||||
if args.production_shadow:
|
||||
result = rebuild_blif_flow_v2_projection(
|
||||
mode="shadow",
|
||||
# Production is deliberately scoped to this invocation; the
|
||||
# module-level default allowlist remains test-only.
|
||||
allowed_databases=frozenset({"clientflow"}),
|
||||
derive_expected_database="clientflow",
|
||||
derive_expected_user=user,
|
||||
expected_opportunity_count=None,
|
||||
require_opportunities=True,
|
||||
)
|
||||
else:
|
||||
result = rebuild_blif_flow_v2_projection()
|
||||
print(json.dumps(result, indent=2, sort_keys=True))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
print(json.dumps(rebuild_blif_flow_v2_projection(), indent=2, sort_keys=True))
|
||||
raise SystemExit(main())
|
||||
|
||||
@@ -99,6 +99,15 @@ def _load(
|
||||
) -> dict[str, Any]:
|
||||
with engine.connect() as conn:
|
||||
conn = conn.execution_options(isolation_level="AUTOCOMMIT")
|
||||
transaction_started = False
|
||||
try:
|
||||
# A fresh connection commonly reports transaction_read_only=off.
|
||||
# Establish the protected transaction on the same connection used
|
||||
# for every factual read before validating a strict audit.
|
||||
if require_read_only:
|
||||
conn.execute(text("BEGIN READ ONLY"))
|
||||
transaction_started = True
|
||||
|
||||
identity = conn.execute(text(
|
||||
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
||||
)).one()
|
||||
@@ -106,8 +115,13 @@ def _load(
|
||||
raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
|
||||
if require_read_only and identity[2] != "on":
|
||||
raise RuntimeError(f"read-only simulation requires transaction_read_only=on: {identity!r}")
|
||||
|
||||
# Preserve the development/test path's established ordering: check
|
||||
# its identity first, then protect the factual reads themselves.
|
||||
if not require_read_only:
|
||||
conn.execute(text("BEGIN READ ONLY"))
|
||||
try:
|
||||
transaction_started = True
|
||||
|
||||
opportunities = [dict(row) for row in conn.execute(text("""
|
||||
SELECT o.*, o.id::text AS id, o.local_customer_id::text,
|
||||
c.name AS linked_customer_name, c.tax_id, c.email AS fiscal_email,
|
||||
@@ -157,6 +171,7 @@ def _load(
|
||||
WHERE status IN ('open','needs_review','conflict') ORDER BY created_at
|
||||
""")).mappings()]
|
||||
finally:
|
||||
if transaction_started:
|
||||
conn.execute(text("ROLLBACK"))
|
||||
return {
|
||||
"identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
|
||||
@@ -530,6 +545,8 @@ def collect(
|
||||
*, expected_database: str = "clientflow_codex_shadow",
|
||||
expected_user: str | None = "clientflow_codex",
|
||||
require_read_only: bool = True,
|
||||
expected_opportunity_count: int | None = None,
|
||||
require_opportunities: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
data = _load(
|
||||
expected_database=expected_database,
|
||||
@@ -537,8 +554,12 @@ def collect(
|
||||
require_read_only=require_read_only,
|
||||
)
|
||||
opportunities = data["opportunities"]
|
||||
if len(opportunities) != 328:
|
||||
raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}")
|
||||
if expected_opportunity_count is not None and len(opportunities) != expected_opportunity_count:
|
||||
raise RuntimeError(
|
||||
f"expected {expected_opportunity_count} opportunities, found {len(opportunities)}"
|
||||
)
|
||||
if require_opportunities and not opportunities:
|
||||
raise RuntimeError("Flow v2 projection requires at least one opportunity")
|
||||
ids = [_s(opp["id"]) for opp in opportunities]
|
||||
v1_decisions = get_opportunity_next_actions(ids)
|
||||
operations = get_operations_summary(limit=200)
|
||||
|
||||
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):
|
||||
|
||||
110
tests/test_blif_flow_v2_production_readonly_audit.py
Normal file
110
tests/test_blif_flow_v2_production_readonly_audit.py
Normal file
@@ -0,0 +1,110 @@
|
||||
from inspect import getsource
|
||||
|
||||
import pytest
|
||||
|
||||
import scripts.audit_blif_flow_v2_cutover as audit
|
||||
|
||||
|
||||
def classify(**overrides):
|
||||
values = dict(
|
||||
v1_state="INFO_SENT", v1_action="SEND_INFO",
|
||||
v2_state="INQUIRY", v2_action="SEND_INFO",
|
||||
v2_projection_present=True, is_duplicate_representation=False,
|
||||
operational_action="SEND_INFO", operational_precedence="business_transition",
|
||||
operational_queue="do_now", v2_confidence="high", v2_diagnostic_status="clear",
|
||||
)
|
||||
values.update(overrides)
|
||||
return audit.classify_semantic_difference(**values)[0]
|
||||
|
||||
|
||||
def test_help_performs_zero_database_work(monkeypatch, capsys):
|
||||
monkeypatch.setattr(audit, "run_audit", lambda **kwargs: pytest.fail("help reached DB audit"))
|
||||
with pytest.raises(SystemExit) as exc:
|
||||
audit.main(["--help"])
|
||||
assert exc.value.code == 0
|
||||
assert "--production-readonly-audit" in capsys.readouterr().out
|
||||
|
||||
|
||||
def test_production_requires_explicit_opt_in():
|
||||
with pytest.raises(RuntimeError, match="explicit --production-readonly-audit"):
|
||||
audit.validate_execution(production_readonly_audit=False, database="clientflow",
|
||||
transaction_read_only="on", mode="shadow")
|
||||
|
||||
|
||||
def test_production_database_must_be_clientflow():
|
||||
with pytest.raises(RuntimeError, match="requires database 'clientflow'"):
|
||||
audit.validate_execution(production_readonly_audit=True, database="clientflow_codex_test",
|
||||
transaction_read_only="on", mode="shadow")
|
||||
|
||||
|
||||
def test_production_transaction_must_be_read_only():
|
||||
with pytest.raises(RuntimeError, match="transaction_read_only=on"):
|
||||
audit.validate_execution(production_readonly_audit=True, database="clientflow",
|
||||
transaction_read_only="off", mode="shadow")
|
||||
|
||||
|
||||
def test_production_audit_contains_no_sql_writes():
|
||||
source = (getsource(audit._database_snapshot) + getsource(audit.run_audit)).upper()
|
||||
for verb in ("INSERT ", "UPDATE ", "DELETE ", "CREATE ", "ALTER ", "DROP ", "TRUNCATE "):
|
||||
assert verb not in source
|
||||
assert 'BEGIN READ ONLY' in source
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["shadow", "compare"])
|
||||
def test_shadow_and_compare_are_allowed(mode):
|
||||
audit.validate_execution(production_readonly_audit=True, database="clientflow",
|
||||
transaction_read_only="on", mode=mode)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["", "off", "authoritative"])
|
||||
def test_unsafe_production_modes_are_refused(mode):
|
||||
with pytest.raises(RuntimeError, match="requires BLIF_FLOW_V2_MODE=shadow or compare"):
|
||||
audit.validate_execution(production_readonly_audit=True, database="clientflow",
|
||||
transaction_read_only="on", mode=mode)
|
||||
|
||||
|
||||
def test_missing_projection_is_reported():
|
||||
assert classify(v2_projection_present=False, v2_state=None, v2_action=None) == "MISSING_PROJECTION"
|
||||
|
||||
|
||||
def test_semantically_equivalent_is_classified():
|
||||
assert classify() == "SEMANTICALLY_EQUIVALENT"
|
||||
assert classify(v1_action="NO_ACTION", v2_state="COMPLETED", v2_action=None,
|
||||
operational_action=None, operational_queue="not_current") == "SEMANTICALLY_EQUIVALENT"
|
||||
|
||||
|
||||
def test_operational_override_is_classified():
|
||||
assert classify(v1_action="FOLLOW_UP_CUSTOMER_REVIEW", v2_state="AWAITING_CUSTOMER",
|
||||
v2_action=None, operational_action="FOLLOW_UP_CUSTOMER_REVIEW",
|
||||
operational_precedence="due_followup") == "V1_OPERATIONAL_OVERRIDE"
|
||||
|
||||
|
||||
def test_v2_correction_is_classified():
|
||||
assert classify(v1_action="SEND_INVOICE", v2_state="PROFORMA_REQUIRED",
|
||||
v2_action="CREATE_PROFORMA") == "V2_CORRECTS_V1"
|
||||
assert classify(v1_action="REVIEW_RECONSTRUCTED_PROCESS", v2_state="REVIEW_REQUIRED",
|
||||
v2_action="REVIEW_REQUIRED", is_duplicate_representation=True) == "V2_CORRECTS_V1"
|
||||
|
||||
|
||||
def test_legacy_only_is_classified():
|
||||
assert classify(v1_action="CREATE_JASMIN_QUOTE", v2_state="AWAITING_CUSTOMER",
|
||||
v2_action=None, v2_confidence="medium") == "LEGACY_ONLY"
|
||||
|
||||
|
||||
def test_real_conflict_is_classified():
|
||||
assert classify(v1_action="UNMAPPED_ACTION", v2_state="INQUIRY",
|
||||
v2_action=None, v2_confidence="medium",
|
||||
v2_diagnostic_status="ambiguous", operational_action=None) == "REAL_CONFLICT"
|
||||
|
||||
|
||||
def test_named_case_ids_and_expected_semantics_are_frozen():
|
||||
assert audit.NAMED_IDS == {
|
||||
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
|
||||
"PANORAMIC": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
|
||||
"ENGEXICON": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
|
||||
"CONSTRURECUP": "e3b23ac5-84db-4763-8a31-a684e873032c",
|
||||
"X_MAT_CANONICAL": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
|
||||
"X_MAT_DUPLICATE": "1816a06e-9a69-4a9b-9279-1263156892d3",
|
||||
"RZSOLAR_CANONICAL": "fd221608-e007-4043-a23d-07e0c119a345",
|
||||
"RZSOLAR_DUPLICATE": "434124fb-ac19-4d78-909a-55761d7e8daa",
|
||||
}
|
||||
104
tests/test_blif_flow_v2_production_shadow_rebuild.py
Normal file
104
tests/test_blif_flow_v2_production_shadow_rebuild.py
Normal file
@@ -0,0 +1,104 @@
|
||||
from inspect import getsource
|
||||
|
||||
import pytest
|
||||
|
||||
import app.blif_flow_v2_projection_service as projection
|
||||
import scripts.rebuild_blif_flow_v2_projection as cli
|
||||
import scripts.simulate_blif_flow_v2 as simulator
|
||||
|
||||
|
||||
def test_help_performs_no_rebuild(monkeypatch, capsys):
|
||||
monkeypatch.setattr(cli, "validate_execution", lambda **kwargs: pytest.fail("help reached execution"))
|
||||
with pytest.raises(SystemExit) as exc:
|
||||
cli.main(["--help"])
|
||||
assert exc.value.code == 0
|
||||
assert "--production-shadow" in capsys.readouterr().out
|
||||
|
||||
|
||||
def test_default_execution_does_not_authorize_production():
|
||||
with pytest.raises(RuntimeError, match="explicit --production-shadow"):
|
||||
cli.validate_execution(production_shadow=False, mode="shadow", database="clientflow")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["", "off", "compare", "authoritative"])
|
||||
def test_production_shadow_requires_exact_shadow_mode(mode):
|
||||
with pytest.raises(RuntimeError, match="requires BLIF_FLOW_V2_MODE=shadow"):
|
||||
cli.validate_execution(production_shadow=True, mode=mode, database="clientflow")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("database", ["clientflow_codex_test", "clientflow_codex_shadow", "other"])
|
||||
def test_production_shadow_requires_exact_production_database(database):
|
||||
with pytest.raises(RuntimeError, match="requires database 'clientflow'"):
|
||||
cli.validate_execution(production_shadow=True, mode="shadow", database=database)
|
||||
|
||||
|
||||
def test_production_shadow_explicit_proof_is_accepted():
|
||||
cli.validate_execution(production_shadow=True, mode="shadow", database="clientflow")
|
||||
|
||||
|
||||
def test_global_write_allowlist_remains_test_only():
|
||||
assert projection.WRITE_DATABASE_ALLOWLIST == frozenset({"clientflow_codex_test"})
|
||||
assert "clientflow" not in projection.WRITE_DATABASE_ALLOWLIST
|
||||
|
||||
|
||||
def test_factual_read_transaction_remains_read_only():
|
||||
source = getsource(simulator._load)
|
||||
assert 'conn.execute(text("BEGIN READ ONLY"))' in source
|
||||
assert 'conn.execute(text("ROLLBACK"))' in source
|
||||
|
||||
|
||||
def test_projection_writer_only_mutates_two_additive_tables():
|
||||
source = getsource(projection.rebuild_blif_flow_v2_projection).upper()
|
||||
assert projection.PROJECTION_WRITE_TABLES == {
|
||||
"opportunity_flow_state_v2", "opportunity_flow_transitions"}
|
||||
assert "INSERT INTO OPPORTUNITY_FLOW_TRANSITIONS" in source
|
||||
assert "INSERT INTO OPPORTUNITY_FLOW_STATE_V2" in source
|
||||
for forbidden in ("OPPORTUNITIES", "TASKS", "MESSAGES", "COMMUNICATIONS",
|
||||
"COMMERCIAL_DOCUMENTS", "CUSTOMERS", "PAYMENTS", "OPERATION_LINKS"):
|
||||
assert f"INSERT INTO {forbidden}" not in source
|
||||
assert f"UPDATE {forbidden}" not in source
|
||||
assert f"DELETE FROM {forbidden}" not in source
|
||||
|
||||
|
||||
def test_production_derivation_has_no_fixed_328_requirement(monkeypatch):
|
||||
captured = {}
|
||||
|
||||
def fake_collect(**kwargs):
|
||||
captured.update(kwargs)
|
||||
return {"opportunities": [{"opportunity_id": "one"}]}
|
||||
|
||||
monkeypatch.setattr(simulator, "collect", fake_collect)
|
||||
rows = projection._derive_all(
|
||||
expected_database="clientflow", expected_user="runtime-role",
|
||||
expected_opportunity_count=None, require_opportunities=True,
|
||||
)
|
||||
assert rows == [{"opportunity_id": "one"}]
|
||||
assert captured["expected_opportunity_count"] is None
|
||||
assert captured["require_opportunities"] is True
|
||||
|
||||
|
||||
def test_snapshot_expected_count_validation_is_still_available(monkeypatch):
|
||||
monkeypatch.setattr(simulator, "_load", lambda **kwargs: {"opportunities": [object()]})
|
||||
with pytest.raises(RuntimeError, match="expected 328 opportunities, found 1"):
|
||||
simulator.collect(expected_opportunity_count=328)
|
||||
|
||||
|
||||
def test_empty_production_universe_is_rejected(monkeypatch):
|
||||
monkeypatch.setattr(simulator, "collect", lambda **kwargs: {"opportunities": []})
|
||||
assert projection._derive_all(
|
||||
expected_database="clientflow", expected_user=None,
|
||||
expected_opportunity_count=None, require_opportunities=True,
|
||||
) == []
|
||||
with pytest.raises(RuntimeError, match="at least one opportunity"):
|
||||
projection.rebuild_blif_flow_v2_projection(
|
||||
mode="shadow", derived_rows=[], require_opportunities=True,
|
||||
)
|
||||
|
||||
|
||||
def test_existing_test_defaults_remain_guarded(monkeypatch):
|
||||
captured = {}
|
||||
monkeypatch.setattr(simulator, "collect", lambda **kwargs: captured.update(kwargs) or {"opportunities": []})
|
||||
projection._derive_all()
|
||||
assert captured["expected_database"] == "clientflow_codex_test"
|
||||
assert captured["expected_user"] == "clientflow_codex_test"
|
||||
assert captured["expected_opportunity_count"] == 328
|
||||
125
tests/test_blif_flow_v2_readonly_transaction_order.py
Normal file
125
tests/test_blif_flow_v2_readonly_transaction_order.py
Normal file
@@ -0,0 +1,125 @@
|
||||
from contextlib import nullcontext
|
||||
|
||||
import pytest
|
||||
|
||||
import scripts.simulate_blif_flow_v2 as simulator
|
||||
|
||||
|
||||
class _Result:
|
||||
def __init__(self, *, identity=None, rows=()):
|
||||
self._identity = identity
|
||||
self._rows = rows
|
||||
|
||||
def one(self):
|
||||
return self._identity
|
||||
|
||||
def mappings(self):
|
||||
return self._rows
|
||||
|
||||
|
||||
class _Connection:
|
||||
def __init__(self, database="clientflow", user="clientflow"):
|
||||
self.database = database
|
||||
self.user = user
|
||||
self.read_only = False
|
||||
self.statements = []
|
||||
self.identity_observations = []
|
||||
|
||||
def execution_options(self, **_kwargs):
|
||||
return self
|
||||
|
||||
def execute(self, statement):
|
||||
sql = " ".join(str(statement).split())
|
||||
self.statements.append(sql)
|
||||
upper = sql.upper()
|
||||
if upper == "BEGIN READ ONLY":
|
||||
self.read_only = True
|
||||
return _Result()
|
||||
if upper == "ROLLBACK":
|
||||
self.read_only = False
|
||||
return _Result()
|
||||
if upper.startswith(("INSERT ", "UPDATE ", "DELETE ")):
|
||||
if self.read_only:
|
||||
raise RuntimeError("cannot execute write in a read-only transaction")
|
||||
return _Result()
|
||||
if "CURRENT_DATABASE()" in upper:
|
||||
identity = (self.database, self.user, "on" if self.read_only else "off")
|
||||
self.identity_observations.append(identity)
|
||||
return _Result(identity=identity)
|
||||
return _Result(rows=())
|
||||
|
||||
|
||||
class _Engine:
|
||||
def __init__(self, connection):
|
||||
self.connection = connection
|
||||
|
||||
def connect(self):
|
||||
return nullcontext(self.connection)
|
||||
|
||||
|
||||
def _load(monkeypatch, connection, **kwargs):
|
||||
monkeypatch.setattr(simulator, "engine", _Engine(connection))
|
||||
return simulator._load(**kwargs)
|
||||
|
||||
|
||||
def test_strict_load_establishes_read_only_before_same_connection_identity_and_reads(monkeypatch):
|
||||
connection = _Connection()
|
||||
assert connection.read_only is False # normal fresh-connection state
|
||||
|
||||
result = _load(
|
||||
monkeypatch, connection, expected_database="clientflow",
|
||||
expected_user="clientflow", require_read_only=True,
|
||||
)
|
||||
|
||||
assert connection.statements[0] == "BEGIN READ ONLY"
|
||||
assert "CURRENT_DATABASE()" in connection.statements[1].upper()
|
||||
assert connection.identity_observations == [("clientflow", "clientflow", "on")]
|
||||
assert result["identity"]["transaction_read_only"] == "on"
|
||||
assert connection.statements[2].upper().startswith("SELECT O.*")
|
||||
assert connection.statements[-1] == "ROLLBACK"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("database", "user", "expected_database", "expected_user"),
|
||||
[
|
||||
("wrong", "clientflow", "clientflow", "clientflow"),
|
||||
("clientflow", "wrong", "clientflow", "clientflow"),
|
||||
],
|
||||
)
|
||||
def test_strict_load_refuses_wrong_identity_before_factual_reads(
|
||||
monkeypatch, database, user, expected_database, expected_user,
|
||||
):
|
||||
connection = _Connection(database=database, user=user)
|
||||
with pytest.raises(RuntimeError, match="unexpected database identity"):
|
||||
_load(
|
||||
monkeypatch, connection, expected_database=expected_database,
|
||||
expected_user=expected_user, require_read_only=True,
|
||||
)
|
||||
|
||||
assert connection.statements[0] == "BEGIN READ ONLY"
|
||||
assert len([sql for sql in connection.statements if sql.upper().startswith("SELECT")]) == 1
|
||||
assert connection.statements[-1] == "ROLLBACK"
|
||||
|
||||
|
||||
def test_strict_factual_transaction_rejects_writes():
|
||||
connection = _Connection()
|
||||
connection.execute("BEGIN READ ONLY")
|
||||
with pytest.raises(RuntimeError, match="read-only transaction"):
|
||||
connection.execute("UPDATE opportunities SET stage = 'forbidden'")
|
||||
assert connection.read_only is True
|
||||
|
||||
|
||||
def test_non_strict_load_preserves_identity_then_read_only_collection_order(monkeypatch):
|
||||
connection = _Connection(database="clientflow_codex_test", user="clientflow_codex_test")
|
||||
result = _load(
|
||||
monkeypatch, connection, expected_database="clientflow_codex_test",
|
||||
expected_user="clientflow_codex_test", require_read_only=False,
|
||||
)
|
||||
|
||||
assert "CURRENT_DATABASE()" in connection.statements[0].upper()
|
||||
assert connection.identity_observations == [
|
||||
("clientflow_codex_test", "clientflow_codex_test", "off")
|
||||
]
|
||||
assert connection.statements[1] == "BEGIN READ ONLY"
|
||||
assert result["identity"]["transaction_read_only"] == "off"
|
||||
assert connection.statements[-1] == "ROLLBACK"
|
||||
@@ -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