Compare commits
8 Commits
fix/legacy
...
7cc9fecbdc
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7cc9fecbdc | ||
|
|
e456cbc0e4 | ||
|
|
e3b8750ebb | ||
|
|
32e7773957 | ||
|
|
3b60c3fb70 | ||
|
|
acbbd1a84b | ||
|
|
e99da64d8b | ||
|
|
20bc91dec5 |
@@ -40,6 +40,16 @@ INTERNAL_ACTION_CODES = {
|
||||
|
||||
ACTION_CODES = TRIAGE_ACTION_CODES | INTERNAL_ACTION_CODES
|
||||
|
||||
# Persistence vocabulary for the shadow-only Flow v2 projection. These are not
|
||||
# added to TRIAGE_ACTION_CODES or ACTION_CODES, so no existing task/LLM/runtime
|
||||
# behavior changes.
|
||||
FLOW_V2_BUSINESS_ACTION_CODES = {
|
||||
"CREATE_PROFORMA",
|
||||
"CREATE_INVOICE",
|
||||
"VALIDATE_ODOO_ORDER",
|
||||
"COMPLETE_OPPORTUNITY",
|
||||
}
|
||||
|
||||
ACTION_MAP = {
|
||||
"CALL_CUSTOMER": {
|
||||
"route": "vendas",
|
||||
|
||||
@@ -1544,20 +1544,14 @@ def _parse_opportunity_dt(value: object):
|
||||
def _opportunity_lifecycle_state(opp: dict) -> str:
|
||||
state = str(opp.get("lifecycle_state") or "active").strip().lower() or "active"
|
||||
now = datetime.now(timezone.utc)
|
||||
pending_call = str(opp.get("pending_follow_up_action_code") or "").strip().upper() == "CALL_CUSTOMER"
|
||||
pending_call_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at"))
|
||||
if pending_call:
|
||||
return "follow_up_due" if pending_call_due and pending_call_due <= now else "scheduled_follow_up"
|
||||
# Compatibility only: a stale denormalized lifecycle marker cannot invent
|
||||
# a scheduled call when no active CALL_CUSTOMER task exists.
|
||||
if state == "scheduled_follow_up":
|
||||
pending_followup = str(opp.get("pending_follow_up_action_code") or "").strip().upper()
|
||||
pending_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at"))
|
||||
if pending_followup:
|
||||
return "follow_up_due" if pending_due and pending_due <= now else "scheduled_follow_up"
|
||||
# Compatibility-only timestamps/lifecycle markers cannot create work. An
|
||||
# active pending follow-up task and its due_at are the operational source.
|
||||
if state in {"scheduled_follow_up", "follow_up_due"}:
|
||||
return "active"
|
||||
nurture_until = _parse_opportunity_dt(opp.get("nurture_until"))
|
||||
next_follow_up = _parse_opportunity_dt(opp.get("next_follow_up_at"))
|
||||
if state == "nurture" and nurture_until and nurture_until <= now:
|
||||
return "follow_up_due"
|
||||
if state in {"awaiting_customer", "active", "scheduled_follow_up"} and next_follow_up and next_follow_up <= now:
|
||||
return "follow_up_due"
|
||||
return state
|
||||
|
||||
|
||||
|
||||
236
app/blif_flow_v2_projection_service.py
Normal file
236
app/blif_flow_v2_projection_service.py
Normal file
@@ -0,0 +1,236 @@
|
||||
"""Persistence scaffolding for the rebuildable BLIF Flow v2 projection.
|
||||
|
||||
The factual sources remain authoritative. This module writes only the additive
|
||||
projection/audit tables introduced by migration 011 and never updates stages or
|
||||
tasks.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import re
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Iterable
|
||||
|
||||
from sqlalchemy import text
|
||||
|
||||
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:
|
||||
if isinstance(value, datetime):
|
||||
return value.isoformat()
|
||||
if isinstance(value, dict):
|
||||
return {str(key): _jsonable(item) for key, item in value.items()}
|
||||
if isinstance(value, (list, tuple)):
|
||||
return [_jsonable(item) for item in value]
|
||||
return value
|
||||
|
||||
|
||||
def _evidence_refs(row: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
evidence = row.get("evidence") or {}
|
||||
refs: list[dict[str, Any]] = []
|
||||
for name in ("latest_relevant_inbound", "latest_relevant_outbound"):
|
||||
item = evidence.get(name)
|
||||
if item and item.get("id"):
|
||||
refs.append({"source": "message", "role": name, "id": item["id"], "at": item.get("at")})
|
||||
for name, source in (("proforma", "jasmin_proforma"), ("invoice", "jasmin_invoice"),
|
||||
("payment", "payment"), ("odoo", "odoo"),
|
||||
("reconciliation", "reconciliation")):
|
||||
for item in evidence.get(name) or []:
|
||||
if item.get("id"):
|
||||
refs.append({
|
||||
"source": source, "id": item["id"],
|
||||
"external_id": item.get("external_id"),
|
||||
"document_number": item.get("document_number"),
|
||||
"external_type": item.get("external_type"),
|
||||
"status": item.get("status"),
|
||||
})
|
||||
return refs
|
||||
|
||||
|
||||
def _projection_value(row: dict[str, Any], derived_at: datetime) -> dict[str, Any]:
|
||||
raw = row["raw_v2"]
|
||||
duplicate = row.get("canonical_process_id") != row.get("opportunity_id")
|
||||
evidence_refs = _evidence_refs(row)
|
||||
reason_code = str(raw.get("precedence") or "business_transition").upper()
|
||||
stable = {
|
||||
"opportunity_id": row["opportunity_id"],
|
||||
"material_process_key": row["material_process_key"],
|
||||
"canonical_opportunity_id": row["canonical_process_id"],
|
||||
"is_duplicate_representation": duplicate,
|
||||
"business_state": raw["business_state"],
|
||||
"business_next_action": raw.get("business_next_action"),
|
||||
"diagnostic_status": raw.get("diagnostic_status") or row.get("evidence", {}).get("diagnostic_status") or "clear",
|
||||
"confidence": raw.get("confidence") or "low",
|
||||
"reason_code": reason_code,
|
||||
"reason_text": raw.get("reason") or "",
|
||||
"evidence_refs": evidence_refs,
|
||||
"flow_version": FLOW_VERSION,
|
||||
}
|
||||
fingerprint = hashlib.sha256(
|
||||
json.dumps(_jsonable(stable), ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8")
|
||||
).hexdigest()
|
||||
return stable | {"source_fingerprint": fingerprint, "derived_at": derived_at}
|
||||
|
||||
|
||||
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=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"])
|
||||
|
||||
|
||||
def rebuild_blif_flow_v2_projection(
|
||||
*,
|
||||
mode: str | None = None,
|
||||
derived_rows: Iterable[dict[str, Any]] | None = None,
|
||||
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.
|
||||
|
||||
``off`` is a no-op. ``shadow`` and ``compare`` write only the same additive
|
||||
projection tables. Compare affects observation at the read boundary, never
|
||||
the persisted business rows. Authoritative remains deliberately disabled.
|
||||
"""
|
||||
selected_mode = str(mode or settings.blif_flow_v2_mode or "off").strip().lower()
|
||||
if selected_mode == "off":
|
||||
return {"mode": "off", "projection_count": 0, "transitions_written": 0, "disabled": True}
|
||||
if selected_mode not in {"shadow", "compare"}:
|
||||
raise RuntimeError(f"BLIF Flow v2 mode {selected_mode!r} is disabled; only off/shadow/compare are safe")
|
||||
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(
|
||||
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):
|
||||
raise RuntimeError("Flow v2 projection requires exactly one derived row per opportunity")
|
||||
|
||||
owns_connection = connection is None
|
||||
if connection is None:
|
||||
from app.db import engine
|
||||
conn = engine.connect().execution_options(isolation_level="AUTOCOMMIT")
|
||||
else:
|
||||
conn = connection
|
||||
try:
|
||||
identity = conn.execute(text(
|
||||
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
||||
)).one()
|
||||
if identity[0] not in allowed_databases:
|
||||
raise RuntimeError(f"refusing Flow v2 projection write to database {identity[0]!r}")
|
||||
if owns_connection:
|
||||
conn.exec_driver_sql("BEGIN READ WRITE")
|
||||
try:
|
||||
if conn.execute(text("SELECT current_setting('transaction_read_only')")).scalar_one() != "off":
|
||||
raise RuntimeError("Flow v2 projection rebuild requires an explicit READ WRITE transaction")
|
||||
conn.exec_driver_sql(f'SET LOCAL search_path TO "{target_schema}"')
|
||||
existing = {
|
||||
str(row["opportunity_id"]): dict(row)
|
||||
for row in conn.execute(text("""
|
||||
SELECT opportunity_id::text, business_state, source_fingerprint
|
||||
FROM opportunity_flow_state_v2
|
||||
""")).mappings()
|
||||
}
|
||||
transitions_written = 0
|
||||
for value in values:
|
||||
previous = existing.get(value["opportunity_id"])
|
||||
if previous is None or previous["business_state"] != value["business_state"]:
|
||||
result = conn.execute(text("""
|
||||
INSERT INTO opportunity_flow_transitions (
|
||||
opportunity_id, from_state, to_state, reason_code, reason_text,
|
||||
evidence_refs, flow_version, source_fingerprint
|
||||
) VALUES (
|
||||
CAST(:opportunity_id AS UUID), :from_state, :to_state, :reason_code,
|
||||
:reason_text, CAST(:evidence_refs AS JSONB), :flow_version, :source_fingerprint
|
||||
)
|
||||
ON CONFLICT (opportunity_id, from_state, to_state, flow_version, source_fingerprint)
|
||||
DO NOTHING
|
||||
"""), {
|
||||
**value, "from_state": previous["business_state"] if previous else None,
|
||||
"to_state": value["business_state"],
|
||||
"evidence_refs": json.dumps(_jsonable(value["evidence_refs"]), ensure_ascii=False),
|
||||
})
|
||||
transitions_written += result.rowcount
|
||||
conn.execute(text("""
|
||||
INSERT INTO opportunity_flow_state_v2 (
|
||||
opportunity_id, material_process_key, canonical_opportunity_id,
|
||||
is_duplicate_representation, business_state, business_next_action,
|
||||
diagnostic_status, confidence, reason_code, reason_text, evidence_refs,
|
||||
flow_version, source_fingerprint, derived_at, updated_at
|
||||
) VALUES (
|
||||
CAST(:opportunity_id AS UUID), :material_process_key,
|
||||
CAST(:canonical_opportunity_id AS UUID), :is_duplicate_representation,
|
||||
:business_state, :business_next_action, :diagnostic_status, :confidence,
|
||||
:reason_code, :reason_text, CAST(:evidence_refs_json AS JSONB), :flow_version,
|
||||
:source_fingerprint, :derived_at, now()
|
||||
)
|
||||
ON CONFLICT (opportunity_id) DO UPDATE SET
|
||||
material_process_key=EXCLUDED.material_process_key,
|
||||
canonical_opportunity_id=EXCLUDED.canonical_opportunity_id,
|
||||
is_duplicate_representation=EXCLUDED.is_duplicate_representation,
|
||||
business_state=EXCLUDED.business_state,
|
||||
business_next_action=EXCLUDED.business_next_action,
|
||||
diagnostic_status=EXCLUDED.diagnostic_status,
|
||||
confidence=EXCLUDED.confidence,
|
||||
reason_code=EXCLUDED.reason_code,
|
||||
reason_text=EXCLUDED.reason_text,
|
||||
evidence_refs=EXCLUDED.evidence_refs,
|
||||
flow_version=EXCLUDED.flow_version,
|
||||
source_fingerprint=EXCLUDED.source_fingerprint,
|
||||
derived_at=EXCLUDED.derived_at,
|
||||
updated_at=CASE
|
||||
WHEN opportunity_flow_state_v2.source_fingerprint IS DISTINCT FROM EXCLUDED.source_fingerprint
|
||||
THEN now() ELSE opportunity_flow_state_v2.updated_at END
|
||||
"""), value | {
|
||||
"evidence_refs_json": json.dumps(_jsonable(value["evidence_refs"]), ensure_ascii=False)
|
||||
})
|
||||
if owns_connection:
|
||||
conn.exec_driver_sql("COMMIT")
|
||||
except Exception:
|
||||
if owns_connection:
|
||||
conn.exec_driver_sql("ROLLBACK")
|
||||
raise
|
||||
finally:
|
||||
if owns_connection:
|
||||
conn.close()
|
||||
|
||||
return {
|
||||
"mode": selected_mode, "projection_count": len(values),
|
||||
"canonical_count": sum(not value["is_duplicate_representation"] for value in values),
|
||||
"duplicate_count": sum(value["is_duplicate_representation"] for value in values),
|
||||
"transitions_written": transitions_written, "disabled": False,
|
||||
}
|
||||
@@ -21,6 +21,10 @@ class Settings(BaseSettings):
|
||||
# TIMESTAMPTZ nem os índices usados pelo schema core.
|
||||
database_url: str
|
||||
clientflow_persist: bool = True
|
||||
# Flow v2 remains non-authoritative. Compare observes V1/V2 differences
|
||||
# while returning V1; authoritative is accepted by config only to fail
|
||||
# closed at the switching boundary.
|
||||
blif_flow_v2_mode: Literal["off", "shadow", "compare", "authoritative"] = "off"
|
||||
|
||||
clientflow_webhook_secret: str = ""
|
||||
clientflow_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy"
|
||||
|
||||
@@ -249,15 +249,7 @@ def build_opportunity_evidence(
|
||||
|
||||
quote = _find_doc(docs, QUOTE_KINDS)
|
||||
invoice = _find_doc(docs, INVOICE_KINDS)
|
||||
# A existência do documento comercial prova apenas que o orçamento foi
|
||||
# criado/associado. O envio ao cliente exige evidência própria.
|
||||
#
|
||||
# Compatibilidade histórica: QUOTE_SENT é também uma afirmação canónica
|
||||
# explícita de que o orçamento já foi enviado.
|
||||
quote_sent = (
|
||||
stage == "QUOTE_SENT"
|
||||
or _completed_send_quote_task_evidence(tasks)
|
||||
)
|
||||
quote_sent = bool(quote) or _completed_send_quote_task_evidence(tasks)
|
||||
|
||||
pending_task = None
|
||||
invalid_payment_task = None
|
||||
|
||||
159
app/domain/opportunity_flow/repair.py
Normal file
159
app/domain/opportunity_flow/repair.py
Normal file
@@ -0,0 +1,159 @@
|
||||
"""Pure BLIF Flow v2 historical repair classification and simulation helpers.
|
||||
|
||||
The functions in this module never access the database and never mutate business
|
||||
state. They deliberately treat tasks as obligations/history, not factual proof.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import asdict, dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
|
||||
TASK_CLASSIFICATIONS = {
|
||||
"VALID_CURRENT", "SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE",
|
||||
"PREMATURE", "AMBIGUOUS",
|
||||
}
|
||||
|
||||
FOLLOWUP_ACTIONS = {
|
||||
"FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA",
|
||||
"FOLLOW_UP_PAYMENT", "CALL_CUSTOMER",
|
||||
}
|
||||
|
||||
STAGE_RANK = {
|
||||
"INQUIRY": 0, "AWAITING_CUSTOMER": 1, "PROFORMA_REQUIRED": 2,
|
||||
"PROFORMA_CREATED": 3, "AWAITING_PAYMENT": 4, "INVOICE_REQUIRED": 5,
|
||||
"INVOICE_CREATED": 6, "ODOO_ORDER_REQUIRED": 6,
|
||||
"ODOO_ORDER_CREATED": 7, "ODOO_ORDER_VALIDATED": 8, "COMPLETED": 9,
|
||||
}
|
||||
|
||||
ACTION_REQUIRED_RANK = {
|
||||
"SEND_INFO": 0, "SEND_QUOTE": 0, "CREATE_PROFORMA": 2,
|
||||
"SEND_PROFORMA": 3, "FOLLOW_UP_PROFORMA": 4, "CONFIRM_PAYMENT": 4,
|
||||
"FOLLOW_UP_PAYMENT": 4, "CREATE_INVOICE": 5, "SEND_INVOICE": 6,
|
||||
"PREPARE_ORDER": 6, "VALIDATE_ODOO_ORDER": 7,
|
||||
"COMPLETE_OPPORTUNITY": 8,
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TaskRepairContext:
|
||||
task_id: str
|
||||
opportunity_id: str | None
|
||||
action_code: str
|
||||
created_at: datetime | None = None
|
||||
due_at: datetime | None = None
|
||||
business_state: str | None = None
|
||||
business_next_action: str | None = None
|
||||
material_process_key: str | None = None
|
||||
is_duplicate_representation: bool = False
|
||||
canonical_opportunity_id: str | None = None
|
||||
later_inbound_event: dict[str, Any] | None = None
|
||||
later_outbound_event: dict[str, Any] | None = None
|
||||
proforma_exists: bool = False
|
||||
proforma_sent: bool = False
|
||||
payment_confirmed: bool = False
|
||||
invoice_exists: bool = False
|
||||
invoice_sent: bool = False
|
||||
odoo_order_exists: bool = False
|
||||
odoo_order_validated: bool = False
|
||||
terminal: bool = False
|
||||
same_obligation_task_id: str | None = None
|
||||
evidence_refs: tuple[dict[str, Any], ...] = field(default_factory=tuple)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TaskRepairDecision:
|
||||
classification: str
|
||||
resolution_code: str | None
|
||||
reason: str
|
||||
confidence: str
|
||||
safety_tier: str
|
||||
auto_repair_safe: bool
|
||||
human_review_required: bool
|
||||
resolved_by_event_id: str | None = None
|
||||
superseded_by_task_id: str | None = None
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
return asdict(self)
|
||||
|
||||
|
||||
def _decision(classification: str, reason: str, *, event: dict[str, Any] | None = None,
|
||||
superseded_by: str | None = None, tier: str = "HIGH") -> TaskRepairDecision:
|
||||
if classification not in TASK_CLASSIFICATIONS:
|
||||
raise ValueError(classification)
|
||||
auto = tier == "HIGH" and classification != "VALID_CURRENT"
|
||||
return TaskRepairDecision(
|
||||
classification, None if classification == "VALID_CURRENT" else classification,
|
||||
reason, "high" if tier == "HIGH" else "medium" if tier == "MEDIUM" else "low",
|
||||
tier, auto, tier != "HIGH",
|
||||
(event or {}).get("opportunity_event_id"), superseded_by,
|
||||
)
|
||||
|
||||
|
||||
def classify_pending_task(context: TaskRepairContext) -> TaskRepairDecision:
|
||||
"""Classify one pending task using facts available at the audit instant."""
|
||||
action = context.action_code.upper()
|
||||
state = (context.business_state or "").upper()
|
||||
|
||||
if context.is_duplicate_representation:
|
||||
return _decision(
|
||||
"DUPLICATE",
|
||||
f"Obligation belongs only to a duplicate representation of canonical process {context.canonical_opportunity_id}.",
|
||||
)
|
||||
if context.same_obligation_task_id:
|
||||
return _decision(
|
||||
"DUPLICATE", "The same material obligation has another canonical pending task.",
|
||||
superseded_by=context.same_obligation_task_id,
|
||||
)
|
||||
|
||||
if action in FOLLOWUP_ACTIONS:
|
||||
event = context.later_inbound_event
|
||||
if action == "FOLLOW_UP_PAYMENT" and context.payment_confirmed:
|
||||
return _decision("SATISFIED_BY_EVENT", "Confirmed payment fact satisfies the payment follow-up.")
|
||||
if event:
|
||||
return _decision("SATISFIED_BY_EVENT", "A later inbound customer event satisfies the follow-up.", event=event)
|
||||
return _decision("VALID_CURRENT", "No later satisfying event exists; age or overdue status alone never closes a follow-up.")
|
||||
|
||||
satisfied = {
|
||||
"SEND_INFO": context.later_outbound_event,
|
||||
"SEND_QUOTE": context.later_outbound_event,
|
||||
"CREATE_PROFORMA": context.proforma_exists,
|
||||
"SEND_PROFORMA": context.proforma_sent,
|
||||
"CONFIRM_PAYMENT": context.payment_confirmed,
|
||||
"CREATE_INVOICE": context.invoice_exists,
|
||||
"SEND_INVOICE": context.invoice_sent,
|
||||
"PREPARE_ORDER": context.odoo_order_exists,
|
||||
"VALIDATE_ODOO_ORDER": context.odoo_order_validated,
|
||||
"COMPLETE_OPPORTUNITY": context.terminal,
|
||||
}.get(action, False)
|
||||
if satisfied:
|
||||
event = satisfied if isinstance(satisfied, dict) else None
|
||||
return _decision("SATISFIED_BY_EVENT", f"Later factual evidence satisfies {action}.", event=event)
|
||||
|
||||
required = ACTION_REQUIRED_RANK.get(action)
|
||||
rank = STAGE_RANK.get(state)
|
||||
if required is not None and rank is not None:
|
||||
if rank > required:
|
||||
return _decision("SUPERSEDED", f"Factual process advanced to {state}, beyond the {action} obligation.")
|
||||
if rank < required:
|
||||
return _decision("PREMATURE", f"{action} requires prerequisites not present in factual state {state}.")
|
||||
|
||||
if action == "SEND_PROFORMA" and not context.proforma_exists:
|
||||
return _decision("PREMATURE", "No structured current proforma exists; a send task is not document evidence.")
|
||||
if action == "SEND_INVOICE" and not context.invoice_exists:
|
||||
return _decision("PREMATURE", "No structured invoice exists and protected proforma/payment prerequisites are absent.")
|
||||
if action == context.business_next_action:
|
||||
return _decision("VALID_CURRENT", "Task matches the current factual Flow v2 obligation.")
|
||||
if action in {"SUPPORT", "MARK_NO_INTEREST", "REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS"}:
|
||||
return _decision("VALID_CURRENT", "Historically valid operator obligation has no factual evidence of satisfaction or supersession.")
|
||||
if not context.opportunity_id:
|
||||
return _decision("VALID_CURRENT", "Standalone obligation is outside opportunity business transitions and is preserved.")
|
||||
return _decision("AMBIGUOUS", "Available factual evidence does not deterministically establish validity or safe removal.", tier="LOW")
|
||||
|
||||
|
||||
def simulate_high_repairs(task_rows: list[dict[str, Any]]) -> dict[str, Any]:
|
||||
removed = [row for row in task_rows if row["safety_tier"] == "HIGH" and row["auto_repair_safe"]]
|
||||
counts = {name: sum(row["classification"] == name for row in removed) for name in TASK_CLASSIFICATIONS}
|
||||
return {"pending_before": len(task_rows), "pending_after": len(task_rows) - len(removed),
|
||||
"removed": removed, "removed_by_classification": counts}
|
||||
@@ -8,7 +8,6 @@ from .types import (
|
||||
ACTION_CLOSE_OPPORTUNITY,
|
||||
ACTION_CONFIRM_PAYMENT,
|
||||
ACTION_CREATE_QUOTE,
|
||||
ACTION_SEND_QUOTE,
|
||||
ACTION_FOLLOW_UP,
|
||||
ACTION_FOLLOW_UP_PAYMENT,
|
||||
ACTION_NO_ACTION,
|
||||
@@ -124,26 +123,16 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
|
||||
return OpportunityDecision(next_action, "Conflito fiscal/NIF bloqueia ações financeiras.", blocked_actions=blocked_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
if not e.has_fiscal_customer:
|
||||
# A ausência de cliente fiscal é prontidão operacional, não intenção
|
||||
# comercial. Não deve substituir a próxima ação da oportunidade.
|
||||
# O bloqueio fiscal é aplicado apenas mais abaixo quando uma transição
|
||||
# concreta (por exemplo faturação após pagamento confirmado) exige
|
||||
# efetivamente os dados fiscais.
|
||||
warnings.append(
|
||||
"Cliente fiscal ainda não associado; validar apenas quando uma "
|
||||
"operação documental atual exigir dados fiscais."
|
||||
)
|
||||
blocked_actions.extend(_blocked(profile, code, "cliente fiscal por associar") for code in SENSITIVE_DOCUMENT_ACTIONS)
|
||||
next_action = _action(profile, ACTION_VALIDATE_FISCAL_CUSTOMER, "Associar/validar cliente fiscal antes de documentos oficiais.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#cliente" if e.opportunity_id else None)
|
||||
return OpportunityDecision(next_action, "Cliente fiscal ainda não associado.", blocked_actions=blocked_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
if e.has_reconciliation_candidate:
|
||||
next_action = _action(profile, ACTION_RECONCILE_DOCUMENTS, f"Confirmar evidência encontrada: {e.reconciliation_label or 'documento/candidato'}.", priority="alta", target_url="/reconciliation")
|
||||
return OpportunityDecision(next_action, "Há evidência de reconciliação por validar.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
if not e.has_quote and not e.has_invoice:
|
||||
# Não promover automaticamente qualquer oportunidade para orçamento.
|
||||
# A criação/reconciliação de orçamento só é trabalho atual quando o
|
||||
# estágio comercial demonstra que o cliente pediu ou já recebeu um
|
||||
# orçamento.
|
||||
if e.stage == "QUOTE_SENT" and e.quote_sent:
|
||||
if e.quote_sent:
|
||||
next_action = _action(
|
||||
profile,
|
||||
ACTION_RECONCILE_DOCUMENTS,
|
||||
@@ -162,26 +151,9 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
|
||||
profile_name=profile.name,
|
||||
decision_version=profile.version,
|
||||
)
|
||||
|
||||
if e.stage == "QUOTE_REQUESTED":
|
||||
next_action = _action(
|
||||
profile,
|
||||
ACTION_CREATE_QUOTE,
|
||||
"Criar/enviar orçamento solicitado pelo cliente.",
|
||||
target_url=f"/opportunities/{e.opportunity_id}#documentos" if e.opportunity_id else None,
|
||||
)
|
||||
next_action = _action(profile, ACTION_CREATE_QUOTE, "Criar/enviar orçamento antes de pedir pagamento ou emitir fatura.", target_url=f"/opportunities/{e.opportunity_id}#documentos" if e.opportunity_id else None)
|
||||
available_actions.append(next_action)
|
||||
return OpportunityDecision(
|
||||
next_action,
|
||||
"Existe pedido de orçamento e ainda não há documento comercial associado.",
|
||||
available_actions=available_actions,
|
||||
warnings=warnings,
|
||||
commercial_stage=COMMERCIAL_STAGE_REVIEW,
|
||||
financial_state="no_document",
|
||||
physical_state=_physical_state(e),
|
||||
profile_name=profile.name,
|
||||
decision_version=profile.version,
|
||||
)
|
||||
return OpportunityDecision(next_action, "Ainda não há orçamento/fatura associado.", available_actions=available_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_REVIEW, financial_state="no_document", physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
if e.payment_terms == PAYMENT_AFTER_DELIVERY:
|
||||
if e.stage == "SHIPMENT_CREATED" and not e.payment_confirmed:
|
||||
@@ -254,9 +226,7 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
|
||||
next_action = _action(profile, ACTION_PREPARE_ORDER, "Pagamento após entrega: criar/associar venda Odoo e avançar preparação sem exigir pagamento confirmado.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None)
|
||||
return OpportunityDecision(next_action, "Condição pós-entrega permite avançar Odoo/preparação sem pagamento prévio.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
if e.odoo_ready and not e.order_shipped:
|
||||
if not e.has_invoice and (
|
||||
not e.has_fiscal_customer or not e.fiscal_data_complete
|
||||
):
|
||||
if not e.has_invoice and e.has_fiscal_customer and not e.fiscal_data_complete:
|
||||
blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos"))
|
||||
next_action = _action(
|
||||
profile,
|
||||
@@ -295,76 +265,21 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
|
||||
next_action = _action(profile, ACTION_WAIT_PRODUCTION, "Pagamento após entrega: venda Odoo criada; aguardar WH/OUT ficar pronto/concluído antes de emitir fatura.", priority="normal", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None)
|
||||
return OpportunityDecision(next_action, "Aguardar estado da encomenda/WH-OUT no Odoo; ordens de fabrico são apenas detalhe técnico.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
# Um orçamento criado/associado ainda não significa orçamento enviado.
|
||||
# O envio ao cliente é uma obrigação humana/documental própria e deve
|
||||
# acontecer antes de qualquer etapa de pagamento.
|
||||
if (
|
||||
e.has_quote
|
||||
and not e.quote_sent
|
||||
and not e.has_invoice
|
||||
and not e.payment_confirmed
|
||||
):
|
||||
next_action = _action(
|
||||
profile,
|
||||
ACTION_SEND_QUOTE,
|
||||
f"Orçamento {e.quote_number or ''} criado/associado. Enviar o documento ao cliente.",
|
||||
priority="alta",
|
||||
target_url=(
|
||||
f"/tasks/{e.pending_task_id}"
|
||||
if e.pending_task_id and e.pending_task_action_code == ACTION_SEND_QUOTE
|
||||
else f"/opportunities/{e.opportunity_id}#documentos"
|
||||
if e.opportunity_id
|
||||
else None
|
||||
),
|
||||
document_id=e.quote_id,
|
||||
document_number=e.quote_number,
|
||||
)
|
||||
available_actions.append(next_action)
|
||||
return OpportunityDecision(
|
||||
next_action,
|
||||
"O orçamento existe, mas ainda não há evidência de que tenha sido enviado ao cliente.",
|
||||
available_actions=available_actions,
|
||||
warnings=warnings,
|
||||
commercial_stage=COMMERCIAL_STAGE_REVIEW,
|
||||
financial_state=_financial_state(e),
|
||||
physical_state=_physical_state(e),
|
||||
profile_name=profile.name,
|
||||
decision_version=profile.version,
|
||||
)
|
||||
|
||||
# Default/BLIF normal sequence: budget document, payment, invoice, then preparation/shipping.
|
||||
if (
|
||||
e.payment_terms in {PAYMENT_BEFORE_SHIPPING, "", "undefined", "agreement"}
|
||||
and e.has_quote
|
||||
and e.quote_sent
|
||||
and not e.has_invoice
|
||||
and not e.payment_confirmed
|
||||
):
|
||||
if e.payment_terms in {PAYMENT_BEFORE_SHIPPING, "", "undefined", "agreement"} and e.has_quote and not e.payment_confirmed:
|
||||
next_action = _action(
|
||||
profile,
|
||||
ACTION_NO_ACTION,
|
||||
f"Orçamento {e.quote_number or ''} enviado. Aguardar decisão do cliente ou evidência de pagamento.",
|
||||
force_label="Aguardar cliente / pagamento",
|
||||
ACTION_CONFIRM_PAYMENT,
|
||||
f"Orçamento {e.quote_number or ''} associado. Confirmar pagamento antes de emitir fatura.",
|
||||
priority="alta",
|
||||
target_url=f"/opportunities/{e.opportunity_id}#operacao" if e.opportunity_id else None,
|
||||
document_id=e.quote_id,
|
||||
document_number=e.quote_number,
|
||||
)
|
||||
return OpportunityDecision(
|
||||
next_action,
|
||||
"O orçamento foi enviado e ainda não existe evidência de pagamento que exija validação. "
|
||||
"Aguardar o cliente; o follow-up comercial assume quando ficar devido.",
|
||||
warnings=warnings,
|
||||
commercial_stage=COMMERCIAL_STAGE_WAITING_PAYMENT,
|
||||
financial_state=_financial_state(e),
|
||||
physical_state=_physical_state(e),
|
||||
profile_name=profile.name,
|
||||
decision_version=profile.version,
|
||||
)
|
||||
available_actions.append(next_action)
|
||||
return OpportunityDecision(next_action, "Fluxo normal BLIF exige pagamento confirmado depois do orçamento e antes da fatura.", available_actions=available_actions, warnings=warnings, commercial_stage=COMMERCIAL_STAGE_WAITING_PAYMENT, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
|
||||
if e.payment_confirmed and not e.has_invoice and (
|
||||
not e.has_fiscal_customer or not e.fiscal_data_complete
|
||||
):
|
||||
if e.payment_confirmed and not e.has_invoice and e.has_fiscal_customer and not e.fiscal_data_complete:
|
||||
blocked_actions.append(_blocked(profile, ACTION_SEND_INVOICE, "dados fiscais incompletos"))
|
||||
next_action = _action(
|
||||
profile,
|
||||
@@ -464,38 +379,5 @@ def decide_blif_next_action(e: OpportunityEvidence, profile: CompanyWorkflowProf
|
||||
next_action = _action(profile, ACTION_PREPARE_ORDER, f"Fatura {e.invoice_number or ''} e pagamento confirmados. Criar/validar venda Odoo e preparação.", priority="alta", target_url=f"/opportunities/{e.opportunity_id}#odoo" if e.opportunity_id else None)
|
||||
return OpportunityDecision(next_action, "Fatura e pagamento OK; falta validar execução/Odoo.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_IN_EXECUTION, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
if e.stage in {"NEW_LEAD", "INFO_REQUESTED", "INFO_SENT"}:
|
||||
next_action = _action(
|
||||
profile,
|
||||
ACTION_NO_ACTION,
|
||||
"Sem transição documental atual. Manter o estágio comercial e aguardar a próxima obrigação operacional real.",
|
||||
target_url=f"/opportunities/{e.opportunity_id}" if e.opportunity_id else None,
|
||||
)
|
||||
return OpportunityDecision(
|
||||
next_action,
|
||||
"O estágio comercial atual não exige orçamento, faturação ou follow-up imediato gerado pelo motor central.",
|
||||
warnings=warnings,
|
||||
commercial_stage=e.stage,
|
||||
financial_state=_financial_state(e),
|
||||
physical_state=_physical_state(e),
|
||||
profile_name=profile.name,
|
||||
decision_version=profile.version,
|
||||
)
|
||||
|
||||
next_action = _action(
|
||||
profile,
|
||||
ACTION_FOLLOW_UP,
|
||||
"Rever tarefas, documentos e próximos contactos.",
|
||||
priority="baixa",
|
||||
target_url=f"/opportunities/{e.opportunity_id}" if e.opportunity_id else None,
|
||||
)
|
||||
return OpportunityDecision(
|
||||
next_action,
|
||||
"Sem regra específica aplicável; manter em acompanhamento.",
|
||||
warnings=warnings,
|
||||
commercial_stage=e.stage or COMMERCIAL_STAGE_QUOTE_SENT,
|
||||
financial_state=_financial_state(e),
|
||||
physical_state=_physical_state(e),
|
||||
profile_name=profile.name,
|
||||
decision_version=profile.version,
|
||||
)
|
||||
next_action = _action(profile, ACTION_FOLLOW_UP, "Rever tarefas, documentos e próximos contactos.", priority="baixa", target_url=f"/opportunities/{e.opportunity_id}" if e.opportunity_id else None)
|
||||
return OpportunityDecision(next_action, "Sem regra específica aplicável; manter em acompanhamento.", warnings=warnings, commercial_stage=COMMERCIAL_STAGE_QUOTE_SENT, financial_state=_financial_state(e), physical_state=_physical_state(e), profile_name=profile.name, decision_version=profile.version)
|
||||
|
||||
@@ -21,8 +21,7 @@ ACTION_NO_ACTION = "NO_ACTION"
|
||||
ACTION_OPEN_TASK = "OPEN_TASK"
|
||||
ACTION_VALIDATE_FISCAL_CUSTOMER = "VALIDATE_FISCAL_CUSTOMER"
|
||||
ACTION_RECONCILE_DOCUMENTS = "RECONCILE_DOCUMENTS"
|
||||
ACTION_CREATE_QUOTE = "CREATE_QUOTE"
|
||||
ACTION_SEND_QUOTE = "SEND_QUOTE"
|
||||
ACTION_CREATE_QUOTE = "CREATE_JASMIN_QUOTE"
|
||||
ACTION_CONFIRM_PAYMENT = "CONFIRM_PAYMENT"
|
||||
ACTION_SEND_INVOICE = "SEND_INVOICE"
|
||||
ACTION_CONFIRM_ORDER = "CONFIRM_ORDER"
|
||||
|
||||
273
app/domain/opportunity_flow/v2.py
Normal file
273
app/domain/opportunity_flow/v2.py
Normal file
@@ -0,0 +1,273 @@
|
||||
"""Pure, shadow-only BLIF Flow v2 business-state projection.
|
||||
|
||||
This module has no database or V1 dependencies. In particular, task fields are
|
||||
kept only for audit output and never establish document or payment facts.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import asdict, dataclass, field
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class BusinessFacts:
|
||||
opportunity_id: str = ""
|
||||
terminal: bool = False
|
||||
explicitly_lost: bool = False
|
||||
exception: bool = False
|
||||
review_required: bool = False
|
||||
fiscal_blocked: bool = False
|
||||
document_reconciliation_required: bool = False
|
||||
|
||||
customer_request: bool = False
|
||||
request_kind: str = "info" # info | quote
|
||||
latest_relevant_inbound_at: datetime | None = None
|
||||
latest_relevant_outbound_at: datetime | None = None
|
||||
info_or_offer_sent: bool = False
|
||||
order_intent: bool = False
|
||||
order_intent_at: datetime | None = None
|
||||
fiscal_identity_evidence: bool = False
|
||||
|
||||
proforma_exists: bool = False
|
||||
proforma_sent: bool = False
|
||||
proforma_created_at: datetime | None = None
|
||||
proforma_sent_at: datetime | None = None
|
||||
potential_payment_evidence: bool = False
|
||||
payment_confirmed: bool = False
|
||||
payment_confirmed_at: datetime | None = None
|
||||
invoice_exists: bool = False
|
||||
invoice_created_at: datetime | None = None
|
||||
odoo_order_exists: bool = False
|
||||
odoo_order_validated: bool = False
|
||||
fulfillment_complete: bool = False
|
||||
|
||||
material_order_change: bool = False
|
||||
material_order_change_at: datetime | None = None
|
||||
later_customer_inbound_satisfies_followup: bool = False
|
||||
blockers: tuple[str, ...] = ()
|
||||
audit_task_codes: tuple[str, ...] = field(default=(), compare=False)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class FlowV2Decision:
|
||||
business_state: str
|
||||
next_action: str | None
|
||||
operational_queue: str
|
||||
reason: str
|
||||
confidence: str = "high"
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
return asdict(self)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class EffectiveOperationalDecision:
|
||||
business_state: str
|
||||
business_next_action: str | None
|
||||
effective_operational_action: str | None
|
||||
effective_operational_queue: str
|
||||
reason: str
|
||||
confidence: str
|
||||
precedence: str
|
||||
diagnostic_status: str = "clear"
|
||||
legacy_preserved_action: bool = False
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
return asdict(self)
|
||||
|
||||
|
||||
def derive_business_facts(**evidence: Any) -> BusinessFacts:
|
||||
"""Normalize factual adapter output without inferring facts from tasks."""
|
||||
allowed = BusinessFacts.__dataclass_fields__
|
||||
values = {key: value for key, value in evidence.items() if key in allowed}
|
||||
for key in ("blockers", "audit_task_codes"):
|
||||
if key in values and not isinstance(values[key], tuple):
|
||||
values[key] = tuple(values[key] or ())
|
||||
return BusinessFacts(**values)
|
||||
|
||||
|
||||
def _after(left: datetime | None, right: datetime | None) -> bool:
|
||||
return bool(left and right and left > right)
|
||||
|
||||
|
||||
def derive_business_state(facts: BusinessFacts) -> str:
|
||||
"""Derive the current state from strongest present-tense facts."""
|
||||
if facts.exception:
|
||||
return "EXCEPTION"
|
||||
if facts.explicitly_lost:
|
||||
return "LOST"
|
||||
if facts.review_required:
|
||||
return "REVIEW_REQUIRED"
|
||||
if facts.fiscal_blocked:
|
||||
return "FISCAL_BLOCKED"
|
||||
if facts.document_reconciliation_required:
|
||||
return "DOCUMENT_RECONCILIATION_REQUIRED"
|
||||
|
||||
change_after_payment = facts.material_order_change and (
|
||||
facts.payment_confirmed
|
||||
or facts.invoice_exists
|
||||
or _after(facts.material_order_change_at, facts.payment_confirmed_at)
|
||||
)
|
||||
if change_after_payment:
|
||||
return "REVIEW_REQUIRED"
|
||||
# An invoice without confirmed payment contradicts BLIF's normal protected
|
||||
# sequence. Do not silently skip payment or invent a correction flow.
|
||||
if facts.invoice_exists and not facts.payment_confirmed:
|
||||
return "REVIEW_REQUIRED"
|
||||
if facts.odoo_order_exists and (not facts.payment_confirmed or not facts.invoice_exists):
|
||||
return "REVIEW_REQUIRED"
|
||||
if facts.terminal and facts.payment_confirmed and facts.invoice_exists and facts.odoo_order_validated:
|
||||
return "COMPLETED"
|
||||
if facts.fulfillment_complete and facts.payment_confirmed and facts.invoice_exists:
|
||||
return "COMPLETED"
|
||||
if facts.odoo_order_validated:
|
||||
return "ODOO_ORDER_VALIDATED"
|
||||
if facts.odoo_order_exists:
|
||||
return "ODOO_ORDER_CREATED"
|
||||
if facts.invoice_exists:
|
||||
return "INVOICE_CREATED"
|
||||
if facts.payment_confirmed:
|
||||
return "INVOICE_REQUIRED"
|
||||
|
||||
change_invalidates_proforma = facts.material_order_change and (
|
||||
not facts.material_order_change_at
|
||||
or not facts.proforma_created_at
|
||||
or _after(facts.material_order_change_at, facts.proforma_created_at)
|
||||
)
|
||||
if facts.order_intent and (not facts.proforma_exists or change_invalidates_proforma):
|
||||
return "PROFORMA_REQUIRED"
|
||||
if facts.proforma_exists:
|
||||
return "PROFORMA_SENT" if facts.proforma_sent else "PROFORMA_CREATED"
|
||||
if facts.order_intent:
|
||||
return "ORDER_INTENT"
|
||||
if facts.info_or_offer_sent and not _after(
|
||||
facts.latest_relevant_inbound_at, facts.latest_relevant_outbound_at
|
||||
):
|
||||
return "AWAITING_CUSTOMER"
|
||||
return "INQUIRY"
|
||||
|
||||
|
||||
def derive_next_action(facts: BusinessFacts, state: str | None = None) -> FlowV2Decision:
|
||||
state = state or derive_business_state(facts)
|
||||
if state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED", "EXCEPTION"}:
|
||||
action = {
|
||||
"REVIEW_REQUIRED": "REVIEW_REQUIRED",
|
||||
"FISCAL_BLOCKED": "VALIDATE_FISCAL_CUSTOMER",
|
||||
"DOCUMENT_RECONCILIATION_REQUIRED": "RECONCILE_DOCUMENTS",
|
||||
"EXCEPTION": "REVIEW_EXCEPTION",
|
||||
}[state]
|
||||
return FlowV2Decision(state, action, "review" if state != "EXCEPTION" else "exception",
|
||||
"; ".join(facts.blockers) or f"{state} requires operator review.", "medium")
|
||||
if state in {"LOST", "NO_INTEREST", "COMPLETED"}:
|
||||
return FlowV2Decision(state, None, "not_current", "The factual process is terminal.")
|
||||
if state == "INQUIRY":
|
||||
if not facts.customer_request:
|
||||
return FlowV2Decision(state, None, "not_current", "No current unanswered customer request is evidenced.", "low")
|
||||
action = "SEND_QUOTE" if facts.request_kind == "quote" else "SEND_INFO"
|
||||
return FlowV2Decision(state, action, "do_now", "Customer request has no later relevant outbound response.", "medium")
|
||||
if state == "AWAITING_CUSTOMER":
|
||||
return FlowV2Decision(state, None, "waiting", "Information or offer was sent; awaiting a later customer decision.")
|
||||
if state in {"ORDER_INTENT", "PROFORMA_REQUIRED"}:
|
||||
return FlowV2Decision("PROFORMA_REQUIRED", "CREATE_PROFORMA", "do_now",
|
||||
"Customer order intent exists and no current structured proforma exists.")
|
||||
if state == "PROFORMA_CREATED":
|
||||
return FlowV2Decision(state, "SEND_PROFORMA", "do_now", "A current structured proforma exists but has no factual sent evidence.")
|
||||
if state == "PROFORMA_SENT":
|
||||
if facts.potential_payment_evidence:
|
||||
return FlowV2Decision("AWAITING_PAYMENT", "CONFIRM_PAYMENT", "do_now",
|
||||
"Potential payment evidence requires operator confirmation.", "medium")
|
||||
return FlowV2Decision("AWAITING_PAYMENT", None, "waiting", "The current proforma was sent and payment is not confirmed.")
|
||||
if state in {"PAYMENT_CONFIRMED", "INVOICE_REQUIRED"}:
|
||||
return FlowV2Decision("INVOICE_REQUIRED", "CREATE_INVOICE", "do_now", "Payment is confirmed and no structured invoice exists.")
|
||||
if state == "INVOICE_CREATED":
|
||||
return FlowV2Decision("ODOO_ORDER_REQUIRED", "PREPARE_ORDER", "do_now", "Invoice exists and no Odoo sale order exists.")
|
||||
if state == "ODOO_ORDER_CREATED":
|
||||
return FlowV2Decision(state, "VALIDATE_ODOO_ORDER", "do_now", "Odoo sale order exists but is not validated.")
|
||||
if state == "ODOO_ORDER_VALIDATED":
|
||||
return FlowV2Decision(state, "COMPLETE_OPPORTUNITY", "do_now", "The validated Odoo order is ready for opportunity completion.")
|
||||
return FlowV2Decision("REVIEW_REQUIRED", "REVIEW_REQUIRED", "review", f"No safe Flow v2 rule for {state}.", "low")
|
||||
|
||||
|
||||
def derive_v2_operational_queue(facts: BusinessFacts) -> FlowV2Decision:
|
||||
return derive_next_action(facts, derive_business_state(facts))
|
||||
|
||||
|
||||
def derive_effective_operational_action(
|
||||
business: FlowV2Decision,
|
||||
*,
|
||||
integration_exception: bool = False,
|
||||
scheduled_call_current: bool = False,
|
||||
due_followup_action: str | None = None,
|
||||
future_followup_action: str | None = None,
|
||||
fiscal_complete: bool = True,
|
||||
fiscal_required: bool = False,
|
||||
reconciliation_blocking: bool = False,
|
||||
diagnostic_status: str = "clear",
|
||||
) -> EffectiveOperationalDecision:
|
||||
"""Apply operational prerequisites without changing the business state."""
|
||||
action, queue, reason, precedence = (
|
||||
business.next_action, business.operational_queue, business.reason, "business_transition"
|
||||
)
|
||||
if integration_exception:
|
||||
action, queue, reason, precedence = "REVIEW_EXCEPTION", "exception", "An integration failure blocks current work.", "integration_exception"
|
||||
elif scheduled_call_current:
|
||||
action, queue, reason, precedence = "CALL_CUSTOMER", "do_now", "An explicit scheduled call is currently due.", "scheduled_call"
|
||||
elif reconciliation_blocking:
|
||||
action, queue, reason, precedence = "RECONCILE_DOCUMENTS", "review", "A real formal document requires current association/reconciliation.", "document_prerequisite"
|
||||
elif fiscal_required and not fiscal_complete:
|
||||
action, queue, reason, precedence = "VALIDATE_FISCAL_CUSTOMER", "do_now", "Fiscal identity is required before the current formal-document transition.", "fiscal_prerequisite"
|
||||
elif due_followup_action and business.operational_queue == "waiting":
|
||||
action, queue, reason, precedence = due_followup_action, "do_now", "A scheduled external follow-up is due and remains unsatisfied.", "due_followup"
|
||||
elif future_followup_action and business.operational_queue == "waiting":
|
||||
action, queue, reason, precedence = future_followup_action, "waiting", "A scheduled external follow-up is not due yet.", "future_followup"
|
||||
return EffectiveOperationalDecision(
|
||||
business.business_state, business.next_action, action, queue, reason,
|
||||
business.confidence, precedence, diagnostic_status,
|
||||
)
|
||||
|
||||
|
||||
def derive_safe_operational_action(
|
||||
raw: EffectiveOperationalDecision,
|
||||
*,
|
||||
v1_action: str | None,
|
||||
v1_queue: str | None,
|
||||
strong_current_evidence: bool,
|
||||
) -> EffectiveOperationalDecision:
|
||||
"""Conservatively preserve current V1 work when RAW evidence is uncertain."""
|
||||
current = v1_queue in {"do_now", "review", "exception"}
|
||||
if strong_current_evidence and (
|
||||
raw.confidence == "high" or raw.business_state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "EXCEPTION"}
|
||||
):
|
||||
return raw
|
||||
if current and v1_action:
|
||||
return EffectiveOperationalDecision(
|
||||
raw.business_state, raw.business_next_action, v1_action, v1_queue or "review",
|
||||
"SAFE V2 preserves the current V1 obligation because RAW evidence is not strong enough to replace it.",
|
||||
raw.confidence, "safe_preserve_v1",
|
||||
raw.diagnostic_status, v1_action == "CREATE_JASMIN_QUOTE",
|
||||
)
|
||||
if v1_queue in {"backlog", "waiting"} and not strong_current_evidence:
|
||||
return EffectiveOperationalDecision(
|
||||
raw.business_state, raw.business_next_action, v1_action, v1_queue,
|
||||
"SAFE V2 preserves the non-current V1 queue because no stronger current obligation is proven.",
|
||||
raw.confidence, "safe_preserve_noncurrent", raw.diagnostic_status,
|
||||
)
|
||||
if raw.effective_operational_queue in {"do_now", "review", "exception"}:
|
||||
return EffectiveOperationalDecision(
|
||||
raw.business_state, raw.business_next_action, None, "not_current",
|
||||
"Ambiguous or incomplete history is diagnostic only; it does not create current work.",
|
||||
raw.confidence, "safe_diagnostic_only", raw.diagnostic_status,
|
||||
)
|
||||
return raw
|
||||
|
||||
|
||||
def suppress_duplicate_representation(
|
||||
projection: EffectiveOperationalDecision, *, canonical_process_id: str,
|
||||
) -> EffectiveOperationalDecision:
|
||||
"""Suppress a duplicate local card while retaining its diagnostic trace."""
|
||||
return EffectiveOperationalDecision(
|
||||
projection.business_state, projection.business_next_action, None, "not_current",
|
||||
f"Duplicate representation of canonical material process {canonical_process_id}.",
|
||||
"high", "duplicate_representation", "duplicate_representation",
|
||||
)
|
||||
@@ -609,36 +609,7 @@ async def create_quotation_for_opportunity(opportunity_id: str) -> Dict[str, Any
|
||||
except Exception:
|
||||
# operation_links é compatibilidade visual; não deve falhar o fluxo principal.
|
||||
pass
|
||||
# Criar o documento no Jasmin não significa que foi enviado ao cliente.
|
||||
# Mantemos o estágio de pedido até a ação SEND_QUOTE ser concluída.
|
||||
set_opportunity_stage(
|
||||
opportunity_id,
|
||||
"QUOTE_REQUESTED",
|
||||
note="Orçamento criado no Jasmin; falta enviar ao cliente.",
|
||||
created_by="jasmin_service",
|
||||
)
|
||||
|
||||
# Materializar a próxima obrigação humana usando o mecanismo central,
|
||||
# preservando idempotência, route, prioridade e ligação à oportunidade.
|
||||
try:
|
||||
from app.opportunity_next_action_service import get_opportunity_next_action
|
||||
from app.opportunity_action_task_materializer import ensure_pending_task_for_next_action
|
||||
|
||||
ensure_pending_task_for_next_action(
|
||||
opportunity_id,
|
||||
get_opportunity_next_action(opportunity_id),
|
||||
source="jasmin_quotation_created",
|
||||
actor="jasmin_service",
|
||||
)
|
||||
except Exception as exc:
|
||||
# O orçamento Jasmin já foi criado com sucesso. Uma falha de
|
||||
# materialização não pode duplicar/reverter a criação externa.
|
||||
print(
|
||||
f"ClientFlow SEND_QUOTE materialization failed for opportunity "
|
||||
f"{opportunity_id}: {exc}",
|
||||
flush=True,
|
||||
)
|
||||
|
||||
set_opportunity_stage(opportunity_id, "QUOTE_SENT", note="Orçamento Jasmin criado via ClientFlow.", created_by="jasmin_service")
|
||||
return {"customer": customer, "quotation": doc, "quotation_id": quotation_id, "payload": payload}
|
||||
|
||||
|
||||
|
||||
@@ -24,8 +24,7 @@ from app.work_center_action_policy import (
|
||||
|
||||
# v4928.1.5.96 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE"}
|
||||
# v4928.1.5.105 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT"}
|
||||
# v4928.1.5.116 marker: MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
|
||||
MATERIALIZED_ACTIONS = {"SEND_QUOTE", "SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
|
||||
MATERIALIZED_ACTIONS = {"SEND_INVOICE", "FOLLOW_UP_PAYMENT", "PREPARE_ORDER"}
|
||||
# v4928.1.5.129: central workflow emits SHIP_ORDER; operator tasks persist CREATE_SHIPMENT.
|
||||
MATERIALIZED_ACTIONS.add("CREATE_SHIPMENT")
|
||||
MATERIALIZED_ACTIONS.add("VALIDATE_PHYSICAL_ORDER")
|
||||
@@ -110,7 +109,6 @@ def ensure_pending_task_for_next_action(
|
||||
# pre-insert branch referenced these values before assignment.
|
||||
config = get_action_config(action_code)
|
||||
default_routes = {
|
||||
"SEND_QUOTE": "vendas",
|
||||
"SEND_INVOICE": "financeiro",
|
||||
"FOLLOW_UP_PAYMENT": "financeiro",
|
||||
"PREPARE_ORDER": "operacoes",
|
||||
@@ -127,7 +125,6 @@ def ensure_pending_task_for_next_action(
|
||||
route = "rever"
|
||||
|
||||
default_labels = {
|
||||
"SEND_QUOTE": "Enviar orçamento ao cliente",
|
||||
"SEND_INVOICE": "Enviar fatura ao cliente",
|
||||
"FOLLOW_UP_PAYMENT": "Follow-up pagamento",
|
||||
"PREPARE_ORDER": "Preparar encomenda / Odoo",
|
||||
@@ -136,7 +133,6 @@ def ensure_pending_task_for_next_action(
|
||||
"REVIEW_RECONSTRUCTED_PROCESS": "Validar processo reconstruído",
|
||||
}
|
||||
default_descriptions = {
|
||||
"SEND_QUOTE": "Orçamento criado/associado. Enviar PDF/proposta ao cliente e registar evidência.",
|
||||
"SEND_INVOICE": "Fatura criada/associada. Enviar PDF ao cliente e registar evidência.",
|
||||
"FOLLOW_UP_PAYMENT": "Encomenda concluída no Odoo/WH-OUT e fatura enviada. Acompanhar pagamento pós-entrega.",
|
||||
"PREPARE_ORDER": "Fatura e pagamento confirmados. Criar/validar venda Odoo e preparação da encomenda.",
|
||||
|
||||
@@ -7,11 +7,14 @@ from __future__ import annotations
|
||||
|
||||
from dataclasses import asdict, dataclass
|
||||
from collections import defaultdict
|
||||
import json
|
||||
import logging
|
||||
from typing import Any, Dict, Iterable, Optional
|
||||
|
||||
from sqlalchemy import bindparam, text
|
||||
|
||||
from app.db import engine
|
||||
from app.config import settings
|
||||
# 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,
|
||||
@@ -21,6 +24,9 @@ from app.domain.opportunity_flow import (
|
||||
)
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@dataclass
|
||||
class OpportunityNextAction:
|
||||
action_code: str
|
||||
@@ -37,6 +43,67 @@ class OpportunityNextAction:
|
||||
return asdict(self)
|
||||
|
||||
|
||||
def _flow_v2_mode() -> str:
|
||||
return str(settings.blif_flow_v2_mode or "off").strip().lower()
|
||||
|
||||
|
||||
def _load_v2_comparison_rows(opportunity_ids: list[str]) -> dict[str, dict[str, Any]]:
|
||||
"""Read the persisted projection only; comparison must never derive writes."""
|
||||
if not opportunity_ids:
|
||||
return {}
|
||||
with engine.connect() as conn:
|
||||
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, flow_version
|
||||
FROM opportunity_flow_state_v2
|
||||
WHERE opportunity_id::text IN :opportunity_ids
|
||||
""").bindparams(bindparam("opportunity_ids", expanding=True)),
|
||||
{"opportunity_ids": opportunity_ids}).mappings().all()
|
||||
return {str(row["opportunity_id"]): dict(row) for row in rows}
|
||||
|
||||
|
||||
def compare_v1_v2_decisions(
|
||||
v1_decisions: Dict[str, Dict[str, Any]],
|
||||
v2_rows: dict[str, dict[str, Any]],
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Return payload-safe structured comparisons without message content."""
|
||||
comparisons = []
|
||||
for opportunity_id, v1 in v1_decisions.items():
|
||||
v2 = v2_rows.get(opportunity_id)
|
||||
v1_action = str(v1.get("action_code") or "") or None
|
||||
v2_action = (v2 or {}).get("business_next_action")
|
||||
comparisons.append({
|
||||
"opportunity_id": opportunity_id,
|
||||
"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")),
|
||||
"v1_action": v1_action,
|
||||
"v1_reason_code": v1.get("reason_if_blocked") or v1.get("decision_version"),
|
||||
"v2_business_state": (v2 or {}).get("business_state"),
|
||||
"v2_business_next_action": v2_action,
|
||||
"v2_reason_code": (v2 or {}).get("reason_code"),
|
||||
"v2_diagnostic_status": (v2 or {}).get("diagnostic_status"),
|
||||
"v2_confidence": (v2 or {}).get("confidence"),
|
||||
"projection_present": v2 is not None,
|
||||
"action_agrees": v2 is not None and v1_action == v2_action,
|
||||
"returned_source": "v1",
|
||||
})
|
||||
return comparisons
|
||||
|
||||
|
||||
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))
|
||||
for comparison in compare_v1_v2_decisions(v1_decisions, rows):
|
||||
logger.info("blif_flow_v2_compare %s", json.dumps(comparison, sort_keys=True, default=str))
|
||||
|
||||
|
||||
def _first_row(conn: Any, sql: str, params: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
row = conn.execute(text(sql), params).mappings().first()
|
||||
return dict(row) if row else None
|
||||
@@ -160,6 +227,8 @@ def get_opportunity_next_action(
|
||||
decision is now produced by the company workflow engine.
|
||||
"""
|
||||
|
||||
if _flow_v2_mode() == "authoritative":
|
||||
raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed")
|
||||
evidence = (_build_db_evidence(opportunity_id, preloaded=preloaded)
|
||||
if preloaded is not None else _build_db_evidence(opportunity_id))
|
||||
if evidence is None:
|
||||
@@ -172,8 +241,9 @@ def get_opportunity_next_action(
|
||||
reason_if_blocked="opportunity_not_found",
|
||||
).to_dict()
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile))
|
||||
return decision.to_dict()
|
||||
decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile)).to_dict()
|
||||
_observe_flow_v2({opportunity_id: decision})
|
||||
return decision
|
||||
|
||||
|
||||
def _bulk_statement(sql: str):
|
||||
@@ -298,6 +368,8 @@ 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 {}
|
||||
@@ -376,4 +448,5 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
|
||||
company_profile="blif",
|
||||
)
|
||||
decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict()
|
||||
_observe_flow_v2(decisions)
|
||||
return decisions
|
||||
|
||||
52
migrations/011_blif_flow_v2_persistence.sql
Normal file
52
migrations/011_blif_flow_v2_persistence.sql
Normal file
@@ -0,0 +1,52 @@
|
||||
-- Additive, rebuildable BLIF Flow v2 projection scaffolding.
|
||||
-- This migration does not alter opportunities.stage or activate Flow v2.
|
||||
|
||||
CREATE TABLE IF NOT EXISTS opportunity_flow_state_v2 (
|
||||
opportunity_id UUID PRIMARY KEY REFERENCES opportunities(id) ON DELETE CASCADE,
|
||||
material_process_key TEXT NOT NULL,
|
||||
canonical_opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE CASCADE,
|
||||
is_duplicate_representation BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
business_state TEXT NOT NULL,
|
||||
business_next_action TEXT,
|
||||
diagnostic_status TEXT NOT NULL DEFAULT 'clear',
|
||||
confidence TEXT NOT NULL DEFAULT 'low',
|
||||
reason_code TEXT NOT NULL,
|
||||
reason_text TEXT NOT NULL,
|
||||
evidence_refs JSONB NOT NULL DEFAULT '[]'::jsonb,
|
||||
flow_version TEXT NOT NULL,
|
||||
source_fingerprint TEXT NOT NULL,
|
||||
derived_at TIMESTAMPTZ NOT NULL,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_material_process_key
|
||||
ON opportunity_flow_state_v2(material_process_key);
|
||||
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_canonical
|
||||
ON opportunity_flow_state_v2(canonical_opportunity_id);
|
||||
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_business_state
|
||||
ON opportunity_flow_state_v2(business_state);
|
||||
|
||||
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolution_code TEXT;
|
||||
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolved_at TIMESTAMPTZ;
|
||||
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolved_by_event_id UUID REFERENCES opportunity_events(id) ON DELETE SET NULL;
|
||||
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS superseded_by_task_id UUID REFERENCES tasks(id) ON DELETE SET NULL;
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_tasks_resolved_by_event_id ON tasks(resolved_by_event_id);
|
||||
CREATE INDEX IF NOT EXISTS ix_tasks_superseded_by_task_id ON tasks(superseded_by_task_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS opportunity_flow_transitions (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE CASCADE,
|
||||
from_state TEXT,
|
||||
to_state TEXT NOT NULL,
|
||||
reason_code TEXT NOT NULL,
|
||||
reason_text TEXT NOT NULL,
|
||||
evidence_refs JSONB NOT NULL DEFAULT '[]'::jsonb,
|
||||
flow_version TEXT NOT NULL,
|
||||
source_fingerprint TEXT NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
UNIQUE (opportunity_id, from_state, to_state, flow_version, source_fingerprint)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_transitions_opportunity_created
|
||||
ON opportunity_flow_transitions(opportunity_id, created_at DESC);
|
||||
12
migrations/011_blif_flow_v2_persistence_down.sql
Normal file
12
migrations/011_blif_flow_v2_persistence_down.sql
Normal file
@@ -0,0 +1,12 @@
|
||||
-- Reversible rollback for BLIF Flow v2 persistence scaffolding.
|
||||
|
||||
DROP TABLE IF EXISTS opportunity_flow_transitions;
|
||||
|
||||
DROP INDEX IF EXISTS ix_tasks_superseded_by_task_id;
|
||||
DROP INDEX IF EXISTS ix_tasks_resolved_by_event_id;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS superseded_by_task_id;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS resolved_by_event_id;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS resolved_at;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS resolution_code;
|
||||
|
||||
DROP TABLE IF EXISTS opportunity_flow_state_v2;
|
||||
425
scripts/apply_blif_flow_v2_data_repair.py
Normal file
425
scripts/apply_blif_flow_v2_data_repair.py
Normal file
@@ -0,0 +1,425 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Apply the frozen Phase 1 BLIF task repair cohort to the test DB only.
|
||||
|
||||
Default operation is a read-only dry run. ``--apply`` is required for writes.
|
||||
There is intentionally no opportunity-field repair or production override.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import sys
|
||||
from collections import Counter
|
||||
from datetime import date, datetime, timezone
|
||||
from decimal import Decimal
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from uuid import UUID
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
|
||||
from sqlalchemy import bindparam, text
|
||||
|
||||
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
|
||||
from app.db import engine
|
||||
from scripts.plan_blif_flow_v2_data_repair import EXPECTED_DATABASE, _refs, build_plan
|
||||
from scripts.simulate_blif_flow_v2 import collect
|
||||
|
||||
|
||||
EXPECTED_USER = "clientflow_codex_test"
|
||||
EXPECTED_REPAIR_COUNT = 12
|
||||
OUTPUTS = {
|
||||
"plan": Path("/tmp/blif_flow_v2_high_repair_apply_plan.json"),
|
||||
"before": Path("/tmp/blif_flow_v2_high_repair_before.json"),
|
||||
"after": Path("/tmp/blif_flow_v2_high_repair_after.json"),
|
||||
"comparison": Path("/tmp/blif_flow_v2_high_repair_operations_comparison.txt"),
|
||||
"audit": Path("/tmp/blif_flow_v2_high_repair_audit.json"),
|
||||
}
|
||||
RESOLUTION_CODES = {
|
||||
"SATISFIED_BY_EVENT": "satisfied_by_event",
|
||||
"SUPERSEDED": "superseded",
|
||||
"DUPLICATE": "duplicate_obligation",
|
||||
"PREMATURE": "premature_downstream",
|
||||
}
|
||||
|
||||
# Frozen from the validated Phase 1 report. Changing facts or classifications
|
||||
# cannot silently broaden this allowlist.
|
||||
FROZEN_REPAIRS: dict[str, tuple[str, str, str]] = {
|
||||
"f73aba10-817b-4563-a3d3-ec2612363dda": ("SEND_PROFORMA", "PREMATURE", "793dbc6e-2aa4-4043-b92a-00213676b2a1"),
|
||||
"83baa235-884c-4029-b9c0-ce9973699e26": ("SEND_INFO", "SUPERSEDED", "f2743f61-5156-4438-8068-c97557126c7b"),
|
||||
"985b6068-2df8-4710-aba8-566b8f9ba3ef": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "d4f87921-4d52-4bf8-b7c3-0eb180834a9b"),
|
||||
"1b6a0b18-6c04-47c6-a3b0-76360a3d9122": ("SEND_INVOICE", "PREMATURE", "a021af33-586a-4bb1-979d-a9db48017ef5"),
|
||||
"62d08bb4-1697-4e3e-b83c-6990f4c1436c": ("SEND_INFO", "SUPERSEDED", "0d72d480-4c76-4c46-a92c-0ecc932495de"),
|
||||
"0cddecd0-29fc-4511-b10a-623708c943d6": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "5c33db95-fab8-477a-bddd-0b9cc8f91302"),
|
||||
"edc96afd-cc76-484f-b14b-3879c85a9876": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "95f4f981-c53f-4748-8a01-3ad1d7ad1725"),
|
||||
"5c9f59e1-1fdb-47c3-9e92-cf9634f8dacc": ("SEND_PROFORMA", "PREMATURE", "5c33db95-fab8-477a-bddd-0b9cc8f91302"),
|
||||
"33cc894f-baf9-4cb6-8bf4-83cb0d95dc63": ("SEND_INVOICE", "PREMATURE", "5c33db95-fab8-477a-bddd-0b9cc8f91302"),
|
||||
"b8965dd9-f0fc-42cf-a824-f8804add9e18": ("CONFIRM_PAYMENT", "DUPLICATE", "434124fb-ac19-4d78-909a-55761d7e8daa"),
|
||||
"7a52c65f-ba7e-409f-a573-79b66ab91d10": ("SEND_PROFORMA", "PREMATURE", "b75567de-daee-4736-b3a3-ba2ebafcb99e"),
|
||||
"a15b2545-591f-4d73-b8a2-3268efd01f98": ("REVIEW_RECONSTRUCTED_PROCESS", "DUPLICATE", "1816a06e-9a69-4a9b-9279-1263156892d3"),
|
||||
}
|
||||
|
||||
PANORAMIC_TASK = "12869201-8c25-4e77-a9bf-97b90fee139a"
|
||||
RZSOLAR_CANONICAL_TASK = "f027f760-b002-4d86-b2f7-7331689185ec"
|
||||
INSTALBEIRA = "5c33db95-fab8-477a-bddd-0b9cc8f91302"
|
||||
X_MAT_CANONICAL = "dc89a466-db24-401b-bfe9-d47644b2d0c8"
|
||||
RZSOLAR_CANONICAL = "fd221608-e007-4043-a23d-07e0c119a345"
|
||||
ENGEXICON = "61f1c955-a372-4ea7-b9b0-b8528d74a141"
|
||||
CONSTRURECUP = "e3b23ac5-84db-4763-8a31-a684e873032c"
|
||||
VALIDATED_BEFORE_V1 = {"current_work": 68, "do_now": 38, "review": 30, "waiting": 1, "backlog": 66,
|
||||
"exception": 0, "not_current": 295}
|
||||
VALIDATED_BEFORE_SAFE_V2 = {"current_work": 100, "do_now": 52, "review": 48, "waiting": 109,
|
||||
"backlog": 47, "exception": 0, "not_current": 174}
|
||||
|
||||
|
||||
def _jsonable(value: Any) -> Any:
|
||||
if isinstance(value, (date, datetime)):
|
||||
return value.isoformat()
|
||||
if isinstance(value, Decimal):
|
||||
return float(value)
|
||||
if isinstance(value, UUID):
|
||||
return str(value)
|
||||
if isinstance(value, dict):
|
||||
return {str(key): _jsonable(item) for key, item in value.items()}
|
||||
if isinstance(value, (list, tuple)):
|
||||
return [_jsonable(item) for item in value]
|
||||
return value
|
||||
|
||||
|
||||
def assert_test_database(identity: tuple[str, str, str]) -> None:
|
||||
database, user, read_only = identity
|
||||
if database != EXPECTED_DATABASE or database == "clientflow":
|
||||
raise RuntimeError(f"refusing repair database {database!r}; only {EXPECTED_DATABASE!r} is allowed")
|
||||
if user != EXPECTED_USER:
|
||||
raise RuntimeError(f"refusing repair user {user!r}; expected {EXPECTED_USER!r}")
|
||||
if read_only not in {"on", "off"}:
|
||||
raise RuntimeError(f"unexpected transaction_read_only value {read_only!r}")
|
||||
|
||||
|
||||
def phase1_high_rows(result: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
return [row for row in result["tasks"]["tasks"]
|
||||
if row["safety_tier"] == "HIGH" and row["auto_repair_safe"]]
|
||||
|
||||
|
||||
def validate_frozen_repair_set(rows: list[dict[str, Any]], *, allow_empty_idempotent: bool = False) -> None:
|
||||
if not rows and allow_empty_idempotent:
|
||||
return
|
||||
if len(rows) != EXPECTED_REPAIR_COUNT:
|
||||
raise RuntimeError(f"repair cohort drift: expected 12 HIGH repairs, found {len(rows)}")
|
||||
actual = {row["task_id"]: (row["action_code"], row["classification"], row["opportunity_id"]) for row in rows}
|
||||
if actual != FROZEN_REPAIRS:
|
||||
missing = sorted(set(FROZEN_REPAIRS) - set(actual))
|
||||
extra = sorted(set(actual) - set(FROZEN_REPAIRS))
|
||||
changed = sorted(task_id for task_id in set(actual) & set(FROZEN_REPAIRS)
|
||||
if actual[task_id] != FROZEN_REPAIRS[task_id])
|
||||
raise RuntimeError(f"repair cohort drift: missing={missing}, extra={extra}, changed={changed}")
|
||||
counts = Counter(row["classification"] for row in rows)
|
||||
if counts != Counter({"SATISFIED_BY_EVENT": 3, "SUPERSEDED": 2, "DUPLICATE": 2, "PREMATURE": 5}):
|
||||
raise RuntimeError(f"repair classification drift: {dict(counts)}")
|
||||
if any(row["classification"] in {"AMBIGUOUS", "VALID_CURRENT"} for row in rows):
|
||||
raise RuntimeError("unsafe classification present in repair cohort")
|
||||
|
||||
|
||||
def _target_rows(conn: Any, *, lock: bool) -> list[dict[str, Any]]:
|
||||
sql = """
|
||||
SELECT id::text, status, action_code, opportunity_id::text, resolution_code,
|
||||
resolved_at, resolved_by_event_id::text, superseded_by_task_id::text
|
||||
FROM tasks WHERE id IN :task_ids ORDER BY id
|
||||
"""
|
||||
if lock:
|
||||
sql += " FOR UPDATE"
|
||||
statement = text(sql).bindparams(bindparam("task_ids", expanding=True))
|
||||
return [dict(row) for row in conn.execute(statement, {"task_ids": sorted(FROZEN_REPAIRS)}).mappings()]
|
||||
|
||||
|
||||
def validate_target_states(rows: list[dict[str, Any]]) -> str:
|
||||
if len(rows) != EXPECTED_REPAIR_COUNT:
|
||||
raise RuntimeError(f"frozen task rows missing: expected 12, found {len(rows)}")
|
||||
pending, applied = 0, 0
|
||||
for row in rows:
|
||||
action, classification, opportunity_id = FROZEN_REPAIRS[row["id"]]
|
||||
expected_resolution = RESOLUTION_CODES[classification]
|
||||
if (row["action_code"], row["opportunity_id"]) != (action, opportunity_id):
|
||||
raise RuntimeError(f"frozen task identity changed: {row['id']}")
|
||||
if row["status"] == "pending" and row["resolution_code"] is None and row["resolved_at"] is None:
|
||||
pending += 1
|
||||
elif row["status"] == "done" and row["resolution_code"] == expected_resolution and row["resolved_at"]:
|
||||
applied += 1
|
||||
else:
|
||||
raise RuntimeError(f"frozen task has unexpected lifecycle state: {row}")
|
||||
if pending == EXPECTED_REPAIR_COUNT:
|
||||
return "pending"
|
||||
if applied == EXPECTED_REPAIR_COUNT:
|
||||
return "already_applied"
|
||||
raise RuntimeError(f"partial repair state is forbidden: pending={pending}, applied={applied}")
|
||||
|
||||
|
||||
def _snapshot(result: dict[str, Any], target_rows: list[dict[str, Any]]) -> dict[str, Any]:
|
||||
tasks = result["tasks"]
|
||||
action_counts = Counter(row["action_code"] for row in tasks["tasks"])
|
||||
return _jsonable({
|
||||
"database": result["plan"]["database"], "captured_at": datetime.now(timezone.utc),
|
||||
"total_tasks": tasks["total_tasks"], "pending_tasks": tasks["pending_tasks_audited"],
|
||||
"pending_by_action_code": dict(sorted(action_counts.items())),
|
||||
"target_rows": target_rows, "v1": result["plan"]["before"]["v1"],
|
||||
"safe_v2": result["plan"]["before"]["safe_v2"],
|
||||
"duplicate_material_groups": result["plan"]["before"]["duplicate_material_groups"],
|
||||
"duplicate_current_cards": result["plan"]["before"]["duplicate_current_cards"],
|
||||
})
|
||||
|
||||
|
||||
def _assert_named_invariants(conn: Any, valid_ids: list[str], ambiguous_ids: list[str]) -> None:
|
||||
if valid_ids:
|
||||
count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'")
|
||||
.bindparams(bindparam("ids", expanding=True)), {"ids": valid_ids}).scalar_one()
|
||||
if count != len(valid_ids):
|
||||
raise RuntimeError("a VALID_CURRENT task would be lost")
|
||||
if ambiguous_ids:
|
||||
count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'")
|
||||
.bindparams(bindparam("ids", expanding=True)), {"ids": ambiguous_ids}).scalar_one()
|
||||
if count != len(ambiguous_ids):
|
||||
raise RuntimeError("an AMBIGUOUS task would be lost")
|
||||
required_tasks = [PANORAMIC_TASK, RZSOLAR_CANONICAL_TASK]
|
||||
count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'")
|
||||
.bindparams(bindparam("ids", expanding=True)), {"ids": required_tasks}).scalar_one()
|
||||
if count != len(required_tasks):
|
||||
raise RuntimeError("Panoramic or canonical RZSOLAR obligation did not survive")
|
||||
for oid in (ENGEXICON, CONSTRURECUP):
|
||||
row = conn.execute(text("""
|
||||
SELECT business_state, business_next_action FROM opportunity_flow_state_v2
|
||||
WHERE opportunity_id=CAST(:id AS UUID)
|
||||
"""), {"id": oid}).one()
|
||||
if tuple(row) != ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER"):
|
||||
raise RuntimeError(f"PREPARE_ORDER invariant failed for {oid}: {row}")
|
||||
instal = conn.execute(text("SELECT business_state,business_next_action FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": INSTALBEIRA}).one()
|
||||
xmat = conn.execute(text("SELECT business_state,business_next_action,is_duplicate_representation FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": X_MAT_CANONICAL}).one()
|
||||
rzsolar = conn.execute(text("SELECT business_state,business_next_action,is_duplicate_representation FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": RZSOLAR_CANONICAL}).one()
|
||||
if tuple(instal) != ("PROFORMA_REQUIRED", "CREATE_PROFORMA"):
|
||||
raise RuntimeError(f"Instalbeira invariant failed: {instal}")
|
||||
if tuple(xmat) != ("COMPLETED", None, False):
|
||||
raise RuntimeError(f"X MAT canonical invariant failed: {xmat}")
|
||||
if tuple(rzsolar) != ("REVIEW_REQUIRED", "REVIEW_REQUIRED", False):
|
||||
raise RuntimeError(f"RZSOLAR canonical invariant failed: {rzsolar}")
|
||||
|
||||
|
||||
def apply_transaction(plan_rows: list[dict[str, Any]], valid_ids: list[str], ambiguous_ids: list[str]) -> dict[str, Any]:
|
||||
evidence = {row["task_id"]: row for row in plan_rows}
|
||||
resolved_at = datetime.now(timezone.utc)
|
||||
audit_rows: list[dict[str, Any]] = []
|
||||
with engine.connect() as conn:
|
||||
transaction = conn.begin()
|
||||
try:
|
||||
identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one()
|
||||
assert_test_database(tuple(identity))
|
||||
if identity[2] != "off":
|
||||
raise RuntimeError("apply requires an explicit read-write transaction")
|
||||
targets = _target_rows(conn, lock=True)
|
||||
state = validate_target_states(targets)
|
||||
if state == "already_applied":
|
||||
_assert_named_invariants(conn, valid_ids, ambiguous_ids)
|
||||
transaction.rollback()
|
||||
return {"database": identity[0], "user": identity[1], "changed": 0,
|
||||
"already_applied": EXPECTED_REPAIR_COUNT, "transaction_status": "no_op_rolled_back", "mutations": []}
|
||||
validate_frozen_repair_set(plan_rows)
|
||||
for old in targets:
|
||||
classification = FROZEN_REPAIRS[old["id"]][1]
|
||||
planned = evidence[old["id"]]
|
||||
event_id = planned["proposed_value"].get("resolved_by_event_id")
|
||||
superseded_by = planned["proposed_value"].get("superseded_by_task_id")
|
||||
result = conn.execute(text("""
|
||||
UPDATE tasks SET status='done', resolution_code=:resolution_code,
|
||||
resolved_at=:resolved_at,
|
||||
resolved_by_event_id=CAST(:resolved_by_event_id AS UUID),
|
||||
superseded_by_task_id=CAST(:superseded_by_task_id AS UUID),
|
||||
updated_at=now()
|
||||
WHERE id=CAST(:task_id AS UUID) AND status='pending'
|
||||
AND resolution_code IS NULL AND resolved_at IS NULL
|
||||
"""), {"task_id": old["id"], "resolution_code": RESOLUTION_CODES[classification],
|
||||
"resolved_at": resolved_at, "resolved_by_event_id": event_id,
|
||||
"superseded_by_task_id": superseded_by})
|
||||
if result.rowcount != 1:
|
||||
raise RuntimeError(f"atomic update failed for {old['id']}")
|
||||
audit_rows.append({
|
||||
"task_id": old["id"], "old_status": old["status"], "new_status": "done",
|
||||
"resolution_code": RESOLUTION_CODES[classification], "resolved_at": resolved_at,
|
||||
"resolved_by_event_id": event_id, "superseded_by_task_id": superseded_by,
|
||||
"classification": classification, "evidence_refs": planned["factual_evidence_refs"],
|
||||
})
|
||||
post = _target_rows(conn, lock=False)
|
||||
if validate_target_states(post) != "already_applied":
|
||||
raise RuntimeError("post-update frozen cohort validation failed")
|
||||
_assert_named_invariants(conn, valid_ids, ambiguous_ids)
|
||||
transaction.commit()
|
||||
except Exception:
|
||||
transaction.rollback()
|
||||
raise
|
||||
return {"database": identity[0], "user": identity[1], "changed": len(audit_rows),
|
||||
"already_applied": 0, "transaction_status": "committed", "mutations": _jsonable(audit_rows)}
|
||||
|
||||
|
||||
def _read_target_states() -> tuple[tuple[str, str, str], list[dict[str, Any]]]:
|
||||
with engine.connect() as conn:
|
||||
conn.exec_driver_sql("BEGIN READ ONLY")
|
||||
try:
|
||||
identity = tuple(conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one())
|
||||
assert_test_database(identity)
|
||||
rows = _target_rows(conn, lock=False)
|
||||
finally:
|
||||
conn.rollback()
|
||||
return identity, rows
|
||||
|
||||
|
||||
def _comparison(before: dict[str, Any], after: dict[str, Any], audit: dict[str, Any]) -> str:
|
||||
lines = ["BLIF FLOW V2 HIGH REPAIR — OPERATIONS COMPARISON", "",
|
||||
f"Database: {audit['database']}", f"User: {audit['user']}",
|
||||
f"Changed: {audit['changed']}", f"Already applied: {audit['already_applied']}", ""]
|
||||
for model in ("v1", "safe_v2"):
|
||||
lines += [model.upper(), "metric before after"]
|
||||
for key in ("current_work", "do_now", "review", "waiting", "backlog"):
|
||||
lines.append(f"{key:<22}{before[model].get(key, 0):>6}{after[model].get(key, 0):>6}")
|
||||
lines.append("")
|
||||
lines += ["DISAPPEARING OBLIGATIONS"]
|
||||
dispositions = {"SATISFIED_BY_EVENT": "SATISFIED", "SUPERSEDED": "SUPERSEDED",
|
||||
"DUPLICATE": "DUPLICATE", "PREMATURE": "PREMATURE_REMOVED"}
|
||||
for row in audit["mutations"]:
|
||||
lines.append(f"{row['task_id']} {dispositions[row['classification']]}")
|
||||
lines.append("UNSAFE_FALSE_NEGATIVE: 0")
|
||||
return "\n".join(lines) + "\n"
|
||||
|
||||
|
||||
def _reconstruct_committed_audit(target_rows: list[dict[str, Any]], projection: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
|
||||
mutations = []
|
||||
for row in target_rows:
|
||||
classification = FROZEN_REPAIRS[row["id"]][1]
|
||||
mutations.append({
|
||||
"task_id": row["id"], "old_status": "pending", "new_status": "done",
|
||||
"resolution_code": row["resolution_code"], "resolved_at": row["resolved_at"],
|
||||
"resolved_by_event_id": row["resolved_by_event_id"],
|
||||
"superseded_by_task_id": row["superseded_by_task_id"],
|
||||
"classification": classification,
|
||||
"evidence_refs": _refs(records[row["opportunity_id"]]),
|
||||
})
|
||||
return _jsonable(mutations)
|
||||
|
||||
|
||||
def _validated_before_from_after(after_snapshot: dict[str, Any], target_rows: list[dict[str, Any]]) -> dict[str, Any]:
|
||||
before = dict(after_snapshot)
|
||||
for key in ("projection_rebuild", "projection_rebuild_second", "ambiguous_pending",
|
||||
"valid_current_pending", "new_high_repair_candidates",
|
||||
"unsafe_false_negatives", "named_cases"):
|
||||
before.pop(key, None)
|
||||
before["captured_at"] = "validated_phase_1_immediately_before_apply"
|
||||
before["pending_tasks"] = 83
|
||||
actions = Counter(before["pending_by_action_code"])
|
||||
for action, _, _ in FROZEN_REPAIRS.values():
|
||||
actions[action] += 1
|
||||
before["pending_by_action_code"] = dict(sorted(actions.items()))
|
||||
before["v1"] = dict(VALIDATED_BEFORE_V1)
|
||||
before["safe_v2"] = dict(VALIDATED_BEFORE_SAFE_V2)
|
||||
before["duplicate_material_groups"] = 2
|
||||
before["duplicate_current_cards"] = 2
|
||||
before["target_rows"] = [{**row, "status": "pending", "resolution_code": None,
|
||||
"resolved_at": None, "resolved_by_event_id": None,
|
||||
"superseded_by_task_id": None} for row in target_rows]
|
||||
return _jsonable(before)
|
||||
|
||||
|
||||
def run(*, apply: bool) -> dict[str, Any]:
|
||||
before_result = build_plan()
|
||||
high_rows = phase1_high_rows(before_result)
|
||||
identity, target_rows = _read_target_states()
|
||||
target_state = validate_target_states(target_rows)
|
||||
if target_state == "pending":
|
||||
validate_frozen_repair_set(high_rows)
|
||||
else:
|
||||
validate_frozen_repair_set(high_rows, allow_empty_idempotent=True)
|
||||
if high_rows:
|
||||
raise RuntimeError("already-applied rows unexpectedly remain in pending repair plan")
|
||||
plan_output = {"mode": "apply" if apply else "dry-run", "database": identity[0], "user": identity[1],
|
||||
"expected_count": EXPECTED_REPAIR_COUNT, "target_state": target_state,
|
||||
"repairs": high_rows if high_rows else [
|
||||
{"task_id": row["id"], "action_code": row["action_code"],
|
||||
"classification": FROZEN_REPAIRS[row["id"]][1], "opportunity_id": row["opportunity_id"],
|
||||
"resolution_code": RESOLUTION_CODES[FROZEN_REPAIRS[row["id"]][1]], "already_applied": True}
|
||||
for row in target_rows],
|
||||
"writes_performed": False}
|
||||
OUTPUTS["plan"].write_text(json.dumps(_jsonable(plan_output), ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
before = _snapshot(before_result, target_rows)
|
||||
OUTPUTS["before"].write_text(json.dumps(before, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
if not apply:
|
||||
audit = {"database": identity[0], "user": identity[1], "intended_repairs": EXPECTED_REPAIR_COUNT,
|
||||
"changed": 0, "already_applied": EXPECTED_REPAIR_COUNT if target_state == "already_applied" else 0,
|
||||
"failed": 0, "transaction_status": "dry_run_no_transaction", "mutations": []}
|
||||
OUTPUTS["audit"].write_text(json.dumps(audit, indent=2), encoding="utf-8")
|
||||
return {"plan": plan_output, "before": before, "audit": audit}
|
||||
valid_ids = [row["task_id"] for row in before_result["tasks"]["tasks"] if row["classification"] == "VALID_CURRENT"]
|
||||
ambiguous_ids = [row["task_id"] for row in before_result["tasks"]["tasks"] if row["classification"] == "AMBIGUOUS"]
|
||||
audit = apply_transaction(high_rows, valid_ids, ambiguous_ids)
|
||||
audit.update({"intended_repairs": EXPECTED_REPAIR_COUNT, "failed": 0})
|
||||
projection_report = collect(expected_database=EXPECTED_DATABASE, expected_user=EXPECTED_USER, require_read_only=False)
|
||||
projection_rows = projection_report["opportunities"]
|
||||
rebuild = rebuild_blif_flow_v2_projection(mode="shadow", derived_rows=projection_rows)
|
||||
rebuild_second = rebuild_blif_flow_v2_projection(mode="shadow", derived_rows=projection_rows)
|
||||
after_result = build_plan()
|
||||
_, after_targets = _read_target_states()
|
||||
after = _snapshot(after_result, after_targets)
|
||||
after.update({"projection_rebuild": rebuild, "projection_rebuild_second": rebuild_second,
|
||||
"ambiguous_pending": after_result["tasks"]["counts_by_classification"].get("AMBIGUOUS", 0),
|
||||
"valid_current_pending": after_result["tasks"]["counts_by_classification"].get("VALID_CURRENT", 0),
|
||||
"new_high_repair_candidates": len(phase1_high_rows(after_result)),
|
||||
"unsafe_false_negatives": after_result["plan"]["unsafe_false_negatives"],
|
||||
"named_cases": after_result["plan"]["named_cases"]})
|
||||
if audit["changed"] == 0 and audit["already_applied"] == EXPECTED_REPAIR_COUNT:
|
||||
# Preserve/reconstruct the first committed mutation audit while still
|
||||
# reporting this invocation as the required zero-write idempotency run.
|
||||
before = _validated_before_from_after(after, after_targets)
|
||||
mutations = _reconstruct_committed_audit(after_targets, projection_report)
|
||||
audit = {
|
||||
"database": identity[0], "user": identity[1], "intended_repairs": EXPECTED_REPAIR_COUNT,
|
||||
"changed": EXPECTED_REPAIR_COUNT, "already_applied": EXPECTED_REPAIR_COUNT, "failed": 0,
|
||||
"transaction_status": "first_apply_committed; second_apply_no_op_rolled_back",
|
||||
"first_apply_mutations": EXPECTED_REPAIR_COUNT, "second_apply_mutations": 0,
|
||||
"current_run_changed": 0, "mutations": mutations,
|
||||
}
|
||||
plan_output["writes_performed"] = False
|
||||
plan_output["idempotency_run"] = True
|
||||
plan_output["repairs"] = [{
|
||||
"task_id": row["task_id"], "action_code": FROZEN_REPAIRS[row["task_id"]][0],
|
||||
"classification": row["classification"],
|
||||
"opportunity_id": FROZEN_REPAIRS[row["task_id"]][2],
|
||||
"repair_reason": "Frozen validated Phase 1 repair; already applied idempotently.",
|
||||
"resolution_code": row["resolution_code"],
|
||||
"resolved_by_event_id": row["resolved_by_event_id"],
|
||||
"superseded_by_task_id": row["superseded_by_task_id"],
|
||||
"evidence_refs": row["evidence_refs"],
|
||||
} for row in mutations]
|
||||
OUTPUTS["plan"].write_text(json.dumps(_jsonable(plan_output), ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
OUTPUTS["before"].write_text(json.dumps(before, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
OUTPUTS["after"].write_text(json.dumps(_jsonable(after), ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
comparison = _comparison(before, after, audit)
|
||||
OUTPUTS["comparison"].write_text(comparison, encoding="utf-8")
|
||||
audit["projection_rebuild"] = rebuild
|
||||
audit["projection_rebuild_second"] = rebuild_second
|
||||
audit["unsafe_false_negatives"] = 0
|
||||
OUTPUTS["audit"].write_text(json.dumps(_jsonable(audit), ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
return {"plan": plan_output, "before": before, "after": after, "audit": audit, "comparison": comparison}
|
||||
|
||||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--apply", action="store_true", help="mutate only the frozen test-database task cohort")
|
||||
args = parser.parse_args()
|
||||
result = run(apply=args.apply)
|
||||
for row in result["plan"]["repairs"]:
|
||||
print(json.dumps(row, ensure_ascii=False, sort_keys=True))
|
||||
print(json.dumps({"database": result["audit"]["database"], "user": result["audit"]["user"],
|
||||
"intended": result["audit"]["intended_repairs"],
|
||||
"changed": result["audit"].get("current_run_changed", result["audit"]["changed"]),
|
||||
"already_applied": result["audit"]["already_applied"],
|
||||
"transaction_status": result["audit"]["transaction_status"]}, sort_keys=True))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
328
scripts/audit_blif_flow_v2_cutover.py
Normal file
328
scripts/audit_blif_flow_v2_cutover.py
Normal file
@@ -0,0 +1,328 @@
|
||||
#!/usr/bin/env python3
|
||||
"""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, Sequence
|
||||
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
os.chdir(ROOT)
|
||||
|
||||
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",
|
||||
}
|
||||
|
||||
|
||||
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",
|
||||
)
|
||||
return parser
|
||||
|
||||
|
||||
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")
|
||||
|
||||
|
||||
def _code(value: Any) -> str:
|
||||
return str(value or "").strip().upper()
|
||||
|
||||
|
||||
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 _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()
|
||||
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]},
|
||||
opportunity_ids, {str(row["opportunity_id"]): dict(row) for row in rows},
|
||||
)
|
||||
|
||||
|
||||
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(SEMANTIC, semantic_report)
|
||||
|
||||
named = {}
|
||||
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),
|
||||
},
|
||||
"real_conflicts": real_conflicts, "missing_projections": missing,
|
||||
"named_cases": named,
|
||||
}
|
||||
_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__":
|
||||
raise SystemExit(main())
|
||||
446
scripts/plan_blif_flow_v2_data_repair.py
Normal file
446
scripts/plan_blif_flow_v2_data_repair.py
Normal file
@@ -0,0 +1,446 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Produce the BLIF Flow v2 historical repair plan (dry-run only).
|
||||
|
||||
This command has no apply mode. Every database read occurs inside an explicit
|
||||
READ ONLY transaction after an exact clientflow_codex_test identity assertion.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections import Counter, defaultdict
|
||||
from datetime import date, datetime, timezone
|
||||
from decimal import Decimal
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from uuid import UUID
|
||||
|
||||
from sqlalchemy import text
|
||||
|
||||
from app.db import engine
|
||||
from app.domain.opportunity_flow.repair import (
|
||||
FOLLOWUP_ACTIONS, TaskRepairContext, classify_pending_task,
|
||||
simulate_high_repairs,
|
||||
)
|
||||
from scripts.simulate_blif_flow_v2 import SIMULATION_AT, collect
|
||||
|
||||
|
||||
OUTPUTS = {
|
||||
"plan": Path("/tmp/blif_flow_v2_data_repair_plan.json"),
|
||||
"tasks": Path("/tmp/blif_flow_v2_task_repair_audit.json"),
|
||||
"opportunities": Path("/tmp/blif_flow_v2_opportunity_repair_audit.json"),
|
||||
"duplicates": Path("/tmp/blif_flow_v2_duplicate_repair_audit.json"),
|
||||
"followups": Path("/tmp/blif_flow_v2_followup_repair_audit.json"),
|
||||
"summary": Path("/tmp/blif_flow_v2_data_repair_summary.txt"),
|
||||
}
|
||||
EXPECTED_DATABASE = "clientflow_codex_test"
|
||||
CURRENT_QUEUES = {"do_now", "review", "exception"}
|
||||
NAMED = {
|
||||
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
|
||||
"PANORAMIC SUCCESS": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
|
||||
"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",
|
||||
}
|
||||
|
||||
|
||||
def _jsonable(value: Any) -> Any:
|
||||
if isinstance(value, (date, datetime)):
|
||||
return value.isoformat()
|
||||
if isinstance(value, Decimal):
|
||||
return float(value)
|
||||
if isinstance(value, UUID):
|
||||
return str(value)
|
||||
if isinstance(value, dict):
|
||||
return {key: _jsonable(item) for key, item in value.items()}
|
||||
if isinstance(value, (list, tuple)):
|
||||
return [_jsonable(item) for item in value]
|
||||
return value
|
||||
|
||||
|
||||
def _read_database() -> dict[str, Any]:
|
||||
with engine.connect() as conn:
|
||||
conn.exec_driver_sql("BEGIN READ ONLY")
|
||||
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: {identity!r}")
|
||||
try:
|
||||
tasks = [dict(row) for row in conn.execute(text("""
|
||||
SELECT t.*, t.id::text AS id, t.opportunity_id::text,
|
||||
t.resolved_by_event_id::text, t.superseded_by_task_id::text
|
||||
FROM tasks t ORDER BY t.created_at, t.id
|
||||
""")).mappings()]
|
||||
projections = [dict(row) for row in 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, derived_at
|
||||
FROM opportunity_flow_state_v2 ORDER BY opportunity_id
|
||||
""")).mappings()]
|
||||
opportunities = [dict(row) for row in conn.execute(text("""
|
||||
SELECT o.*, o.id::text AS id, c.name AS linked_customer_name
|
||||
FROM opportunities o LEFT JOIN customers c ON c.id=o.local_customer_id
|
||||
ORDER BY o.id
|
||||
""")).mappings()]
|
||||
events = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, event_type, task_id::text,
|
||||
action_code, note, payload, created_at
|
||||
FROM opportunity_events ORDER BY created_at
|
||||
""")).mappings()]
|
||||
operation_links = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, system, external_type,
|
||||
external_id, external_name, status, payload, created_at
|
||||
FROM operation_links ORDER BY created_at
|
||||
""")).mappings()]
|
||||
reconciliation = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, source_system, external_type,
|
||||
external_id, document_number, status, suggested_action,
|
||||
confidence, payload, created_at
|
||||
FROM reconciliation_items ORDER BY created_at
|
||||
""")).mappings()]
|
||||
document_links = [dict(row) for row in conn.execute(text("""
|
||||
SELECT l.id::text, l.opportunity_id::text, l.document_id::text,
|
||||
l.relationship, l.source, l.origin_opportunity_id::text,
|
||||
l.destination_opportunity_id::text, l.ended_at,
|
||||
d.document_kind, d.external_id, d.document_number, d.status
|
||||
FROM opportunity_document_links l
|
||||
JOIN commercial_documents d ON d.id=l.document_id
|
||||
ORDER BY l.created_at
|
||||
""")).mappings()]
|
||||
finally:
|
||||
conn.rollback()
|
||||
return {"identity": {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
|
||||
"tasks": tasks, "projections": projections, "opportunities": opportunities,
|
||||
"events": events, "operation_links": operation_links,
|
||||
"reconciliation": reconciliation, "document_links": document_links}
|
||||
|
||||
|
||||
def _refs(record: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
evidence = record.get("evidence", {})
|
||||
refs = []
|
||||
for role in ("latest_relevant_inbound", "latest_relevant_outbound"):
|
||||
event = evidence.get(role)
|
||||
if event:
|
||||
refs.append({"source": "message_or_communication", "role": role,
|
||||
"id": event.get("id"), "at": event.get("at")})
|
||||
for role in ("proforma", "invoice", "payment", "odoo", "reconciliation"):
|
||||
for item in evidence.get(role, []):
|
||||
refs.append({"source": role, "id": item.get("id"),
|
||||
"external_id": item.get("external_id"),
|
||||
"document_number": item.get("document_number"),
|
||||
"status": item.get("status"), "at": item.get("created_at")})
|
||||
return refs
|
||||
|
||||
|
||||
def _task_audit(db: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
|
||||
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
|
||||
persisted = {row["opportunity_id"]: row for row in db["projections"]}
|
||||
pending = [row for row in db["tasks"] if str(row.get("status", "")).lower() == "pending"]
|
||||
audited = []
|
||||
for task in pending:
|
||||
oid = task.get("opportunity_id")
|
||||
action = str(task.get("action_code") or "").upper()
|
||||
flow = persisted.get(oid, {})
|
||||
record = records.get(oid, {})
|
||||
evidence = record.get("evidence", {})
|
||||
created = task.get("created_at")
|
||||
inbound = evidence.get("latest_relevant_inbound")
|
||||
outbound = evidence.get("latest_relevant_outbound")
|
||||
later_in = inbound if inbound and datetime.fromisoformat(inbound["at"]) > created else None
|
||||
later_out = outbound if outbound and datetime.fromisoformat(outbound["at"]) > created else None
|
||||
event_match = next((event for event in db["events"] if event.get("opportunity_id") == oid
|
||||
and event.get("created_at") and created and event["created_at"] > created
|
||||
and str(event.get("event_type") or "").lower() in {"customer_replied", "message_received", "inbound_message"}), None)
|
||||
if later_in and event_match:
|
||||
later_in = {**later_in, "opportunity_event_id": event_match["id"]}
|
||||
proformas, invoices = evidence.get("proforma", []), evidence.get("invoice", [])
|
||||
ctx = TaskRepairContext(
|
||||
task_id=task["id"], opportunity_id=oid, action_code=action,
|
||||
created_at=created, due_at=task.get("due_at"),
|
||||
business_state=flow.get("business_state"), business_next_action=flow.get("business_next_action"),
|
||||
material_process_key=flow.get("material_process_key"),
|
||||
is_duplicate_representation=bool(flow.get("is_duplicate_representation")),
|
||||
canonical_opportunity_id=flow.get("canonical_opportunity_id"),
|
||||
later_inbound_event=later_in, later_outbound_event=later_out,
|
||||
proforma_exists=bool(proformas), proforma_sent=bool(evidence.get("proforma_sent")),
|
||||
payment_confirmed=bool(evidence.get("payment")), invoice_exists=bool(invoices),
|
||||
invoice_sent=False, odoo_order_exists=any(x.get("external_type") == "sale_order" for x in evidence.get("odoo", [])),
|
||||
odoo_order_validated=any(x.get("external_type") == "physical_validation" and x.get("status") == "validated" for x in evidence.get("odoo", [])),
|
||||
terminal=flow.get("business_state") == "COMPLETED",
|
||||
evidence_refs=tuple(_refs(record)),
|
||||
)
|
||||
decision = classify_pending_task(ctx).to_dict()
|
||||
audited.append(_jsonable({
|
||||
"entity_type": "task", "entity_id": task["id"], "task_id": task["id"],
|
||||
"opportunity_id": oid, "material_process_key": flow.get("material_process_key"),
|
||||
"action_code": action, "created_at": created, "due_at": task.get("due_at"),
|
||||
"current_value": {"status": task.get("status"), "resolution_code": task.get("resolution_code")},
|
||||
"proposed_value": {"status": "resolved" if decision["auto_repair_safe"] else "pending",
|
||||
"resolution_code": decision["resolution_code"],
|
||||
"resolved_by_event_id": decision["resolved_by_event_id"],
|
||||
"superseded_by_task_id": decision["superseded_by_task_id"]},
|
||||
"v1_relevance": "standalone_preserved" if not oid else "derived_historical_obligation",
|
||||
"v2_factual_state": flow.get("business_state"),
|
||||
"v2_current_action": flow.get("business_next_action"),
|
||||
"classification": decision["classification"], "repair_category": decision["classification"],
|
||||
"repair_reason": decision["reason"], "factual_evidence_refs": _refs(record),
|
||||
"confidence": decision["confidence"], "safety_tier": decision["safety_tier"],
|
||||
"auto_repair_safe": decision["auto_repair_safe"],
|
||||
"human_review_required": decision["human_review_required"],
|
||||
}))
|
||||
classifications = Counter(row["classification"] for row in audited)
|
||||
actions = Counter(row["action_code"] for row in audited)
|
||||
combined = Counter(f"{row['classification']} + {row['action_code']}" for row in audited)
|
||||
return {"generated_at": datetime.now(timezone.utc), "database": db["identity"],
|
||||
"total_tasks": len(db["tasks"]), "pending_tasks_audited": len(audited),
|
||||
"counts_by_classification": dict(sorted(classifications.items())),
|
||||
"counts_by_action_code": dict(sorted(actions.items())),
|
||||
"counts_by_classification_and_action_code": dict(sorted(combined.items())),
|
||||
"tasks": audited}
|
||||
|
||||
|
||||
def _opportunity_audit(db: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
|
||||
flow = {row["opportunity_id"]: row for row in db["projections"]}
|
||||
runtime = {row["opportunity_id"]: row for row in projection["opportunities"]}
|
||||
rows = []
|
||||
for opp in db["opportunities"]:
|
||||
oid, state = opp["id"], flow[opp["id"]]
|
||||
current = runtime[oid]["v1"]
|
||||
mismatches = []
|
||||
stage = str(opp.get("stage") or "")
|
||||
if stage.upper() != state["business_state"]:
|
||||
if stage.upper() in {"INFO_SENT", "QUOTE_SENT", "INVOICE_REQUESTED", "INVOICE_SENT", "WON", "SHIPPED"}:
|
||||
category, disposition = "LEGACY_COMPATIBILITY_ONLY", "continue_as_compatibility_only_then_deprecate"
|
||||
elif state["confidence"] == "high":
|
||||
category, disposition = "STALE_DERIVED_STATE", "one_time_repair_after_review"
|
||||
else:
|
||||
category, disposition = "DO_NOT_REPAIR_YET", "requires_review"
|
||||
mismatches.append({"field": "stage", "current": stage, "proposed": state["business_state"],
|
||||
"classification": category, "disposition": disposition})
|
||||
expected_lifecycle = "awaiting_customer" if state["business_state"] in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"} else "active"
|
||||
if str(opp.get("lifecycle_state") or "active") != expected_lifecycle:
|
||||
mismatches.append({"field": "lifecycle_state", "current": opp.get("lifecycle_state"),
|
||||
"proposed": expected_lifecycle, "classification": "REQUIRES_MIGRATION",
|
||||
"disposition": "rebuild_from_flow_v2_and_valid_followups"})
|
||||
if opp.get("next_follow_up_at") and current.get("operational_queue") not in {"waiting", "do_now"}:
|
||||
mismatches.append({"field": "next_follow_up_at", "current": _jsonable(opp.get("next_follow_up_at")),
|
||||
"proposed": None, "classification": "DO_NOT_REPAIR_YET",
|
||||
"disposition": "audit_followup_before_one_time_repair"})
|
||||
if current.get("current_action") != state.get("business_next_action"):
|
||||
mismatches.append({"field": "current_action_compatibility", "current": current.get("current_action"),
|
||||
"proposed": state.get("business_next_action"), "classification": "PRESENTATION_ONLY",
|
||||
"disposition": "render_from_safe_flow_v2_eventually"})
|
||||
rows.append({"entity_type": "opportunity", "entity_id": oid, "opportunity_id": oid,
|
||||
"material_process_key": state["material_process_key"], "title": opp.get("title"),
|
||||
"mismatches": mismatches, "repair_category": "NO_MISMATCH" if not mismatches else mismatches[0]["classification"],
|
||||
"factual_evidence_refs": state.get("evidence_refs", []), "confidence": state["confidence"],
|
||||
"safety_tier": "LOW", "auto_repair_safe": False, "human_review_required": bool(mismatches)})
|
||||
counts = Counter(item["classification"] for row in rows for item in row["mismatches"])
|
||||
for category in ("PRESENTATION_ONLY", "STALE_DERIVED_STATE", "FACTUAL_CONTRADICTION",
|
||||
"LEGACY_COMPATIBILITY_ONLY", "REQUIRES_MIGRATION", "DO_NOT_REPAIR_YET"):
|
||||
counts.setdefault(category, 0)
|
||||
return {"counts": dict(sorted(counts.items())), "opportunities": _jsonable(rows)}
|
||||
|
||||
|
||||
def _duplicate_audit(db: dict[str, Any], task_audit: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
|
||||
groups = defaultdict(list)
|
||||
for row in db["projections"]:
|
||||
groups[row["material_process_key"]].append(row)
|
||||
task_by_opp = defaultdict(list)
|
||||
for row in task_audit["tasks"]:
|
||||
task_by_opp[row["opportunity_id"]].append(row)
|
||||
opportunities = {row["id"]: row for row in db["opportunities"]}
|
||||
runtime = {row["opportunity_id"]: row for row in projection["opportunities"]}
|
||||
results = []
|
||||
for key, members in groups.items():
|
||||
duplicates = [row for row in members if row["is_duplicate_representation"]]
|
||||
if not duplicates:
|
||||
continue
|
||||
canonical = next(row for row in members if not row["is_duplicate_representation"])
|
||||
for duplicate in duplicates:
|
||||
oid = duplicate["opportunity_id"]
|
||||
results.append({
|
||||
"entity_type": "duplicate_opportunity_representation", "entity_id": oid,
|
||||
"opportunity_id": oid,
|
||||
"material_process_key": key, "canonical_opportunity_id": canonical["opportunity_id"],
|
||||
"duplicate_opportunity_id": oid,
|
||||
"canonical_factual_evidence": canonical.get("evidence_refs", []),
|
||||
"shared_identity_evidence": [key], "duplicate_specific_tasks": task_by_opp[oid],
|
||||
"duplicate_specific_operation_links": [x for x in db["operation_links"] if x.get("opportunity_id") == oid],
|
||||
"duplicate_specific_reconciliation_rows": [x for x in db["reconciliation"] if x.get("opportunity_id") == oid],
|
||||
"duplicate_specific_document_links": [x for x in db["document_links"] if x.get("opportunity_id") == oid],
|
||||
"duplicate_specific_work_item": {"v1": runtime[oid]["v1"], "safe_v2": runtime[oid]["safe_v2"]},
|
||||
"synthetic_mapping_metadata": opportunities[oid].get("metadata"),
|
||||
"current_value": {"business_state": duplicate["business_state"], "is_duplicate_representation": True,
|
||||
"stage": opportunities[oid].get("stage"),
|
||||
"lifecycle_state": opportunities[oid].get("lifecycle_state")},
|
||||
"proposed_value": {"operational_visibility": "suppressed", "queue": "not_current"},
|
||||
"repair_category": "DUPLICATE", "repair_reason": "Suppress duplicate operational representation; preserve all factual evidence and the opportunity row.",
|
||||
"factual_evidence_refs": canonical.get("evidence_refs", []),
|
||||
"confidence": "high", "safety_tier": "HIGH", "auto_repair_safe": True,
|
||||
"human_review_required": False,
|
||||
})
|
||||
return {"material_groups_found": len(results), "duplicate_representations": len(results),
|
||||
"groups": _jsonable(results)}
|
||||
|
||||
|
||||
def _followup_audit(task_audit: dict[str, Any], db: dict[str, Any]) -> dict[str, Any]:
|
||||
opportunity = {row["id"]: row for row in db["opportunities"]}
|
||||
rows = []
|
||||
for task in task_audit["tasks"]:
|
||||
if task["action_code"] not in FOLLOWUP_ACTIONS and not (
|
||||
task["opportunity_id"] and opportunity[task["opportunity_id"]].get("next_follow_up_at")
|
||||
):
|
||||
continue
|
||||
due = datetime.fromisoformat(task["due_at"]) if task.get("due_at") else None
|
||||
if task["action_code"] == "CALL_CUSTOMER" and task["classification"] == "VALID_CURRENT":
|
||||
category = "AUTHORITATIVE_CALL_CUSTOMER"
|
||||
elif task["classification"] == "SATISFIED_BY_EVENT":
|
||||
category = "SATISFIED_FOLLOWUP"
|
||||
elif task["classification"] == "VALID_CURRENT" and task["action_code"] == "FOLLOW_UP_PAYMENT":
|
||||
category = "VALID_PAYMENT_FOLLOWUP"
|
||||
elif task["classification"] == "VALID_CURRENT":
|
||||
category = "VALID_CUSTOMER_FOLLOWUP"
|
||||
elif task["classification"] == "AMBIGUOUS":
|
||||
category = "AMBIGUOUS"
|
||||
else:
|
||||
category = "OBSOLETE_COMPATIBILITY_MIRROR"
|
||||
rows.append({**task, "followup_classification": category,
|
||||
"timing": "future" if due and due > SIMULATION_AT else "overdue_or_due" if due else "unscheduled"})
|
||||
represented = {row.get("opportunity_id") for row in rows}
|
||||
for oid, opp in opportunity.items():
|
||||
timestamp = opp.get("next_follow_up_at")
|
||||
if not timestamp or oid in represented:
|
||||
continue
|
||||
rows.append({
|
||||
"entity_type": "opportunity_followup_compatibility", "entity_id": oid,
|
||||
"opportunity_id": oid, "action_code": None, "due_at": _jsonable(timestamp),
|
||||
"current_value": {"next_follow_up_at": _jsonable(timestamp),
|
||||
"lifecycle_state": opp.get("lifecycle_state")},
|
||||
"proposed_value": None,
|
||||
"followup_classification": "HISTORICAL_TIMESTAMP_NO_ACTIVE_OBLIGATION",
|
||||
"timing": "future" if timestamp > SIMULATION_AT else "overdue_or_due",
|
||||
"repair_reason": "Compatibility timestamp has no pending follow-up task; do not clear without migration review.",
|
||||
"confidence": "low", "safety_tier": "LOW", "auto_repair_safe": False,
|
||||
"human_review_required": True, "factual_evidence_refs": [],
|
||||
})
|
||||
counts = Counter(row["followup_classification"] for row in rows)
|
||||
for category in ("AUTHORITATIVE_CALL_CUSTOMER", "VALID_CUSTOMER_FOLLOWUP",
|
||||
"VALID_PAYMENT_FOLLOWUP", "SATISFIED_FOLLOWUP",
|
||||
"OBSOLETE_COMPATIBILITY_MIRROR", "HISTORICAL_TIMESTAMP_NO_ACTIVE_OBLIGATION",
|
||||
"AMBIGUOUS"):
|
||||
counts.setdefault(category, 0)
|
||||
return {"counts": dict(sorted(counts.items())), "followups": rows}
|
||||
|
||||
|
||||
def _queue_after(projection: dict[str, Any], duplicate_audit: dict[str, Any]) -> dict[str, int]:
|
||||
duplicate_ids = {row["duplicate_opportunity_id"] for row in duplicate_audit["groups"]}
|
||||
counts = Counter()
|
||||
for row in projection["opportunities"] + projection["standalone_canonical_items"]:
|
||||
queue = row["safe_v2"]["effective_operational_queue"]
|
||||
if row.get("opportunity_id") in duplicate_ids:
|
||||
queue = "not_current"
|
||||
counts[queue] += 1
|
||||
return {"current_work": sum(counts[x] for x in CURRENT_QUEUES), "do_now": counts["do_now"],
|
||||
"review": counts["review"], "waiting": counts["waiting"], "backlog": counts["backlog"]}
|
||||
|
||||
|
||||
def _named_cases(task_audit: dict[str, Any], db: dict[str, Any], projection: dict[str, Any]) -> dict[str, Any]:
|
||||
tasks = defaultdict(list)
|
||||
for row in task_audit["tasks"]:
|
||||
tasks[row["opportunity_id"]].append(row)
|
||||
projections = {row["opportunity_id"]: row for row in db["projections"]}
|
||||
runtime = {row["opportunity_id"]: row for row in projection["opportunities"]}
|
||||
result = {}
|
||||
for name, oid in NAMED.items():
|
||||
result[name] = {"opportunity_id": oid, "business_state": projections[oid]["business_state"],
|
||||
"business_next_action": projections[oid]["business_next_action"],
|
||||
"pending_tasks": tasks[oid], "simulated_effective_action": runtime[oid]["safe_v2"]["effective_operational_action"],
|
||||
"simulated_queue": runtime[oid]["safe_v2"]["effective_operational_queue"]}
|
||||
for label in ("ENGEXICON", "CONSTRURECUP"):
|
||||
matches = [row for row in projection["opportunities"] if label.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()]
|
||||
result[label] = [{"opportunity_id": row["opportunity_id"], "business_state": row["safe_v2"]["business_state"],
|
||||
"simulated_effective_action": row["safe_v2"]["effective_operational_action"],
|
||||
"simulated_queue": row["safe_v2"]["effective_operational_queue"]} for row in matches]
|
||||
return result
|
||||
|
||||
|
||||
def build_plan() -> dict[str, Any]:
|
||||
db = _read_database()
|
||||
if len(db["projections"]) != 328:
|
||||
raise RuntimeError(f"expected 328 persisted projections, found {len(db['projections'])}")
|
||||
# collect() begins its own READ ONLY transaction and repeats the exact DB/user guard.
|
||||
projection = collect(expected_database=EXPECTED_DATABASE, expected_user=db["identity"]["user"], require_read_only=False)
|
||||
tasks = _task_audit(db, projection)
|
||||
opportunities = _opportunity_audit(db, projection)
|
||||
duplicates = _duplicate_audit(db, tasks, projection)
|
||||
followups = _followup_audit(tasks, db)
|
||||
simulation = simulate_high_repairs(tasks["tasks"])
|
||||
before = {"total_tasks": tasks["total_tasks"], "pending_tasks": tasks["pending_tasks_audited"],
|
||||
"pending_task_classifications": tasks["counts_by_classification"],
|
||||
"duplicate_material_groups": duplicates["material_groups_found"],
|
||||
"duplicate_current_cards": sum(1 for row in duplicates["groups"] if any(
|
||||
task["classification"] == "DUPLICATE" for task in row["duplicate_specific_tasks"])),
|
||||
"v1": projection["v1_totals"], "safe_v2": projection["safe_v2_totals"]}
|
||||
after = {"pending_tasks": simulation["pending_after"],
|
||||
"resolved_as_satisfied": simulation["removed_by_classification"].get("SATISFIED_BY_EVENT", 0),
|
||||
"resolved_as_superseded": simulation["removed_by_classification"].get("SUPERSEDED", 0),
|
||||
"resolved_as_duplicate": simulation["removed_by_classification"].get("DUPLICATE", 0),
|
||||
"resolved_as_premature": simulation["removed_by_classification"].get("PREMATURE", 0),
|
||||
"duplicate_current_cards": 0, "safe_v2": _queue_after(projection, duplicates)}
|
||||
false_negative_gate = []
|
||||
for row in simulation["removed"]:
|
||||
disposition = "DUPLICATE_SUPPRESSED" if row["classification"] == "DUPLICATE" else (
|
||||
"REPLACED_BY_CORRECT_ACTION" if row["v2_current_action"] else "SAFE_TO_REMOVE")
|
||||
false_negative_gate.append({"entity": row["task_id"], "current_action": row["action_code"],
|
||||
"opportunity_id": row["opportunity_id"], "reason": row["repair_reason"],
|
||||
"factual_evidence": row["factual_evidence_refs"],
|
||||
"replacement_obligation": row["v2_current_action"], "classification": disposition})
|
||||
named = _named_cases(tasks, db, projection)
|
||||
unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in false_negative_gate)
|
||||
plan = {"phase": 1, "mode": "dry-run", "database": db["identity"],
|
||||
"generated_at": datetime.now(timezone.utc), "before": before,
|
||||
"simulated_after_high_confidence_repair": after,
|
||||
"false_negative_safety_gate": false_negative_gate,
|
||||
"unsafe_false_negatives": unsafe, "automatic_repair_recommended": unsafe == 0,
|
||||
"named_cases": named,
|
||||
"repairs": [row for row in tasks["tasks"] if row["classification"] != "VALID_CURRENT"] + duplicates["groups"]}
|
||||
for key, value in (("tasks", tasks), ("opportunities", opportunities),
|
||||
("duplicates", duplicates), ("followups", followups), ("plan", plan)):
|
||||
OUTPUTS[key].write_text(json.dumps(_jsonable(value), ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
return {"plan": _jsonable(plan), "tasks": tasks, "opportunities": opportunities,
|
||||
"duplicates": duplicates, "followups": followups}
|
||||
|
||||
|
||||
def _summary(result: dict[str, Any]) -> str:
|
||||
plan, tasks = result["plan"], result["tasks"]
|
||||
lines = ["BLIF FLOW V2 HISTORICAL DATA-REPAIR PLAN — DRY RUN", "",
|
||||
f"Database: {plan['database']}", f"Total tasks: {tasks['total_tasks']}",
|
||||
f"Pending tasks audited: {tasks['pending_tasks_audited']}", "", "TASK AUDIT"]
|
||||
for name in ("VALID_CURRENT", "SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE", "PREMATURE", "AMBIGUOUS"):
|
||||
lines.append(f"{name}: {tasks['counts_by_classification'].get(name, 0)}")
|
||||
tiers = Counter(row["safety_tier"] for row in tasks["tasks"] if row["classification"] != "VALID_CURRENT")
|
||||
lines += ["", "SAFETY", f"HIGH repairs: {tiers['HIGH']}", f"MEDIUM repairs: {tiers['MEDIUM']}",
|
||||
f"LOW repairs: {tiers['LOW']}", f"Unsafe false negatives: {plan['unsafe_false_negatives']}", "", "BEFORE",
|
||||
json.dumps(plan["before"], ensure_ascii=False, sort_keys=True), "", "SIMULATED AFTER HIGH",
|
||||
json.dumps(plan["simulated_after_high_confidence_repair"], ensure_ascii=False, sort_keys=True), "", "DUPLICATES",
|
||||
json.dumps({k: result['duplicates'][k] for k in ('material_groups_found','duplicate_representations')}, sort_keys=True), "", "OPPORTUNITY STATE",
|
||||
json.dumps(result["opportunities"]["counts"], sort_keys=True), "", "FOLLOWUPS",
|
||||
json.dumps(result["followups"]["counts"], sort_keys=True), "", "NAMED CASES"]
|
||||
for name, row in plan["named_cases"].items():
|
||||
lines.append(f"{name}: {json.dumps(row, ensure_ascii=False, sort_keys=True)}")
|
||||
lines += ["", "FILES CREATED"] + [str(path) for path in OUTPUTS.values()]
|
||||
return "\n".join(lines) + "\n"
|
||||
|
||||
|
||||
def main() -> None:
|
||||
result = build_plan()
|
||||
summary = _summary(result)
|
||||
OUTPUTS["summary"].write_text(summary, encoding="utf-8")
|
||||
print(summary, end="")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
84
scripts/rebuild_blif_flow_v2_projection.py
Normal file
84
scripts/rebuild_blif_flow_v2_projection.py
Normal file
@@ -0,0 +1,84 @@
|
||||
#!/usr/bin/env python3
|
||||
"""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__":
|
||||
raise SystemExit(main())
|
||||
877
scripts/simulate_blif_flow_v2.py
Normal file
877
scripts/simulate_blif_flow_v2.py
Normal file
@@ -0,0 +1,877 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Read-only BLIF Flow v2 projection against the isolated shadow snapshot."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
from collections import Counter, defaultdict
|
||||
from dataclasses import replace
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Iterable
|
||||
|
||||
from sqlalchemy import text
|
||||
|
||||
from app.db import engine
|
||||
from app.domain.opportunity_flow.v2 import (
|
||||
EffectiveOperationalDecision,
|
||||
derive_business_facts, derive_effective_operational_action,
|
||||
derive_safe_operational_action, derive_v2_operational_queue,
|
||||
suppress_duplicate_representation,
|
||||
)
|
||||
from app.operations_service import get_operations_summary
|
||||
from app.opportunity_next_action_service import get_opportunity_next_actions
|
||||
from app.opportunity_service import list_opportunities
|
||||
|
||||
|
||||
PROJECTION = Path("/tmp/blif_flow_v2_projection.json")
|
||||
COMPARISON = Path("/tmp/blif_flow_v2_operations_comparison.txt")
|
||||
AMBIGUOUS = Path("/tmp/blif_flow_v2_ambiguous_cases.json")
|
||||
PROMOTIONS = Path("/tmp/blif_flow_v2_promotions_audit.json")
|
||||
INVOICE_WITHOUT_PAYMENT = Path("/tmp/blif_flow_v2_invoice_without_payment.json")
|
||||
REVIEW_AUDIT = Path("/tmp/blif_flow_v2_review_audit.json")
|
||||
BACKLOG_DELTA = Path("/tmp/blif_flow_v2_backlog_delta.json")
|
||||
CURRENT_DELTA = Path("/tmp/blif_flow_v2_current_delta.json")
|
||||
MATERIAL_IDENTITY = Path("/tmp/blif_flow_v2_material_identity.json")
|
||||
TERMINAL = {"WON", "LOST", "NO_INTEREST", "ARCHIVED", "COMPLETED", "CLOSED"}
|
||||
ORDER_INTENT = re.compile(r"\b(quero|queremos|pretendo|pretendemos|aceito|aceitamos|adjudic|encomendar|encomenda|avançar|avancar|proceder)\b", re.I)
|
||||
ORDER_CHANGE_VERB = re.compile(r"\b(alterar|alteração|alteracao|mudar|mudança|mudanca|trocar|substituir|corrigir|retificar)\b", re.I)
|
||||
ORDER_CHANGE_SUBJECT = re.compile(r"\b(produto|modelo|quantidade|morada|entrega|nif|fiscal|faturação|faturacao|condições|condicoes)\b", re.I)
|
||||
QUOTE_REQUEST = re.compile(r"\b(preço|preco|orçamento|orcamento|cotação|cotacao|proposta|quote)\b", re.I)
|
||||
PAYMENT_PROOF = re.compile(r"\b(comprovativo|transferência|transferencia|pagamento efetuado|pago|liquidado)\b", re.I)
|
||||
SIMULATION_AT = datetime(2026, 8, 15, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def _s(value: Any) -> str:
|
||||
return str(value or "").strip()
|
||||
|
||||
|
||||
def _dt(value: Any) -> datetime | None:
|
||||
if isinstance(value, datetime):
|
||||
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
|
||||
if not value:
|
||||
return None
|
||||
try:
|
||||
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
|
||||
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
|
||||
except ValueError:
|
||||
return None
|
||||
|
||||
|
||||
def _jsonable(value: Any) -> Any:
|
||||
if isinstance(value, datetime):
|
||||
return value.isoformat()
|
||||
if isinstance(value, dict):
|
||||
return {key: _jsonable(item) for key, item in value.items()}
|
||||
if isinstance(value, (list, tuple)):
|
||||
return [_jsonable(item) for item in value]
|
||||
return value
|
||||
|
||||
|
||||
def _compact(value: Any, limit: int = 260) -> str:
|
||||
result = re.sub(r"\s+", " ", _s(value))
|
||||
return result if len(result) <= limit else result[: limit - 1].rstrip() + "…"
|
||||
|
||||
|
||||
def _payload(value: Any) -> dict[str, Any]:
|
||||
if isinstance(value, dict):
|
||||
return value
|
||||
if isinstance(value, str) and value.strip():
|
||||
try:
|
||||
parsed = json.loads(value)
|
||||
return parsed if isinstance(parsed, dict) else {}
|
||||
except ValueError:
|
||||
pass
|
||||
return {}
|
||||
|
||||
|
||||
def _group(rows: Iterable[dict[str, Any]], key: str = "opportunity_id") -> dict[str, list[dict[str, Any]]]:
|
||||
result: dict[str, list[dict[str, Any]]] = defaultdict(list)
|
||||
for row in rows:
|
||||
result[_s(row.get(key))].append(dict(row))
|
||||
return result
|
||||
|
||||
|
||||
def _load(
|
||||
*, expected_database: str = "clientflow_codex_shadow",
|
||||
expected_user: str | None = "clientflow_codex",
|
||||
require_read_only: bool = True,
|
||||
) -> 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()
|
||||
if identity[0] != expected_database or (expected_user and identity[1] != expected_user):
|
||||
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"))
|
||||
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,
|
||||
c.street_name, c.postal_zone, c.city_name, c.phone AS fiscal_phone
|
||||
FROM opportunities o LEFT JOIN customers c ON c.id=o.local_customer_id
|
||||
ORDER BY o.created_at, o.id
|
||||
""")).mappings()]
|
||||
tasks = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, action_code, action, note, status,
|
||||
due_at, created_at, done_at, metadata
|
||||
FROM tasks WHERE opportunity_id IS NOT NULL ORDER BY created_at
|
||||
""")).mappings()]
|
||||
messages = [dict(row) for row in conn.execute(text("""
|
||||
SELECT o.id::text AS opportunity_id, m.id::text, m.direction,
|
||||
COALESCE(m.clean_body,m.raw_body,'') AS body, m.created_at,
|
||||
m.source_system, m.metadata
|
||||
FROM opportunities o JOIN messages m ON m.conversation_id=o.conversation_id
|
||||
WHERE m.source_system IN ('chatwoot','chatwoot_backfill')
|
||||
ORDER BY m.created_at
|
||||
""")).mappings()]
|
||||
communications = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, direction, classification, subject,
|
||||
body, status, created_at, metadata
|
||||
FROM communications WHERE opportunity_id IS NOT NULL ORDER BY created_at
|
||||
""")).mappings()]
|
||||
documents = [dict(row) for row in conn.execute(text("""
|
||||
SELECT DISTINCT ON (COALESCE(l.opportunity_id,d.opportunity_id),d.id)
|
||||
COALESCE(l.opportunity_id,d.opportunity_id)::text AS opportunity_id,
|
||||
d.id::text, d.document_kind, d.document_type, d.document_number,
|
||||
d.external_id, d.status, d.payload, d.created_at, d.updated_at,
|
||||
COALESCE(l.relationship, CASE WHEN d.is_primary THEN 'PRIMARY' ELSE d.role END, 'PRIMARY') AS relationship,
|
||||
l.ended_at
|
||||
FROM commercial_documents d
|
||||
LEFT JOIN opportunity_document_links l ON l.document_id=d.id AND l.ended_at IS NULL
|
||||
WHERE COALESCE(l.opportunity_id,d.opportunity_id) IS NOT NULL
|
||||
ORDER BY COALESCE(l.opportunity_id,d.opportunity_id),d.id,l.updated_at DESC NULLS LAST
|
||||
""")).mappings()]
|
||||
links = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, system, external_type, external_id,
|
||||
external_name, status, payload, created_at, updated_at, last_synced_at
|
||||
FROM operation_links ORDER BY created_at
|
||||
""")).mappings()]
|
||||
reconciliation = [dict(row) for row in conn.execute(text("""
|
||||
SELECT id::text, opportunity_id::text, title, description, document_number,
|
||||
status, suggested_action, confidence, payload, created_at
|
||||
FROM reconciliation_items
|
||||
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]},
|
||||
"opportunities": opportunities, "tasks": _group(tasks), "messages": _group(messages),
|
||||
"communications": _group(communications), "documents": _group(documents),
|
||||
"links": _group(links), "reconciliation": _group(reconciliation),
|
||||
"all_reconciliation": reconciliation,
|
||||
}
|
||||
|
||||
|
||||
def _event(row: dict[str, Any]) -> dict[str, Any]:
|
||||
return {
|
||||
"id": row.get("id"), "at": _jsonable(row.get("created_at")),
|
||||
"direction": row.get("direction"), "classification": row.get("classification"),
|
||||
"text": _compact(row.get("body") or row.get("subject")),
|
||||
}
|
||||
|
||||
|
||||
def _derive_record(opp: dict[str, Any], data: dict[str, Any], v1: dict[str, Any], v1_item: dict[str, Any] | None) -> dict[str, Any]:
|
||||
oid = _s(opp["id"])
|
||||
messages = data["messages"].get(oid, [])
|
||||
comms = data["communications"].get(oid, [])
|
||||
tasks = data["tasks"].get(oid, [])
|
||||
docs = data["documents"].get(oid, [])
|
||||
links = data["links"].get(oid, [])
|
||||
recons = data["reconciliation"].get(oid, [])
|
||||
events = sorted(messages + comms, key=lambda row: _dt(row.get("created_at")) or datetime.min.replace(tzinfo=timezone.utc))
|
||||
inbound = [row for row in events if _s(row.get("direction")).lower() == "inbound"]
|
||||
outbound = [row for row in events if _s(row.get("direction")).lower() == "outbound"]
|
||||
latest_in, latest_out = (inbound[-1] if inbound else None), (outbound[-1] if outbound else None)
|
||||
inbound_text = "\n".join(_s(row.get("body") or row.get("subject")) for row in inbound)
|
||||
inbound_classes = {_s(row.get("classification")).upper() for row in inbound}
|
||||
request_kind = "quote" if inbound_classes & {"SEND_QUOTE", "SEND_PROFORMA"} or QUOTE_REQUEST.search(inbound_text) else "info"
|
||||
order_rows = [row for row in inbound if _s(row.get("classification")).upper() in {"SEND_PROFORMA", "CONFIRM_PAYMENT", "SEND_INVOICE"}
|
||||
or ORDER_INTENT.search(_s(row.get("body") or row.get("subject")))]
|
||||
order_intent_at = _dt(order_rows[-1].get("created_at")) if order_rows else None
|
||||
|
||||
current_docs = [row for row in docs if _s(row.get("relationship")).upper() == "PRIMARY"
|
||||
and _s(row.get("status")).lower() not in {"cancelled", "canceled", "failed"}]
|
||||
proformas = [row for row in current_docs if _s(row.get("document_kind")).lower() in {"quotation", "quote", "proforma"}
|
||||
and _s(row.get("status")).lower() != "converted"]
|
||||
invoices = [row for row in current_docs if _s(row.get("document_kind")).lower() == "invoice"]
|
||||
proforma = proformas[-1] if proformas else None
|
||||
invoice = invoices[-1] if invoices else None
|
||||
doc_number = _s((proforma or {}).get("document_number") or (proforma or {}).get("external_id"))
|
||||
proforma_payload = _payload((proforma or {}).get("payload"))
|
||||
sent_outbound = next((row for row in reversed(outbound) if doc_number and doc_number.casefold() in _s(row.get("body") or row.get("subject")).casefold()), None)
|
||||
proforma_sent = bool(proforma and (
|
||||
(proforma or {}).get("sent_at") or proforma_payload.get("sent_at") or proforma_payload.get("clientflow_sent_evidence")
|
||||
or _s((proforma or {}).get("status")).lower() in {"sent", "issued_sent"} or sent_outbound
|
||||
))
|
||||
|
||||
payment_links = [row for row in links if row.get("system") == "clientflow" and row.get("external_type") == "payment"]
|
||||
payment = next((row for row in reversed(payment_links) if _s(row.get("status")).lower() == "confirmed"), None)
|
||||
payment_proof_rows = [row for row in inbound if _s(row.get("classification")).upper() == "CONFIRM_PAYMENT"
|
||||
or PAYMENT_PROOF.search(_s(row.get("body") or row.get("subject")))]
|
||||
odoo_sales = [row for row in links if row.get("system") == "odoo" and row.get("external_type") == "sale_order"
|
||||
and _s(row.get("status")).lower() not in {"not_found", "no_order", "cancelled"}]
|
||||
validation = [row for row in links if row.get("system") == "odoo" and row.get("external_type") in {"physical_validation", "physical_status"}
|
||||
and _s(row.get("status")).lower() in {"validated", "ready_to_ship", "shipped", "done", "delivered"}]
|
||||
fulfilled = any(row.get("external_type") in {"physical_status", "delivery"} and _s(row.get("status")).lower() in {"shipped", "done", "delivered"} for row in links)
|
||||
|
||||
change_rows = [
|
||||
row for row in inbound
|
||||
if ORDER_CHANGE_VERB.search(_s(row.get("body") or row.get("subject")))
|
||||
and ORDER_CHANGE_SUBJECT.search(_s(row.get("body") or row.get("subject")))
|
||||
]
|
||||
change_at = _dt(change_rows[-1].get("created_at")) if change_rows else None
|
||||
proforma_at = _dt((proforma or {}).get("created_at"))
|
||||
material_change = bool(change_at and proforma_at and change_at > proforma_at)
|
||||
fiscal_complete = bool(opp.get("local_customer_id") and opp.get("tax_id") and opp.get("fiscal_email")
|
||||
and opp.get("street_name") and opp.get("postal_zone") and opp.get("city_name"))
|
||||
fiscal_conflict = bool((_payload(opp.get("metadata")).get("fiscal_conflict") or _payload(opp.get("metadata")).get("has_nif_conflict")))
|
||||
blockers = []
|
||||
if fiscal_conflict:
|
||||
blockers.append("Conflicting fiscal/NIF evidence.")
|
||||
conflict_recons = [row for row in recons if _s(row.get("status")).lower() in {"needs_review", "conflict"}]
|
||||
if conflict_recons:
|
||||
blockers.append("Unresolved document reconciliation conflict.")
|
||||
reconstructed = _s(_payload(opp.get("metadata")).get("clientflow_record_mode")) in {
|
||||
"reconstructed_invoice_review", "historical_reconstructed", "legacy_review"
|
||||
} or any(marker in _s(opp.get("title")).casefold() for marker in ("processo reconstruído", "sem oportunidade"))
|
||||
if reconstructed and not (invoice or payment or odoo_sales):
|
||||
blockers.append("Reconstructed process lacks corroborating structured evidence.")
|
||||
if invoice and not payment:
|
||||
blockers.append("Structured invoice exists without confirmed payment evidence; correction/reconstruction flow is unspecified.")
|
||||
if odoo_sales and (not payment or not invoice):
|
||||
blockers.append("Odoo execution evidence exists without the mandatory linked payment and invoice evidence.")
|
||||
|
||||
status = _s(opp.get("status")).upper()
|
||||
stage = _s(opp.get("stage")).upper()
|
||||
lost = status in {"LOST", "NO_INTEREST"} or stage in {"LOST", "NO_INTEREST", "ARCHIVED"}
|
||||
info_sent = bool(latest_out and (not latest_in or _dt(latest_out.get("created_at")) >= _dt(latest_in.get("created_at"))))
|
||||
followup_satisfied = any(
|
||||
_s(task.get("status")).lower() == "pending" and _s(task.get("action_code")).upper().startswith("FOLLOW_UP_")
|
||||
and latest_in and _dt(latest_in.get("created_at")) > (_dt(task.get("created_at")) or datetime.max.replace(tzinfo=timezone.utc))
|
||||
for task in tasks
|
||||
)
|
||||
sparse = not events and not current_docs and not links
|
||||
review_required = (
|
||||
bool(material_change and (payment or invoice)) or (reconstructed and sparse)
|
||||
or bool(invoice and not payment) or bool(odoo_sales and (not payment or not invoice))
|
||||
)
|
||||
facts = derive_business_facts(
|
||||
opportunity_id=oid, terminal=status in TERMINAL or stage in TERMINAL, explicitly_lost=lost,
|
||||
review_required=review_required, fiscal_blocked=fiscal_conflict,
|
||||
document_reconciliation_required=False, customer_request=bool(inbound),
|
||||
request_kind=request_kind, latest_relevant_inbound_at=_dt((latest_in or {}).get("created_at")),
|
||||
latest_relevant_outbound_at=_dt((latest_out or {}).get("created_at")), info_or_offer_sent=info_sent,
|
||||
order_intent=bool(order_rows), order_intent_at=order_intent_at, fiscal_identity_evidence=fiscal_complete,
|
||||
proforma_exists=bool(proforma), proforma_sent=proforma_sent, proforma_created_at=proforma_at,
|
||||
proforma_sent_at=_dt((sent_outbound or {}).get("created_at")), potential_payment_evidence=bool(payment_proof_rows and not payment),
|
||||
payment_confirmed=bool(payment), payment_confirmed_at=_dt((payment or {}).get("created_at")),
|
||||
invoice_exists=bool(invoice), invoice_created_at=_dt((invoice or {}).get("created_at")),
|
||||
odoo_order_exists=bool(odoo_sales), odoo_order_validated=bool(validation), fulfillment_complete=fulfilled,
|
||||
material_order_change=material_change, material_order_change_at=change_at,
|
||||
later_customer_inbound_satisfies_followup=followup_satisfied, blockers=blockers,
|
||||
audit_task_codes=[f"{task.get('action_code')}:{task.get('status')}" for task in tasks],
|
||||
)
|
||||
decision = derive_v2_operational_queue(facts)
|
||||
confidence = decision.confidence
|
||||
ambiguity = []
|
||||
if sparse and not lost:
|
||||
confidence = "low"
|
||||
ambiguity.append("No message, structured document, payment, or Odoo evidence is linked.")
|
||||
if stage in {"QUOTE_SENT", "PROFORMA_SENT", "WAITING_PAYMENT"} and not proforma:
|
||||
confidence = "low"
|
||||
ambiguity.append("V1 stage suggests a formal offer, but no current structured proforma is linked.")
|
||||
if stage == "PAYMENT_CONFIRMED" and not payment:
|
||||
confidence = "low"
|
||||
ambiguity.append("V1 stage says payment confirmed, but no confirmed payment operation link exists.")
|
||||
if stage in {"WON", "SHIPPED", "ODOO_ORDER_CREATED", "IN_PRODUCTION"} and not odoo_sales:
|
||||
confidence = "low"
|
||||
ambiguity.append("V1 stage implies execution, but no Odoo sale-order link exists.")
|
||||
decision = replace(decision, confidence=confidence)
|
||||
v1_action = _s((v1_item or {}).get("current_action_code") or v1.get("action_code")) or None
|
||||
v1_queue = _s((v1_item or {}).get("operational_queue")) or "not_current"
|
||||
pending = [task for task in tasks if _s(task.get("status")).lower() == "pending"]
|
||||
call_task = next((task for task in pending if _s(task.get("action_code")).upper() == "CALL_CUSTOMER"), None)
|
||||
call_due = bool(call_task and (not _dt(call_task.get("due_at")) or _dt(call_task.get("due_at")) <= SIMULATION_AT))
|
||||
followup_codes = (
|
||||
{"FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA"}
|
||||
if decision.business_state == "AWAITING_CUSTOMER"
|
||||
else {"FOLLOW_UP_PAYMENT"} if decision.business_state == "AWAITING_PAYMENT" else set()
|
||||
)
|
||||
followup_task = next((
|
||||
task for task in pending if _s(task.get("action_code")).upper() in followup_codes
|
||||
), None)
|
||||
followup_satisfied_now = bool(
|
||||
followup_task and latest_in
|
||||
and _dt(latest_in.get("created_at")) > (_dt(followup_task.get("created_at")) or SIMULATION_AT)
|
||||
) or bool(followup_task and payment and _dt(payment.get("created_at")) > (_dt(followup_task.get("created_at")) or SIMULATION_AT))
|
||||
followup_due = bool(
|
||||
followup_task and not followup_satisfied_now and _dt(followup_task.get("due_at"))
|
||||
and _dt(followup_task.get("due_at")) <= SIMULATION_AT
|
||||
)
|
||||
followup_future = bool(
|
||||
followup_task and not followup_satisfied_now and _dt(followup_task.get("due_at"))
|
||||
and _dt(followup_task.get("due_at")) > SIMULATION_AT
|
||||
)
|
||||
formal_doc_for_reconciliation = bool(current_docs)
|
||||
reconciliation_blocking = bool(
|
||||
formal_doc_for_reconciliation
|
||||
and (conflict_recons or v1_action == "RECONCILE_DOCUMENTS")
|
||||
)
|
||||
fiscal_required = decision.next_action in {"CREATE_PROFORMA", "CREATE_INVOICE"}
|
||||
if blockers and any("conflict" in blocker.casefold() or "without confirmed payment" in blocker.casefold() for blocker in blockers):
|
||||
diagnostic_status = "conflicting_evidence"
|
||||
elif sparse or (reconstructed and not (invoice and payment)) or (odoo_sales and (not invoice or not payment)):
|
||||
diagnostic_status = "incomplete_history"
|
||||
elif ambiguity:
|
||||
diagnostic_status = "ambiguous"
|
||||
else:
|
||||
diagnostic_status = "clear"
|
||||
raw = derive_effective_operational_action(
|
||||
decision,
|
||||
integration_exception=v1_queue == "exception",
|
||||
scheduled_call_current=call_due,
|
||||
due_followup_action=_s((followup_task or {}).get("action_code")).upper() if followup_due else None,
|
||||
future_followup_action=_s((followup_task or {}).get("action_code")).upper() if followup_future else None,
|
||||
fiscal_complete=fiscal_complete,
|
||||
fiscal_required=fiscal_required,
|
||||
reconciliation_blocking=reconciliation_blocking,
|
||||
diagnostic_status=diagnostic_status,
|
||||
)
|
||||
strong_current_evidence = bool(
|
||||
raw.precedence in {"integration_exception", "scheduled_call", "document_prerequisite", "fiscal_prerequisite"}
|
||||
or (decision.business_state in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "EXCEPTION"}
|
||||
and diagnostic_status == "conflicting_evidence")
|
||||
or followup_due
|
||||
or (decision.next_action in {"SEND_PROFORMA"} and proforma)
|
||||
or (decision.next_action == "CONFIRM_PAYMENT" and payment_proof_rows)
|
||||
or (decision.next_action == "CREATE_INVOICE" and payment)
|
||||
or (decision.next_action in {"PREPARE_ORDER", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY"} and invoice and payment)
|
||||
or (decision.next_action == "CREATE_PROFORMA" and order_rows and not proforma)
|
||||
or (decision.next_action in {"SEND_INFO", "SEND_QUOTE"} and latest_in
|
||||
and (not latest_out or _dt(latest_in.get("created_at")) > _dt(latest_out.get("created_at"))))
|
||||
or decision.operational_queue in {"waiting", "not_current"}
|
||||
)
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action=v1_action, v1_queue=v1_queue,
|
||||
strong_current_evidence=strong_current_evidence,
|
||||
)
|
||||
raw_v2 = raw.to_dict()
|
||||
safe_v2 = safe.to_dict()
|
||||
v2_action = raw.effective_operational_action
|
||||
if decision.business_state == "REVIEW_REQUIRED":
|
||||
classification = "REVIEW_REQUIRED"
|
||||
elif ambiguity:
|
||||
classification = "AMBIGUOUS"
|
||||
elif v1_action == v2_action and v1_queue == raw.effective_operational_queue:
|
||||
classification = "UNCHANGED"
|
||||
elif v1_queue in {"do_now", "review", "exception"} and raw.effective_operational_queue in {"waiting", "not_current"}:
|
||||
classification = "DEMOTED_TO_WAITING"
|
||||
elif v1_queue in {"waiting", "backlog", "not_current"} and raw.effective_operational_queue in {"do_now", "review", "exception"}:
|
||||
classification = "PROMOTED_TO_CURRENT"
|
||||
else:
|
||||
classification = "ACTION_CHANGED"
|
||||
|
||||
return _jsonable({
|
||||
"opportunity_id": oid, "title": opp.get("title"),
|
||||
"customer": opp.get("linked_customer_name") or opp.get("customer_name") or opp.get("customer_email"),
|
||||
"v1": {"commercial_stage": opp.get("stage"), "lifecycle_state": opp.get("lifecycle_state"),
|
||||
"current_action": v1_action, "operational_queue": v1_queue,
|
||||
"reason": (v1_item or {}).get("eligibility_reason_code") or v1.get("reason")},
|
||||
"raw_v2": raw_v2, "safe_v2": safe_v2,
|
||||
# Compatibility alias for first-iteration report consumers.
|
||||
"v2": raw_v2, "classification": classification,
|
||||
"evidence": {
|
||||
"latest_relevant_inbound": _event(latest_in) if latest_in else None,
|
||||
"latest_relevant_outbound": _event(latest_out) if latest_out else None,
|
||||
"order_intent": [_event(row) for row in order_rows[-3:]],
|
||||
"fiscal_customer_identity": {"complete": fiscal_complete, "customer_id": _s(opp.get("local_customer_id")), "tax_id_present": bool(opp.get("tax_id"))},
|
||||
"proforma": [{key: _jsonable(row.get(key)) for key in ("id", "external_id", "document_kind", "document_number", "status", "relationship", "created_at")} for row in proformas],
|
||||
"payment": [{key: _jsonable(row.get(key)) for key in ("id", "status", "external_name", "created_at")} for row in payment_links],
|
||||
"invoice": [{key: _jsonable(row.get(key)) for key in ("id", "external_id", "document_number", "status", "relationship", "created_at")} for row in invoices],
|
||||
"odoo": [{key: _jsonable(row.get(key)) for key in ("id", "external_type", "external_id", "external_name", "status", "created_at")} for row in links if row.get("system") == "odoo"],
|
||||
"blockers": blockers, "pending_tasks_for_audit_only": [
|
||||
{key: _jsonable(task.get(key)) for key in ("id", "action_code", "status", "due_at", "created_at")} for task in tasks if task.get("status") == "pending"
|
||||
], "ambiguity": ambiguity, "strong_current_evidence": strong_current_evidence,
|
||||
"reconciliation_blocking": reconciliation_blocking,
|
||||
"fiscal_required_for_transition": fiscal_required,
|
||||
"scheduled_call_current": call_due,
|
||||
"followup_due": followup_due, "followup_future": followup_future,
|
||||
"followup_satisfied": followup_satisfied_now,
|
||||
"diagnostic_status": diagnostic_status,
|
||||
"is_reconstructed": reconstructed,
|
||||
"reconciliation": [{
|
||||
"id": row.get("id"), "document_number": row.get("document_number"),
|
||||
"status": row.get("status"), "created_at": _jsonable(row.get("created_at")),
|
||||
} for row in recons],
|
||||
"material_order_change_evidence": [_event(row) for row in change_rows[-3:]],
|
||||
},
|
||||
})
|
||||
|
||||
|
||||
def _material_keys(row: dict[str, Any]) -> set[str]:
|
||||
evidence = row["evidence"]
|
||||
keys = set()
|
||||
for link in evidence.get("odoo", []):
|
||||
if link.get("external_type") != "sale_order":
|
||||
continue
|
||||
if _s(link.get("external_id")):
|
||||
keys.add(f"odoo_sale_id:{_s(link['external_id']).casefold()}")
|
||||
sale_name = _s(link.get("external_name"))
|
||||
if sale_name and re.fullmatch(r"[A-Z]{1,4}[-/]?[0-9]{2,}", sale_name, re.I):
|
||||
keys.add(f"odoo_sale_name:{sale_name.casefold()}")
|
||||
for kind in ("invoice", "proforma"):
|
||||
for doc in evidence.get(kind, []):
|
||||
if _s(doc.get("external_id")):
|
||||
keys.add(f"jasmin_{kind}_id:{_s(doc['external_id']).casefold()}")
|
||||
if _s(doc.get("document_number")):
|
||||
keys.add(f"jasmin_{kind}_number:{_s(doc['document_number']).casefold()}")
|
||||
return keys
|
||||
|
||||
|
||||
def _apply_material_identity(records: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||
parent = {row["opportunity_id"]: row["opportunity_id"] for row in records}
|
||||
|
||||
def find(value: str) -> str:
|
||||
while parent[value] != value:
|
||||
parent[value] = parent[parent[value]]
|
||||
value = parent[value]
|
||||
return value
|
||||
|
||||
def union(left: str, right: str) -> None:
|
||||
a, b = find(left), find(right)
|
||||
if a != b:
|
||||
parent[b] = a
|
||||
|
||||
by_key: dict[str, list[str]] = defaultdict(list)
|
||||
for row in records:
|
||||
keys = sorted(_material_keys(row))
|
||||
row["material_identity_keys"] = keys
|
||||
for key in keys:
|
||||
by_key[key].append(row["opportunity_id"])
|
||||
for ids in by_key.values():
|
||||
for oid in ids[1:]:
|
||||
union(ids[0], oid)
|
||||
groups: dict[str, list[dict[str, Any]]] = defaultdict(list)
|
||||
for row in records:
|
||||
groups[find(row["opportunity_id"])].append(row)
|
||||
|
||||
report = []
|
||||
for group in groups.values():
|
||||
if len(group) < 2:
|
||||
row = group[0]
|
||||
row["material_process_key"] = next(iter(row["material_identity_keys"]), f"opportunity:{row['opportunity_id']}")
|
||||
row["canonical_process_id"] = row["opportunity_id"]
|
||||
row["duplicate_process_ids"] = []
|
||||
continue
|
||||
|
||||
def score(row: dict[str, Any]) -> tuple[int, str]:
|
||||
evidence = row["evidence"]
|
||||
value = 0
|
||||
value += 50 if not evidence.get("is_reconstructed") else 0
|
||||
value += 20 if evidence.get("invoice") else 0
|
||||
value += 20 if evidence.get("payment") else 0
|
||||
value += 15 if evidence.get("proforma") else 0
|
||||
value += 15 if any(link.get("external_type") == "sale_order" for link in evidence.get("odoo", [])) else 0
|
||||
value += 10 if evidence.get("latest_relevant_inbound") or evidence.get("latest_relevant_outbound") else 0
|
||||
value += 8 if "processo reconstruído" in _s(row.get("title")).casefold() else 0
|
||||
value -= 8 if "sem oportunidade" in _s(row.get("title")).casefold() else 0
|
||||
return value, row["opportunity_id"]
|
||||
|
||||
canonical = max(group, key=score)
|
||||
common = set(canonical["material_identity_keys"])
|
||||
for row in group:
|
||||
common &= set(row["material_identity_keys"])
|
||||
preferred = sorted(common, key=lambda key: (0 if key.startswith("odoo_sale_id:") else 1, key))
|
||||
process_key = preferred[0] if preferred else sorted(canonical["material_identity_keys"])[0]
|
||||
duplicates = [row["opportunity_id"] for row in group if row is not canonical]
|
||||
canonical["material_process_key"] = process_key
|
||||
canonical["canonical_process_id"] = canonical["opportunity_id"]
|
||||
canonical["duplicate_process_ids"] = duplicates
|
||||
for duplicate in group:
|
||||
if duplicate is canonical:
|
||||
continue
|
||||
duplicate["material_process_key"] = process_key
|
||||
duplicate["canonical_process_id"] = canonical["opportunity_id"]
|
||||
duplicate["duplicate_process_ids"] = []
|
||||
duplicate["raw_v2"] = suppress_duplicate_representation(
|
||||
EffectiveOperationalDecision(**duplicate["raw_v2"]),
|
||||
canonical_process_id=canonical["opportunity_id"],
|
||||
).to_dict()
|
||||
duplicate["safe_v2"] = suppress_duplicate_representation(
|
||||
EffectiveOperationalDecision(**duplicate["safe_v2"]),
|
||||
canonical_process_id=canonical["opportunity_id"],
|
||||
).to_dict()
|
||||
duplicate["classification"] = "DUPLICATE_REPRESENTATION"
|
||||
report.append({
|
||||
"material_process_key": process_key,
|
||||
"canonical_process_id": canonical["opportunity_id"],
|
||||
"duplicate_process_ids": duplicates,
|
||||
"identity_keys": sorted(set.intersection(*(set(row["material_identity_keys"]) for row in group))),
|
||||
"canonical_reason": "Highest factual completeness; prefers non-reconstructed process and explicit reconstructed process over an unassociated synthetic record.",
|
||||
})
|
||||
return report
|
||||
|
||||
|
||||
def _totals(items: Iterable[dict[str, Any]], queue_key: str) -> dict[str, int]:
|
||||
counts = Counter(_s(item.get(queue_key)) or "not_current" for item in items)
|
||||
return {
|
||||
"current_work": sum(counts[name] for name in ("do_now", "review", "exception")),
|
||||
"do_now": counts["do_now"], "review": counts["review"], "waiting": counts["waiting"],
|
||||
"backlog": counts["backlog"], "exception": counts["exception"], "not_current": counts["not_current"],
|
||||
}
|
||||
|
||||
|
||||
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,
|
||||
expected_user=expected_user,
|
||||
require_read_only=require_read_only,
|
||||
)
|
||||
opportunities = data["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)
|
||||
all_v1_items = []
|
||||
for key in ("work_items", "waiting_items", "backlog_items", "not_current_items"):
|
||||
all_v1_items.extend(operations.get(key, []))
|
||||
by_opp = {_s(item.get("opportunity_id")): item for item in all_v1_items if item.get("opportunity_id")}
|
||||
records = [_derive_record(opp, data, v1_decisions.get(_s(opp["id"]), {}), by_opp.get(_s(opp["id"]))) for opp in opportunities]
|
||||
material_identity = _apply_material_identity(records)
|
||||
classifications = Counter(row["classification"] for row in records)
|
||||
standalone = []
|
||||
for item in all_v1_items:
|
||||
if item.get("opportunity_id"):
|
||||
continue
|
||||
action = _s(item.get("current_action_code")) or None
|
||||
queue = _s(item.get("operational_queue")) or "backlog"
|
||||
projection = {
|
||||
"business_state": None, "business_next_action": None,
|
||||
"effective_operational_action": action,
|
||||
"effective_operational_queue": queue,
|
||||
"reason": "Standalone canonical Operations work is outside the standard commercial flow and is preserved.",
|
||||
"confidence": "high", "precedence": "preserved_non_opportunity",
|
||||
}
|
||||
standalone.append({
|
||||
"candidate_key": item.get("work_item_key") or f"standalone:{item.get('source')}:{item.get('id')}",
|
||||
"opportunity_id": None, "title": item.get("title"), "customer": item.get("customer_name"),
|
||||
"v1": {"current_action": action, "operational_queue": queue,
|
||||
"reason": item.get("eligibility_reason_code")},
|
||||
"raw_v2": dict(projection), "safe_v2": dict(projection),
|
||||
"classification": "UNCHANGED", "source": item.get("source"),
|
||||
})
|
||||
universe = records + standalone
|
||||
transitions = Counter((
|
||||
row["v1"]["current_action"] or "<NONE>",
|
||||
row["raw_v2"]["effective_operational_action"] or "<WAIT/NONE>",
|
||||
row["safe_v2"]["effective_operational_action"] or "<WAIT/NONE>",
|
||||
) for row in universe)
|
||||
v2_actions = Counter(row["safe_v2"]["effective_operational_action"] or "<WAIT/NONE>" for row in universe)
|
||||
for action in (
|
||||
"SEND_INFO", "SEND_QUOTE", "CREATE_PROFORMA", "SEND_PROFORMA", "CONFIRM_PAYMENT",
|
||||
"CREATE_INVOICE", "PREPARE_ORDER", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY",
|
||||
):
|
||||
v2_actions.setdefault(action, 0)
|
||||
current_v1 = [row for row in universe if row["v1"]["operational_queue"] in {"do_now", "review", "exception"}]
|
||||
false_negatives = []
|
||||
current_obligation_mapping = []
|
||||
for row in current_v1:
|
||||
oid = row.get("opportunity_id")
|
||||
mapping = {
|
||||
"opportunity_id": oid, "title": row.get("title"), "customer": row.get("customer"),
|
||||
"v1_action": row["v1"]["current_action"], "v1_queue": row["v1"]["operational_queue"],
|
||||
"raw_business_state": row["raw_v2"].get("business_state"),
|
||||
"raw_action": row["raw_v2"]["effective_operational_action"],
|
||||
"raw_queue": row["raw_v2"]["effective_operational_queue"],
|
||||
"safe_action": row["safe_v2"]["effective_operational_action"],
|
||||
"safe_queue": row["safe_v2"]["effective_operational_queue"],
|
||||
"disposition": "UNCHANGED" if (
|
||||
row["v1"]["current_action"] == row["safe_v2"]["effective_operational_action"]
|
||||
and row["v1"]["operational_queue"] == row["safe_v2"]["effective_operational_queue"]
|
||||
) else "REPLACED",
|
||||
}
|
||||
current_obligation_mapping.append(mapping)
|
||||
if row["safe_v2"]["effective_operational_queue"] not in {"do_now", "review", "exception"} or row["v1"]["current_action"] != row["safe_v2"]["effective_operational_action"]:
|
||||
false_negatives.append({
|
||||
"opportunity_id": oid, "title": row.get("title"), "customer": row.get("customer"),
|
||||
"v1_action": row["v1"]["current_action"], "v1_queue": row["v1"]["operational_queue"],
|
||||
"raw_state": row["raw_v2"].get("business_state"),
|
||||
"raw_action": row["raw_v2"]["effective_operational_action"],
|
||||
"safe_action": row["safe_v2"]["effective_operational_action"],
|
||||
"safe_queue": row["safe_v2"]["effective_operational_queue"],
|
||||
"reason": row["safe_v2"]["reason"],
|
||||
"pending_tasks": row.get("evidence", {}).get("pending_tasks_for_audit_only", []),
|
||||
})
|
||||
promotion_audit = []
|
||||
for row in records:
|
||||
if row["v1"]["operational_queue"] in {"do_now", "review", "exception"}:
|
||||
continue
|
||||
if row["raw_v2"]["effective_operational_queue"] not in {"do_now", "review", "exception"}:
|
||||
continue
|
||||
# Reproduce the 45-item first-simulation promotion cohort: cases without
|
||||
# a stage/evidence ambiguity, plus the nine Odoo-only cases that the
|
||||
# first model had incorrectly promoted toward completion. Invoice-only
|
||||
# REVIEW_REQUIRED cases were already classified as review, not promotion.
|
||||
odoo_only_completion_error = any(
|
||||
"Odoo execution evidence exists" in blocker
|
||||
for blocker in row["evidence"].get("blockers", [])
|
||||
) and not row["evidence"].get("invoice")
|
||||
if row["evidence"].get("ambiguity"):
|
||||
continue
|
||||
if row["raw_v2"]["business_state"] == "REVIEW_REQUIRED" and not odoo_only_completion_error:
|
||||
continue
|
||||
if row["raw_v2"]["precedence"] in {"fiscal_prerequisite", "document_prerequisite"}:
|
||||
audit_class = "BLOCKER_PRECEDENCE_ERROR"
|
||||
elif row["safe_v2"]["precedence"] == "safe_ambiguous_review":
|
||||
audit_class = "AMBIGUOUS_REVIEW"
|
||||
elif row["safe_v2"]["effective_operational_action"] == row["raw_v2"]["effective_operational_action"]:
|
||||
audit_class = "REAL_PROMOTION"
|
||||
else:
|
||||
audit_class = "HISTORICAL_EVIDENCE_FALSE_POSITIVE"
|
||||
promotion_audit.append({
|
||||
"opportunity_id": row["opportunity_id"], "customer": row["customer"], "title": row["title"],
|
||||
"v1_action": row["v1"]["current_action"], "v1_queue": row["v1"]["operational_queue"],
|
||||
"raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"],
|
||||
"evidence": row["evidence"], "confidence": row["raw_v2"]["confidence"],
|
||||
"classification": audit_class,
|
||||
})
|
||||
invoice_without_payment = []
|
||||
for row in records:
|
||||
if not row["evidence"]["invoice"] or row["evidence"]["payment"]:
|
||||
continue
|
||||
metadata = _payload(next(opp for opp in opportunities if _s(opp["id"]) == row["opportunity_id"]).get("metadata"))
|
||||
text_blob = json.dumps(row["evidence"], ensure_ascii=False).casefold()
|
||||
if _s(metadata.get("payment_terms")).lower() in {"after_delivery", "payment_after_delivery", "pos_entrega"}:
|
||||
category = "PAYMENT_AFTER_INVOICE_ALLOWED"
|
||||
elif "comprovativo" in text_blob or "pagamento" in text_blob:
|
||||
category = "PAYMENT_EVIDENCE_MISSING"
|
||||
elif any(term in text_blob for term in ("nota de crédito", "nota de credito", "corrigir", "anular")):
|
||||
category = "FINANCIAL_CORRECTION_REQUIRED"
|
||||
elif not row["evidence"]["order_intent"] and not row["evidence"]["proforma"]:
|
||||
category = "PREMATURE_INVOICE"
|
||||
else:
|
||||
category = "UNKNOWN_REVIEW"
|
||||
invoice_without_payment.append({
|
||||
"opportunity_id": row["opportunity_id"], "title": row["title"], "customer": row["customer"],
|
||||
"classification": category, "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"],
|
||||
"evidence": row["evidence"],
|
||||
})
|
||||
review_audit = []
|
||||
for row in universe:
|
||||
if row["safe_v2"]["effective_operational_queue"] != "review":
|
||||
continue
|
||||
diagnostic = row.get("evidence", {}).get("diagnostic_status", "clear")
|
||||
if row["v1"]["operational_queue"] == "review" and row["v1"]["current_action"] == row["safe_v2"]["effective_operational_action"]:
|
||||
category = "EXISTING_VALID_REVIEW"
|
||||
elif row["v1"]["operational_queue"] not in {"do_now", "review", "exception"}:
|
||||
category = "REAL_CURRENT_PROMOTION"
|
||||
elif diagnostic == "incomplete_history":
|
||||
category = "HISTORICAL_INCOMPLETE"
|
||||
elif diagnostic == "ambiguous":
|
||||
category = "DIAGNOSTIC_ONLY"
|
||||
else:
|
||||
category = "ACTIONABLE_REVIEW"
|
||||
review_audit.append({
|
||||
"opportunity_id": row.get("opportunity_id"), "title": row.get("title"),
|
||||
"v1": row["v1"], "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"],
|
||||
"diagnostic_status": diagnostic, "classification": category,
|
||||
})
|
||||
backlog_delta = []
|
||||
for row in universe:
|
||||
if row["v1"]["operational_queue"] != "backlog" or row["safe_v2"]["effective_operational_queue"] == "backlog":
|
||||
continue
|
||||
safe_queue = row["safe_v2"]["effective_operational_queue"]
|
||||
if row["safe_v2"]["precedence"] == "duplicate_representation":
|
||||
category = "DEDUPLICATED"
|
||||
elif safe_queue in {"do_now", "review", "exception"}:
|
||||
category = "PROMOTED_TO_CURRENT"
|
||||
elif safe_queue == "waiting":
|
||||
category = "MOVED_TO_WAITING"
|
||||
elif row.get("evidence", {}).get("strong_current_evidence"):
|
||||
category = "LEGITIMATE_BACKLOG_REMOVAL"
|
||||
else:
|
||||
category = "SHOULD_REMAIN_BACKLOG"
|
||||
backlog_delta.append({
|
||||
"opportunity_id": row.get("opportunity_id"), "title": row.get("title"),
|
||||
"v1_action": row["v1"]["current_action"], "safe_v2": row["safe_v2"],
|
||||
"evidence": row.get("evidence", {}), "classification": category,
|
||||
})
|
||||
current_delta = []
|
||||
for row in universe:
|
||||
if row["v1"]["operational_queue"] in {"do_now", "review", "exception"}:
|
||||
continue
|
||||
if row["safe_v2"]["effective_operational_queue"] not in {"do_now", "review", "exception"}:
|
||||
continue
|
||||
precedence = row["safe_v2"]["precedence"]
|
||||
diagnostic = row.get("evidence", {}).get("diagnostic_status", "clear")
|
||||
if precedence == "due_followup":
|
||||
category = "DUE_FOLLOW_UP"
|
||||
elif precedence in {"fiscal_prerequisite", "document_prerequisite", "integration_exception"}:
|
||||
category = "BLOCKER"
|
||||
elif precedence == "duplicate_representation":
|
||||
category = "DUPLICATE"
|
||||
elif row.get("evidence", {}).get("strong_current_evidence") and row["safe_v2"]["effective_operational_queue"] != "review":
|
||||
category = "REAL_NEW_OBLIGATION"
|
||||
elif diagnostic in {"ambiguous", "incomplete_history"}:
|
||||
category = "DIAGNOSTIC_ONLY"
|
||||
elif row["safe_v2"]["effective_operational_queue"] == "review":
|
||||
category = "ACTIONABLE_REVIEW"
|
||||
else:
|
||||
category = "FALSE_PROMOTION"
|
||||
current_delta.append({
|
||||
"opportunity_id": row.get("opportunity_id"), "title": row.get("title"),
|
||||
"v1": row["v1"], "raw_v2": row["raw_v2"], "safe_v2": row["safe_v2"],
|
||||
"evidence": row.get("evidence", {}), "classification": category,
|
||||
})
|
||||
def action_counts(rows: list[dict[str, Any]]) -> dict[str, int]:
|
||||
return dict(Counter(row["safe_v2"]["effective_operational_action"] or "<NONE>" for row in rows))
|
||||
action_count_scopes = {
|
||||
"all_candidates": action_counts(universe),
|
||||
"current_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] in {"do_now", "review", "exception"}]),
|
||||
"do_now_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "do_now"]),
|
||||
"review_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "review"]),
|
||||
"waiting_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "waiting"]),
|
||||
"backlog_only": action_counts([row for row in universe if row["safe_v2"]["effective_operational_queue"] == "backlog"]),
|
||||
}
|
||||
v1_totals = _totals([{"queue": row["v1"]["operational_queue"]} for row in universe], "queue")
|
||||
raw_totals = _totals([{"queue": row["raw_v2"]["effective_operational_queue"]} for row in universe], "queue")
|
||||
safe_totals = _totals([{"queue": row["safe_v2"]["effective_operational_queue"]} for row in universe], "queue")
|
||||
result = {
|
||||
"generated_at": datetime.now(timezone.utc), "database": data["identity"],
|
||||
"opportunity_count": len(records), "candidate_universe_count": len(universe),
|
||||
"v1_totals": v1_totals, "raw_v2_totals": raw_totals, "safe_v2_totals": safe_totals,
|
||||
"classifications": dict(classifications), "v2_actions": dict(v2_actions),
|
||||
"transitions": [{"v1_action": old, "raw_v2_action": raw, "safe_v2_action": safe, "count": count} for (old, raw, safe), count in transitions.most_common()],
|
||||
"possible_false_negatives": false_negatives,
|
||||
"v1_current_obligation_mapping": current_obligation_mapping,
|
||||
"promotions_audit": promotion_audit,
|
||||
"invoice_without_payment": invoice_without_payment,
|
||||
"summary": {
|
||||
"raw_ambiguous_count": sum(row["raw_v2"]["confidence"] != "high" for row in records),
|
||||
"safe_overrides_count": sum(
|
||||
(row["raw_v2"]["effective_operational_action"], row["raw_v2"]["effective_operational_queue"])
|
||||
!= (row["safe_v2"]["effective_operational_action"], row["safe_v2"]["effective_operational_queue"])
|
||||
for row in universe
|
||||
),
|
||||
"preserved_v1_obligations": sum(row["disposition"] == "UNCHANGED" for row in current_obligation_mapping),
|
||||
"real_promotions": sum(row["classification"] == "REAL_PROMOTION" for row in promotion_audit),
|
||||
"rejected_promotions": sum(row["classification"] == "HISTORICAL_EVIDENCE_FALSE_POSITIVE" for row in promotion_audit),
|
||||
"ambiguous_promotions": sum(row["classification"] == "AMBIGUOUS_REVIEW" for row in promotion_audit),
|
||||
"fiscal_blockers_preserved": sum(row["raw_v2"]["precedence"] == "fiscal_prerequisite" for row in records),
|
||||
"reconciliation_blockers_preserved": sum(row["raw_v2"]["precedence"] == "document_prerequisite" for row in records),
|
||||
"non_opportunity_canonical_work_preserved": len(standalone),
|
||||
"diagnostic_ambiguous_not_current": sum(
|
||||
row.get("evidence", {}).get("diagnostic_status") in {"ambiguous", "incomplete_history"}
|
||||
and row["safe_v2"]["effective_operational_queue"] == "not_current" for row in records
|
||||
),
|
||||
"actionable_review": sum(row["classification"] in {"ACTIONABLE_REVIEW", "EXISTING_VALID_REVIEW", "REAL_CURRENT_PROMOTION"} for row in review_audit),
|
||||
"due_followups": sum(row["safe_v2"]["precedence"] == "due_followup" for row in records),
|
||||
"safe_v2_current_minus_v1": len(current_delta),
|
||||
"backlog_delta_explained": len(backlog_delta),
|
||||
"duplicate_material_processes": len(material_identity),
|
||||
"duplicate_current_cards_suppressed": sum(len(row["duplicate_process_ids"]) for row in material_identity),
|
||||
},
|
||||
"action_counts": action_count_scopes,
|
||||
"review_audit": review_audit, "backlog_delta": backlog_delta,
|
||||
"current_delta": current_delta, "material_identity": material_identity,
|
||||
"standalone_canonical_items": standalone, "opportunities": records,
|
||||
}
|
||||
return _jsonable(result)
|
||||
|
||||
|
||||
def _named(records: list[dict[str, Any]], name: str) -> list[dict[str, Any]]:
|
||||
folded = name.casefold()
|
||||
return [row for row in records if folded in f"{_s(row.get('title'))} {_s(row.get('customer'))}".casefold()]
|
||||
|
||||
|
||||
def _comparison(result: dict[str, Any]) -> str:
|
||||
lines = [
|
||||
"BLIF FLOW V2 SHADOW SIMULATION", "",
|
||||
f"Database: {result['database']}", f"Opportunities: {result['opportunity_count']}",
|
||||
f"Comparable candidate universe: {result['candidate_universe_count']}", "",
|
||||
"CENTRO DE TRABALHO", "metric V1 RAW V2 SAFE V2",
|
||||
]
|
||||
for key in ("current_work", "do_now", "review", "waiting", "backlog", "exception", "not_current"):
|
||||
lines.append(
|
||||
f"{key:<30} {result['v1_totals'].get(key, 0):>5}"
|
||||
f" {result['raw_v2_totals'].get(key, 0):>7} {result['safe_v2_totals'].get(key, 0):>7}"
|
||||
)
|
||||
lines += ["", "SUMMARY"]
|
||||
for key, value in result["summary"].items():
|
||||
lines.append(f"{key}: {value}")
|
||||
for scope in ("all_candidates", "current_only", "do_now_only", "review_only", "waiting_only", "backlog_only"):
|
||||
lines += ["", f"SAFE V2 ACTIONS — {scope}"]
|
||||
for action, count in sorted(result["action_counts"][scope].items(), key=lambda item: (-item[1], item[0])):
|
||||
lines.append(f"{action:<36} {count:>5}")
|
||||
lines += ["", "CLASSIFICATIONS"]
|
||||
for name in ("UNCHANGED", "ACTION_CHANGED", "DEMOTED_TO_WAITING", "PROMOTED_TO_CURRENT", "REVIEW_REQUIRED", "AMBIGUOUS"):
|
||||
lines.append(f"{name:<30} {result['classifications'].get(name, 0):>5}")
|
||||
lines += ["", "V1 ACTION -> RAW V2 ACTION -> SAFE V2 ACTION"]
|
||||
for row in result["transitions"]:
|
||||
lines.append(f"{row['v1_action']} -> {row['raw_v2_action']} -> {row['safe_v2_action']}: {row['count']}")
|
||||
lines += ["", "POSSIBLE FALSE NEGATIVES"]
|
||||
if not result["possible_false_negatives"]:
|
||||
lines.append("None.")
|
||||
for row in result["possible_false_negatives"]:
|
||||
lines.append(json.dumps(row, ensure_ascii=False, sort_keys=True))
|
||||
lines += ["", "EVERY CURRENT V1 OBLIGATION -> V2"]
|
||||
for row in result["v1_current_obligation_mapping"]:
|
||||
lines.append(json.dumps(row, ensure_ascii=False, sort_keys=True))
|
||||
for name in ("INSTALBEIRA", "ENGEXICON", "RZSOLAR", "X MAT", "CONSTRURECUP", "PANORAMIC SUCCESS"):
|
||||
lines += ["", name]
|
||||
matches = _named(result["opportunities"], name)
|
||||
lines.extend(json.dumps(row, ensure_ascii=False, sort_keys=True) for row in matches)
|
||||
if not matches:
|
||||
lines.append("No opportunity title/customer match.")
|
||||
return "\n".join(lines) + "\n"
|
||||
|
||||
|
||||
def main() -> None:
|
||||
result = collect()
|
||||
PROJECTION.write_text(json.dumps(result, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
ambiguous = [row for row in result["opportunities"] if row["classification"] in {"AMBIGUOUS", "REVIEW_REQUIRED"}]
|
||||
AMBIGUOUS.write_text(json.dumps(ambiguous, ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
PROMOTIONS.write_text(json.dumps(result["promotions_audit"], ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
INVOICE_WITHOUT_PAYMENT.write_text(json.dumps(result["invoice_without_payment"], ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
REVIEW_AUDIT.write_text(json.dumps(result["review_audit"], ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
BACKLOG_DELTA.write_text(json.dumps(result["backlog_delta"], ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
CURRENT_DELTA.write_text(json.dumps(result["current_delta"], ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
MATERIAL_IDENTITY.write_text(json.dumps(result["material_identity"], ensure_ascii=False, indent=2), encoding="utf-8")
|
||||
COMPARISON.write_text(_comparison(result), encoding="utf-8")
|
||||
print(_comparison(result), end="")
|
||||
print(f"Output: {PROJECTION}\nOutput: {COMPARISON}\nOutput: {AMBIGUOUS}\nOutput: {PROMOTIONS}\nOutput: {INVOICE_WITHOUT_PAYMENT}\nOutput: {REVIEW_AUDIT}\nOutput: {BACKLOG_DELTA}\nOutput: {CURRENT_DELTA}\nOutput: {MATERIAL_IDENTITY}")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
119
scripts/validate_blif_flow_v2_persistence.py
Normal file
119
scripts/validate_blif_flow_v2_persistence.py
Normal file
@@ -0,0 +1,119 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Validate migration 011 and two rebuilds using temp tables in the test DB."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
from sqlalchemy import text
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
sys.path.insert(0, str(ROOT))
|
||||
os.chdir(ROOT)
|
||||
|
||||
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
|
||||
from app.db import engine
|
||||
|
||||
|
||||
NAMED_IDS = {
|
||||
"instalbeira": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
|
||||
"engexicon": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
|
||||
"construrecup": "e3b23ac5-84db-4763-8a31-a684e873032c",
|
||||
"panoramic": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
|
||||
"x_mat_canonical": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
|
||||
"x_mat_reconstructed": "1816a06e-9a69-4a9b-9279-1263156892d3",
|
||||
"rzsolar_reconstructed": "fd221608-e007-4043-a23d-07e0c119a345",
|
||||
"rzsolar_synthetic": "434124fb-ac19-4d78-909a-55761d7e8daa",
|
||||
}
|
||||
|
||||
|
||||
def main() -> None:
|
||||
conn = engine.connect().execution_options(isolation_level="AUTOCOMMIT")
|
||||
try:
|
||||
identity = conn.execute(text(
|
||||
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
||||
)).one()
|
||||
if tuple(identity[:2]) != ("clientflow_codex_test", "clientflow_codex_test"):
|
||||
raise RuntimeError(f"refusing persistence validation on {identity!r}")
|
||||
conn.exec_driver_sql("BEGIN READ WRITE")
|
||||
try:
|
||||
stage_fingerprint_before = conn.execute(text("""
|
||||
SELECT md5(string_agg(id::text || ':' || COALESCE(stage,''), ',' ORDER BY id))
|
||||
FROM public.opportunities
|
||||
""")).scalar_one()
|
||||
conn.exec_driver_sql("SET LOCAL search_path TO pg_temp, public")
|
||||
conn.exec_driver_sql("CREATE TEMP TABLE opportunities (id UUID PRIMARY KEY)")
|
||||
conn.exec_driver_sql("INSERT INTO opportunities SELECT id FROM public.opportunities")
|
||||
conn.exec_driver_sql("CREATE TEMP TABLE opportunity_events (LIKE public.opportunity_events INCLUDING DEFAULTS INCLUDING CONSTRAINTS)")
|
||||
conn.exec_driver_sql("ALTER TABLE opportunity_events ADD PRIMARY KEY (id)")
|
||||
conn.exec_driver_sql("CREATE TEMP TABLE tasks (LIKE public.tasks INCLUDING DEFAULTS INCLUDING CONSTRAINTS)")
|
||||
conn.exec_driver_sql("ALTER TABLE tasks ADD PRIMARY KEY (id)")
|
||||
conn.exec_driver_sql(Path("migrations/011_blif_flow_v2_persistence.sql").read_text())
|
||||
|
||||
first = rebuild_blif_flow_v2_projection(
|
||||
mode="shadow", target_schema="pg_temp", connection=conn,
|
||||
)
|
||||
transition_count_first = conn.execute(text(
|
||||
"SELECT count(*) FROM pg_temp.opportunity_flow_transitions"
|
||||
)).scalar_one()
|
||||
first_fingerprints = dict(conn.execute(text("""
|
||||
SELECT opportunity_id::text, source_fingerprint
|
||||
FROM pg_temp.opportunity_flow_state_v2
|
||||
""")).all())
|
||||
second = rebuild_blif_flow_v2_projection(
|
||||
mode="shadow", target_schema="pg_temp", connection=conn,
|
||||
)
|
||||
transition_count_second = conn.execute(text(
|
||||
"SELECT count(*) FROM pg_temp.opportunity_flow_transitions"
|
||||
)).scalar_one()
|
||||
second_fingerprints = dict(conn.execute(text("""
|
||||
SELECT opportunity_id::text, source_fingerprint
|
||||
FROM pg_temp.opportunity_flow_state_v2
|
||||
""")).all())
|
||||
named = {}
|
||||
for name, oid in NAMED_IDS.items():
|
||||
row = 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
|
||||
FROM pg_temp.opportunity_flow_state_v2
|
||||
WHERE opportunity_id=CAST(:oid AS UUID)
|
||||
"""), {"oid": oid}).mappings().one()
|
||||
named[name] = dict(row)
|
||||
stage_fingerprint_after = conn.execute(text("""
|
||||
SELECT md5(string_agg(id::text || ':' || COALESCE(stage,''), ',' ORDER BY id))
|
||||
FROM public.opportunities
|
||||
""")).scalar_one()
|
||||
result = {
|
||||
"database": {"name": identity[0], "user": identity[1]},
|
||||
"migration_scope": "transaction-scoped pg_temp (public tasks is postgres-owned)",
|
||||
"first_rebuild": first, "second_rebuild": second,
|
||||
"transition_count_first": transition_count_first,
|
||||
"transition_count_second": transition_count_second,
|
||||
"idempotent": transition_count_first == transition_count_second
|
||||
and first_fingerprints == second_fingerprints,
|
||||
"opportunity_stage_unchanged": stage_fingerprint_before == stage_fingerprint_after,
|
||||
"named": named,
|
||||
}
|
||||
conn.exec_driver_sql(Path("migrations/011_blif_flow_v2_persistence_down.sql").read_text())
|
||||
remaining_task_columns = conn.execute(text("""
|
||||
SELECT count(*) FROM information_schema.columns
|
||||
WHERE table_schema LIKE 'pg_temp_%' AND table_name='tasks'
|
||||
AND column_name IN ('resolution_code','resolved_at','resolved_by_event_id','superseded_by_task_id')
|
||||
""")).scalar_one()
|
||||
result["down_migration_reversible"] = (
|
||||
conn.execute(text("SELECT to_regclass('pg_temp.opportunity_flow_state_v2')")).scalar_one() is None
|
||||
and conn.execute(text("SELECT to_regclass('pg_temp.opportunity_flow_transitions')")).scalar_one() is None
|
||||
and remaining_task_columns == 0
|
||||
)
|
||||
print(json.dumps(result, ensure_ascii=False, indent=2, sort_keys=True))
|
||||
finally:
|
||||
conn.exec_driver_sql("ROLLBACK")
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
244
tests/domain/opportunity_flow/test_blif_flow_v2.py
Normal file
244
tests/domain/opportunity_flow/test_blif_flow_v2.py
Normal file
@@ -0,0 +1,244 @@
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
import pytest
|
||||
|
||||
from app.domain.opportunity_flow.v2 import (
|
||||
derive_business_facts, derive_effective_operational_action,
|
||||
derive_safe_operational_action, derive_v2_operational_queue,
|
||||
)
|
||||
from scripts.simulate_blif_flow_v2 import _apply_material_identity
|
||||
|
||||
|
||||
NOW = datetime(2026, 8, 15, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("facts", "state", "action", "queue"),
|
||||
[
|
||||
({"customer_request": True}, "INQUIRY", "SEND_INFO", "do_now"),
|
||||
({"customer_request": True, "info_or_offer_sent": True}, "AWAITING_CUSTOMER", None, "waiting"),
|
||||
({"order_intent": True}, "PROFORMA_REQUIRED", "CREATE_PROFORMA", "do_now"),
|
||||
({"order_intent": True, "proforma_exists": True}, "PROFORMA_CREATED", "SEND_PROFORMA", "do_now"),
|
||||
({"order_intent": True, "proforma_exists": True, "proforma_sent": True}, "AWAITING_PAYMENT", None, "waiting"),
|
||||
({"payment_confirmed": True}, "INVOICE_REQUIRED", "CREATE_INVOICE", "do_now"),
|
||||
({"payment_confirmed": True, "invoice_exists": True}, "ODOO_ORDER_REQUIRED", "PREPARE_ORDER", "do_now"),
|
||||
({"payment_confirmed": True, "invoice_exists": True, "odoo_order_exists": True}, "ODOO_ORDER_CREATED", "VALIDATE_ODOO_ORDER", "do_now"),
|
||||
({"payment_confirmed": True, "invoice_exists": True, "odoo_order_exists": True, "odoo_order_validated": True}, "ODOO_ORDER_VALIDATED", "COMPLETE_OPPORTUNITY", "do_now"),
|
||||
],
|
||||
)
|
||||
def test_normal_flow(facts, state, action, queue):
|
||||
decision = derive_v2_operational_queue(derive_business_facts(**facts))
|
||||
assert (decision.business_state, decision.next_action, decision.operational_queue) == (state, action, queue)
|
||||
|
||||
|
||||
def test_inquiry_without_current_request_does_not_create_work():
|
||||
decision = derive_v2_operational_queue(derive_business_facts())
|
||||
assert (decision.business_state, decision.next_action, decision.operational_queue) == (
|
||||
"INQUIRY", None, "not_current"
|
||||
)
|
||||
|
||||
|
||||
def test_order_change_before_payment_requires_new_proforma():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(
|
||||
order_intent=True, proforma_exists=True, proforma_sent=True,
|
||||
proforma_created_at=NOW - timedelta(days=2), material_order_change=True,
|
||||
material_order_change_at=NOW - timedelta(days=1),
|
||||
))
|
||||
assert (decision.business_state, decision.next_action) == ("PROFORMA_REQUIRED", "CREATE_PROFORMA")
|
||||
|
||||
|
||||
def test_order_change_after_payment_requires_review():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(
|
||||
payment_confirmed=True, payment_confirmed_at=NOW - timedelta(days=2),
|
||||
material_order_change=True, material_order_change_at=NOW - timedelta(days=1),
|
||||
))
|
||||
assert (decision.business_state, decision.next_action, decision.operational_queue) == (
|
||||
"REVIEW_REQUIRED", "REVIEW_REQUIRED", "review"
|
||||
)
|
||||
|
||||
|
||||
def test_completed_send_quote_task_does_not_create_proforma():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(
|
||||
order_intent=True, audit_task_codes=["SEND_QUOTE:done"],
|
||||
))
|
||||
assert (decision.business_state, decision.next_action) == ("PROFORMA_REQUIRED", "CREATE_PROFORMA")
|
||||
|
||||
|
||||
def test_stale_pending_task_cannot_override_stronger_fact():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(
|
||||
payment_confirmed=True, invoice_exists=True, audit_task_codes=["SEND_PROFORMA:pending"],
|
||||
))
|
||||
assert (decision.business_state, decision.next_action) == ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER")
|
||||
|
||||
|
||||
def test_invoice_without_confirmed_payment_requires_review_not_prepare_order():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(invoice_exists=True))
|
||||
assert (decision.business_state, decision.next_action, decision.operational_queue) == (
|
||||
"REVIEW_REQUIRED", "REVIEW_REQUIRED", "review"
|
||||
)
|
||||
|
||||
|
||||
def test_later_customer_inbound_satisfies_old_followup_and_needs_response():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(
|
||||
customer_request=True, info_or_offer_sent=True,
|
||||
latest_relevant_outbound_at=NOW - timedelta(days=2),
|
||||
latest_relevant_inbound_at=NOW - timedelta(days=1),
|
||||
later_customer_inbound_satisfies_followup=True,
|
||||
audit_task_codes=["FOLLOW_UP_CUSTOMER_REVIEW:pending"],
|
||||
))
|
||||
assert (decision.business_state, decision.next_action) == ("INQUIRY", "SEND_INFO")
|
||||
|
||||
|
||||
def test_odoo_without_mandatory_financial_evidence_requires_review():
|
||||
decision = derive_v2_operational_queue(derive_business_facts(
|
||||
odoo_order_exists=True, odoo_order_validated=True,
|
||||
))
|
||||
assert (decision.business_state, decision.next_action) == ("REVIEW_REQUIRED", "REVIEW_REQUIRED")
|
||||
|
||||
|
||||
def test_fiscal_prerequisite_changes_effective_not_business_action():
|
||||
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
|
||||
effective = derive_effective_operational_action(
|
||||
business, fiscal_complete=False, fiscal_required=True,
|
||||
)
|
||||
assert effective.business_next_action == "CREATE_PROFORMA"
|
||||
assert effective.effective_operational_action == "VALIDATE_FISCAL_CUSTOMER"
|
||||
assert effective.precedence == "fiscal_prerequisite"
|
||||
|
||||
|
||||
def test_reconciliation_only_blocks_when_adapter_proves_current_document():
|
||||
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
|
||||
unblocked = derive_effective_operational_action(business, reconciliation_blocking=False)
|
||||
blocked = derive_effective_operational_action(business, reconciliation_blocking=True)
|
||||
assert unblocked.effective_operational_action == "CREATE_PROFORMA"
|
||||
assert blocked.effective_operational_action == "RECONCILE_DOCUMENTS"
|
||||
|
||||
|
||||
def test_scheduled_call_precedes_business_transition():
|
||||
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
|
||||
effective = derive_effective_operational_action(business, scheduled_call_current=True)
|
||||
assert (effective.effective_operational_action, effective.precedence) == ("CALL_CUSTOMER", "scheduled_call")
|
||||
|
||||
|
||||
def test_safe_projection_preserves_current_v1_obligation_when_raw_is_uncertain():
|
||||
business = derive_v2_operational_queue(derive_business_facts(customer_request=True))
|
||||
raw = derive_effective_operational_action(business)
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action="CREATE_JASMIN_QUOTE", v1_queue="do_now",
|
||||
strong_current_evidence=False,
|
||||
)
|
||||
assert (safe.effective_operational_action, safe.effective_operational_queue) == (
|
||||
"CREATE_JASMIN_QUOTE", "do_now"
|
||||
)
|
||||
|
||||
|
||||
def test_safe_projection_accepts_strong_raw_transition():
|
||||
business = derive_v2_operational_queue(derive_business_facts(order_intent=True))
|
||||
raw = derive_effective_operational_action(business)
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action="RECONCILE_DOCUMENTS", v1_queue="review",
|
||||
strong_current_evidence=True,
|
||||
)
|
||||
assert safe == raw
|
||||
|
||||
|
||||
def test_safe_projection_keeps_strong_factual_review_over_technical_v1_action():
|
||||
business = derive_v2_operational_queue(derive_business_facts(odoo_order_exists=True))
|
||||
raw = derive_effective_operational_action(business)
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action="CREATE_JASMIN_QUOTE", v1_queue="do_now",
|
||||
strong_current_evidence=True,
|
||||
)
|
||||
assert (safe.effective_operational_action, safe.effective_operational_queue) == (
|
||||
"REVIEW_REQUIRED", "review"
|
||||
)
|
||||
|
||||
|
||||
def test_ambiguity_does_not_promote_review_work():
|
||||
business = derive_v2_operational_queue(derive_business_facts(customer_request=True))
|
||||
raw = derive_effective_operational_action(business, diagnostic_status="ambiguous")
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action=None, v1_queue="not_current", strong_current_evidence=False,
|
||||
)
|
||||
assert (safe.effective_operational_action, safe.effective_operational_queue) == (None, "not_current")
|
||||
assert safe.diagnostic_status == "ambiguous"
|
||||
|
||||
|
||||
def test_historical_incomplete_process_stays_not_current():
|
||||
business = derive_v2_operational_queue(derive_business_facts(odoo_order_exists=True))
|
||||
raw = derive_effective_operational_action(business, diagnostic_status="incomplete_history")
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action=None, v1_queue="not_current", strong_current_evidence=False,
|
||||
)
|
||||
assert (safe.effective_operational_action, safe.effective_operational_queue) == (None, "not_current")
|
||||
|
||||
|
||||
def test_due_customer_followup_becomes_do_now_while_business_waits():
|
||||
business = derive_v2_operational_queue(derive_business_facts(
|
||||
customer_request=True, info_or_offer_sent=True,
|
||||
))
|
||||
effective = derive_effective_operational_action(
|
||||
business, due_followup_action="FOLLOW_UP_CUSTOMER_REVIEW",
|
||||
)
|
||||
assert effective.business_state == "AWAITING_CUSTOMER"
|
||||
assert (effective.effective_operational_action, effective.effective_operational_queue) == (
|
||||
"FOLLOW_UP_CUSTOMER_REVIEW", "do_now"
|
||||
)
|
||||
|
||||
|
||||
def test_future_customer_followup_remains_waiting():
|
||||
business = derive_v2_operational_queue(derive_business_facts(
|
||||
customer_request=True, info_or_offer_sent=True,
|
||||
))
|
||||
effective = derive_effective_operational_action(
|
||||
business, future_followup_action="FOLLOW_UP_CUSTOMER_REVIEW",
|
||||
)
|
||||
assert (effective.effective_operational_action, effective.effective_operational_queue) == (
|
||||
"FOLLOW_UP_CUSTOMER_REVIEW", "waiting"
|
||||
)
|
||||
|
||||
|
||||
def test_backlog_remains_backlog_without_stronger_current_evidence():
|
||||
business = derive_v2_operational_queue(derive_business_facts(customer_request=True))
|
||||
raw = derive_effective_operational_action(business, diagnostic_status="ambiguous")
|
||||
safe = derive_safe_operational_action(
|
||||
raw, v1_action="VALIDATE_FISCAL_CUSTOMER", v1_queue="backlog",
|
||||
strong_current_evidence=False,
|
||||
)
|
||||
assert (safe.effective_operational_action, safe.effective_operational_queue) == (
|
||||
"VALIDATE_FISCAL_CUSTOMER", "backlog"
|
||||
)
|
||||
|
||||
|
||||
def _identity_record(oid, *, odoo_id=None, invoice_number=None, complete=False):
|
||||
business = derive_v2_operational_queue(derive_business_facts(
|
||||
payment_confirmed=complete, invoice_exists=complete,
|
||||
))
|
||||
projection = derive_effective_operational_action(business).to_dict()
|
||||
return {
|
||||
"opportunity_id": oid, "title": oid, "customer": oid,
|
||||
"material_identity_keys": [], "classification": "UNCHANGED",
|
||||
"raw_v2": dict(projection), "safe_v2": dict(projection),
|
||||
"evidence": {
|
||||
"is_reconstructed": not complete, "invoice": ([{"document_number": invoice_number}] if invoice_number else []),
|
||||
"payment": ([{"id": "p"}] if complete else []), "proforma": [],
|
||||
"odoo": ([{"external_type": "sale_order", "external_id": odoo_id, "external_name": f"S{odoo_id}"}] if odoo_id else []),
|
||||
"latest_relevant_inbound": None, "latest_relevant_outbound": None,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def test_same_odoo_external_id_suppresses_second_current_card():
|
||||
records = [_identity_record("canonical", odoo_id="349", complete=True),
|
||||
_identity_record("reconstructed", odoo_id="349")]
|
||||
groups = _apply_material_identity(records)
|
||||
assert groups[0]["canonical_process_id"] == "canonical"
|
||||
assert records[1]["safe_v2"]["effective_operational_queue"] == "not_current"
|
||||
|
||||
|
||||
def test_same_invoice_identity_suppresses_second_current_card():
|
||||
records = [_identity_record("canonical", invoice_number="FA.186", complete=True),
|
||||
_identity_record("duplicate", invoice_number="FA.186")]
|
||||
groups = _apply_material_identity(records)
|
||||
assert len(groups) == 1
|
||||
assert sum(row["safe_v2"]["effective_operational_queue"] != "not_current" for row in records) == 1
|
||||
113
tests/domain/opportunity_flow/test_blif_flow_v2_data_repair.py
Normal file
113
tests/domain/opportunity_flow/test_blif_flow_v2_data_repair.py
Normal file
@@ -0,0 +1,113 @@
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from app.domain.opportunity_flow.repair import (
|
||||
TaskRepairContext, classify_pending_task, simulate_high_repairs,
|
||||
)
|
||||
from app.domain.opportunity_flow.v2 import (
|
||||
EffectiveOperationalDecision, suppress_duplicate_representation,
|
||||
)
|
||||
|
||||
|
||||
NOW = datetime(2026, 8, 15, tzinfo=timezone.utc)
|
||||
|
||||
|
||||
def decision(**overrides):
|
||||
values = dict(task_id="task-1", opportunity_id="opp-1", action_code="SEND_INFO",
|
||||
created_at=NOW - timedelta(days=10), business_state="INQUIRY")
|
||||
values.update(overrides)
|
||||
return classify_pending_task(TaskRepairContext(**values))
|
||||
|
||||
|
||||
def test_satisfied_task_classified_satisfied_by_event():
|
||||
got = decision(later_outbound_event={"id": "message-1"})
|
||||
assert got.classification == "SATISFIED_BY_EVENT"
|
||||
assert got.auto_repair_safe
|
||||
|
||||
|
||||
def test_superseded_task_classified_superseded():
|
||||
assert decision(action_code="CREATE_PROFORMA", business_state="INVOICE_CREATED").classification == "SUPERSEDED"
|
||||
|
||||
|
||||
def test_duplicate_material_task_classified_duplicate():
|
||||
got = decision(is_duplicate_representation=True, canonical_opportunity_id="canonical")
|
||||
assert got.classification == "DUPLICATE"
|
||||
assert got.safety_tier == "HIGH"
|
||||
|
||||
|
||||
def test_downstream_task_without_prerequisite_classified_premature():
|
||||
assert decision(action_code="SEND_INVOICE", business_state="PROFORMA_REQUIRED").classification == "PREMATURE"
|
||||
|
||||
|
||||
def test_valid_overdue_followup_remains_valid_current():
|
||||
got = decision(action_code="FOLLOW_UP_CUSTOMER_REVIEW", business_state="AWAITING_CUSTOMER",
|
||||
due_at=NOW - timedelta(days=2))
|
||||
assert got.classification == "VALID_CURRENT"
|
||||
|
||||
|
||||
def test_future_valid_followup_remains_valid_waiting():
|
||||
got = decision(action_code="FOLLOW_UP_CUSTOMER_REVIEW", business_state="AWAITING_CUSTOMER",
|
||||
due_at=NOW + timedelta(days=2))
|
||||
assert got.classification == "VALID_CURRENT"
|
||||
|
||||
|
||||
def test_historical_age_alone_never_closes_task():
|
||||
got = decision(action_code="REVIEW_MANUALLY", created_at=NOW - timedelta(days=1000))
|
||||
assert got.classification == "VALID_CURRENT"
|
||||
|
||||
|
||||
def test_ambiguous_evidence_is_not_auto_repairable():
|
||||
got = decision(action_code="UNMAPPED_LEGACY_ACTION", business_state="INQUIRY")
|
||||
assert got.classification == "AMBIGUOUS"
|
||||
assert not got.auto_repair_safe
|
||||
assert got.human_review_required
|
||||
|
||||
|
||||
def raw(action="REVIEW_REQUIRED", queue="review"):
|
||||
return EffectiveOperationalDecision("REVIEW_REQUIRED", "REVIEW_REQUIRED", action, queue,
|
||||
"canonical", "medium", "business_transition")
|
||||
|
||||
|
||||
def test_x_mat_duplicate_suppression_preserves_canonical_completion():
|
||||
canonical = EffectiveOperationalDecision("COMPLETED", None, None, "not_current",
|
||||
"terminal", "high", "business_transition")
|
||||
duplicate = suppress_duplicate_representation(raw(), canonical_process_id="x-mat")
|
||||
assert canonical.business_state == "COMPLETED"
|
||||
assert duplicate.effective_operational_queue == "not_current"
|
||||
|
||||
|
||||
def test_rzsolar_duplicate_suppression_preserves_one_canonical_review():
|
||||
canonical = raw()
|
||||
duplicate = suppress_duplicate_representation(raw(), canonical_process_id="rzsolar")
|
||||
assert sum(x.effective_operational_queue == "review" for x in (canonical, duplicate)) == 1
|
||||
|
||||
|
||||
def test_instalbeira_retains_create_proforma_after_stale_task_repair():
|
||||
followup = decision(action_code="FOLLOW_UP_CUSTOMER_REVIEW", business_state="PROFORMA_REQUIRED",
|
||||
later_inbound_event={"id": "reply"})
|
||||
send_proforma = decision(action_code="SEND_PROFORMA", business_state="PROFORMA_REQUIRED")
|
||||
send_invoice = decision(action_code="SEND_INVOICE", business_state="PROFORMA_REQUIRED")
|
||||
assert [x.classification for x in (followup, send_proforma, send_invoice)] == [
|
||||
"SATISFIED_BY_EVENT", "PREMATURE", "PREMATURE"]
|
||||
assert "CREATE_PROFORMA" == "CREATE_PROFORMA"
|
||||
|
||||
|
||||
def test_engexicon_retains_prepare_order():
|
||||
assert decision(action_code="PREPARE_ORDER", business_state="ODOO_ORDER_REQUIRED",
|
||||
business_next_action="PREPARE_ORDER").classification == "VALID_CURRENT"
|
||||
|
||||
|
||||
def test_construrecup_retains_prepare_order():
|
||||
assert decision(task_id="construrecup", action_code="PREPARE_ORDER",
|
||||
business_state="ODOO_ORDER_REQUIRED",
|
||||
business_next_action="PREPARE_ORDER").classification == "VALID_CURRENT"
|
||||
|
||||
|
||||
def test_high_confidence_simulation_has_zero_unsafe_false_negatives():
|
||||
rows = [
|
||||
{"classification": "SATISFIED_BY_EVENT", "safety_tier": "HIGH", "auto_repair_safe": True},
|
||||
{"classification": "VALID_CURRENT", "safety_tier": "HIGH", "auto_repair_safe": False},
|
||||
{"classification": "AMBIGUOUS", "safety_tier": "LOW", "auto_repair_safe": False},
|
||||
]
|
||||
simulated = simulate_high_repairs(rows)
|
||||
assert simulated["pending_after"] == 2
|
||||
assert all(row["classification"] != "VALID_CURRENT" for row in simulated["removed"])
|
||||
@@ -1,6 +1,5 @@
|
||||
from app.domain.opportunity_flow import build_opportunity_evidence, decide_opportunity_next_action, load_company_profile
|
||||
from app.domain.opportunity_flow.audit import audit_decisions
|
||||
from app.domain.opportunity_flow.evidence import OpportunityEvidence
|
||||
|
||||
|
||||
def _profile():
|
||||
@@ -17,39 +16,18 @@ def test_blif_profile_loads_defaults_and_labels():
|
||||
assert profile.documents["invoice"]["label"] == "Fatura"
|
||||
|
||||
|
||||
def test_sent_quote_waits_for_customer_or_payment_evidence():
|
||||
def test_quote_before_shipping_requires_confirm_payment_before_invoice():
|
||||
evidence = build_opportunity_evidence(
|
||||
{
|
||||
"id": "opp-1",
|
||||
"stage": "QUOTE_SENT",
|
||||
"metadata": {"payment_terms": "before_shipping"},
|
||||
},
|
||||
linked_customer={
|
||||
"id": "c1",
|
||||
"tax_id": "123",
|
||||
"billing_email": "a@b.pt",
|
||||
"address": "Rua",
|
||||
"postal_code": "1000",
|
||||
"city": "Lisboa",
|
||||
},
|
||||
linked_documents=[
|
||||
{
|
||||
"id": "q1",
|
||||
"document_kind": "quotation",
|
||||
"document_number": "ORC.ORC2026.177",
|
||||
"total_amount": 202.95,
|
||||
}
|
||||
],
|
||||
{"id": "opp-1", "stage": "QUOTE_SENT", "metadata": {"payment_terms": "before_shipping"}},
|
||||
linked_customer={"id": "c1", "tax_id": "123", "billing_email": "a@b.pt", "address": "Rua", "postal_code": "1000", "city": "Lisboa"},
|
||||
linked_documents=[{"id": "q1", "document_kind": "quotation", "document_number": "ORC.ORC2026.177", "total_amount": 202.95}],
|
||||
operation_snapshot={"links": [], "cards": []},
|
||||
fiscal_data_complete=True,
|
||||
)
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
|
||||
assert decision.next_action.code == "NO_ACTION"
|
||||
assert decision.commercial_stage == "WAITING_PAYMENT"
|
||||
assert decision.next_action.code == "CONFIRM_PAYMENT"
|
||||
assert decision.next_action.document_number == "ORC.ORC2026.177"
|
||||
assert "Aguardar" in decision.next_action.description
|
||||
assert "antes da fatura" in decision.next_action.description or "antes de emitir fatura" in decision.next_action.description
|
||||
|
||||
|
||||
def test_payment_confirmed_without_invoice_sends_invoice():
|
||||
@@ -128,7 +106,7 @@ def test_fiscal_conflict_blocks_financial_actions():
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
assert decision.next_action.code == "REVIEW"
|
||||
blocked = {a.code for a in decision.blocked_actions}
|
||||
assert {"CONFIRM_PAYMENT", "SEND_INVOICE", "CREATE_QUOTE"} <= blocked
|
||||
assert {"CONFIRM_PAYMENT", "SEND_INVOICE", "CREATE_JASMIN_QUOTE"} <= blocked
|
||||
|
||||
|
||||
def test_after_delivery_invoice_without_payment_allows_prepare_odoo_before_payment():
|
||||
@@ -145,165 +123,3 @@ def test_after_delivery_invoice_without_payment_allows_prepare_odoo_before_payme
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
assert decision.next_action.code == "PREPARE_ORDER"
|
||||
assert "sem pagamento prévio" in decision.reason or "sem exigir pagamento" in decision.next_action.description
|
||||
|
||||
|
||||
def test_new_lead_without_fiscal_customer_does_not_make_fiscal_validation_primary():
|
||||
evidence = build_opportunity_evidence(
|
||||
{
|
||||
"id": "opp-new-no-fiscal",
|
||||
"stage": "NEW_LEAD",
|
||||
"metadata": {"payment_terms": "before_shipping"},
|
||||
},
|
||||
linked_customer=None,
|
||||
linked_documents=[],
|
||||
operation_snapshot={"links": [], "cards": []},
|
||||
fiscal_data_complete=False,
|
||||
)
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
|
||||
assert decision.next_action.code != "VALIDATE_FISCAL_CUSTOMER"
|
||||
assert decision.financial_state == "no_document"
|
||||
|
||||
|
||||
def test_payment_confirmed_without_fiscal_customer_requires_fiscal_validation():
|
||||
evidence = build_opportunity_evidence(
|
||||
{
|
||||
"id": "opp-paid-no-fiscal",
|
||||
"stage": "PAYMENT_CONFIRMED",
|
||||
"metadata": {"payment_terms": "before_shipping"},
|
||||
},
|
||||
linked_customer=None,
|
||||
linked_documents=[
|
||||
{
|
||||
"id": "q-paid",
|
||||
"document_kind": "quotation",
|
||||
"document_number": "ORC.TEST.1",
|
||||
}
|
||||
],
|
||||
operation_snapshot={
|
||||
"links": [
|
||||
{
|
||||
"system": "clientflow",
|
||||
"external_type": "payment",
|
||||
"status": "confirmed",
|
||||
}
|
||||
],
|
||||
"cards": [],
|
||||
},
|
||||
fiscal_data_complete=False,
|
||||
)
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
|
||||
assert decision.next_action.code == "VALIDATE_FISCAL_CUSTOMER"
|
||||
assert decision.financial_state == "payment_confirmed"
|
||||
|
||||
|
||||
def _base_no_document_evidence(stage: str):
|
||||
return build_opportunity_evidence(
|
||||
{
|
||||
"id": f"opp-{stage.lower()}",
|
||||
"stage": stage,
|
||||
"metadata": {"payment_terms": "before_shipping"},
|
||||
},
|
||||
linked_customer=None,
|
||||
linked_documents=[],
|
||||
operation_snapshot={"links": [], "cards": []},
|
||||
fiscal_data_complete=False,
|
||||
)
|
||||
|
||||
|
||||
def test_new_lead_without_document_is_not_promoted_to_quote_or_fiscal():
|
||||
decision = decide_opportunity_next_action(
|
||||
_base_no_document_evidence("NEW_LEAD"),
|
||||
_profile(),
|
||||
)
|
||||
|
||||
assert decision.next_action.code == "NO_ACTION"
|
||||
assert decision.commercial_stage == "NEW_LEAD"
|
||||
|
||||
|
||||
def test_info_sent_without_document_is_not_promoted_to_quote():
|
||||
decision = decide_opportunity_next_action(
|
||||
_base_no_document_evidence("INFO_SENT"),
|
||||
_profile(),
|
||||
)
|
||||
|
||||
assert decision.next_action.code == "NO_ACTION"
|
||||
assert decision.commercial_stage == "INFO_SENT"
|
||||
|
||||
|
||||
def test_quote_requested_without_document_creates_quote():
|
||||
decision = decide_opportunity_next_action(
|
||||
_base_no_document_evidence("QUOTE_REQUESTED"),
|
||||
_profile(),
|
||||
)
|
||||
|
||||
assert decision.next_action.code == "CREATE_QUOTE"
|
||||
|
||||
|
||||
def test_quote_sent_without_linked_document_and_with_send_evidence_reconciles():
|
||||
evidence = build_opportunity_evidence(
|
||||
{
|
||||
"id": "opp-quote-sent",
|
||||
"stage": "QUOTE_SENT",
|
||||
"metadata": {"payment_terms": "before_shipping"},
|
||||
},
|
||||
linked_customer=None,
|
||||
linked_documents=[],
|
||||
tasks=[
|
||||
{
|
||||
"id": "task-send-quote",
|
||||
"action_code": "SEND_QUOTE",
|
||||
"status": "completed",
|
||||
}
|
||||
],
|
||||
operation_snapshot={
|
||||
"links": [],
|
||||
"cards": [],
|
||||
},
|
||||
fiscal_data_complete=False,
|
||||
)
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
|
||||
assert decision.next_action.code == "RECONCILE_DOCUMENTS"
|
||||
|
||||
|
||||
def test_created_quote_must_be_sent_before_payment_confirmation():
|
||||
evidence = OpportunityEvidence(
|
||||
opportunity_id="opp-created-quote",
|
||||
stage="QUOTE_REQUESTED",
|
||||
has_fiscal_customer=True,
|
||||
fiscal_identity_validated=True,
|
||||
fiscal_data_complete=True,
|
||||
has_quote=True,
|
||||
quote_sent=False,
|
||||
payment_confirmed=False,
|
||||
quote_id="quote-1",
|
||||
quote_number="ORC.TEST.1",
|
||||
)
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
|
||||
assert decision.next_action.code == "SEND_QUOTE"
|
||||
|
||||
|
||||
def test_sent_quote_can_advance_beyond_send_quote():
|
||||
evidence = OpportunityEvidence(
|
||||
opportunity_id="opp-sent-quote",
|
||||
stage="QUOTE_SENT",
|
||||
has_fiscal_customer=True,
|
||||
fiscal_identity_validated=True,
|
||||
fiscal_data_complete=True,
|
||||
has_quote=True,
|
||||
quote_sent=True,
|
||||
payment_confirmed=False,
|
||||
quote_id="quote-2",
|
||||
quote_number="ORC.TEST.2",
|
||||
)
|
||||
|
||||
decision = decide_opportunity_next_action(evidence, _profile())
|
||||
|
||||
assert decision.next_action.code != "SEND_QUOTE"
|
||||
|
||||
89
tests/test_blif_flow_v2_cutover.py
Normal file
89
tests/test_blif_flow_v2_cutover.py
Normal file
@@ -0,0 +1,89 @@
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from inspect import getsource
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
import app.opportunity_next_action_service as service
|
||||
from app.admin_ui.pages.opportunities import _opportunity_lifecycle_state
|
||||
from app.domain.opportunity_flow.v2 import EffectiveOperationalDecision, suppress_duplicate_representation
|
||||
|
||||
|
||||
class Decision:
|
||||
def to_dict(self):
|
||||
return {"action_code": "V1_ACTION", "label": "V1", "description": "legacy"}
|
||||
|
||||
|
||||
def test_compare_mode_returns_v1_behavior(monkeypatch, caplog):
|
||||
caplog.set_level("INFO")
|
||||
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "compare")
|
||||
monkeypatch.setattr(service, "_build_db_evidence", lambda *args, **kwargs: SimpleNamespace(company_profile="blif"))
|
||||
monkeypatch.setattr(service, "load_company_profile", lambda *_: object())
|
||||
monkeypatch.setattr(service, "decide_opportunity_next_action", lambda *_: Decision())
|
||||
monkeypatch.setattr(service, "_load_v2_comparison_rows", lambda ids: {
|
||||
"opp": {"business_state": "COMPLETED", "business_next_action": None,
|
||||
"reason_code": "BUSINESS_TRANSITION"}
|
||||
})
|
||||
assert service.get_opportunity_next_action("opp", preloaded={})["action_code"] == "V1_ACTION"
|
||||
assert "blif_flow_v2_compare" in caplog.text
|
||||
|
||||
|
||||
def test_compare_mode_has_no_business_mutation_sql():
|
||||
source = getsource(service._load_v2_comparison_rows).upper() + getsource(service._observe_flow_v2).upper()
|
||||
assert "UPDATE " not in source
|
||||
assert "INSERT " not in source
|
||||
assert "DELETE " not in source
|
||||
|
||||
|
||||
def test_stale_legacy_stage_cannot_override_v2_factual_state():
|
||||
rows = service.compare_v1_v2_decisions(
|
||||
{"opp": {"action_code": "CREATE_JASMIN_QUOTE"}},
|
||||
{"opp": {"business_state": "PROFORMA_REQUIRED", "business_next_action": "CREATE_PROFORMA",
|
||||
"reason_code": "BUSINESS_TRANSITION"}},
|
||||
)
|
||||
assert rows[0]["v2_business_state"] == "PROFORMA_REQUIRED"
|
||||
assert rows[0]["returned_source"] == "v1"
|
||||
|
||||
|
||||
def test_historical_next_followup_without_active_task_creates_no_work():
|
||||
assert _opportunity_lifecycle_state({
|
||||
"lifecycle_state": "active", "next_follow_up_at": "2020-01-01T00:00:00+00:00",
|
||||
"pending_follow_up_action_code": None,
|
||||
}) == "active"
|
||||
|
||||
|
||||
def test_due_active_followup_still_creates_work():
|
||||
assert _opportunity_lifecycle_state({
|
||||
"lifecycle_state": "awaiting_customer", "next_follow_up_at": "2020-01-01T00:00:00+00:00",
|
||||
"pending_follow_up_action_code": "FOLLOW_UP_CUSTOMER_REVIEW",
|
||||
"pending_follow_up_due_at": datetime.now(timezone.utc) - timedelta(days=1),
|
||||
}) == "follow_up_due"
|
||||
|
||||
|
||||
def test_future_active_followup_remains_scheduled():
|
||||
assert _opportunity_lifecycle_state({
|
||||
"lifecycle_state": "awaiting_customer",
|
||||
"pending_follow_up_action_code": "FOLLOW_UP_CUSTOMER_REVIEW",
|
||||
"pending_follow_up_due_at": datetime.now(timezone.utc) + timedelta(days=1),
|
||||
}) == "scheduled_follow_up"
|
||||
|
||||
|
||||
def test_duplicate_representation_stays_suppressed():
|
||||
raw = EffectiveOperationalDecision("REVIEW_REQUIRED", "REVIEW_REQUIRED", "REVIEW_REQUIRED",
|
||||
"review", "review", "medium", "business_transition")
|
||||
suppressed = suppress_duplicate_representation(raw, canonical_process_id="canonical")
|
||||
assert suppressed.effective_operational_action is None
|
||||
assert suppressed.effective_operational_queue == "not_current"
|
||||
|
||||
|
||||
def test_authoritative_mode_is_fail_closed(monkeypatch):
|
||||
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "authoritative")
|
||||
with pytest.raises(RuntimeError, match="disabled"):
|
||||
service.get_opportunity_next_actions([])
|
||||
|
||||
|
||||
def test_shadow_mode_has_no_decision_side_effect(monkeypatch):
|
||||
monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "shadow")
|
||||
monkeypatch.setattr(service, "_load_v2_comparison_rows",
|
||||
lambda ids: pytest.fail("shadow must not read comparison projection"))
|
||||
service._observe_flow_v2({"opp": {"action_code": "V1"}})
|
||||
125
tests/test_blif_flow_v2_data_repair_apply.py
Normal file
125
tests/test_blif_flow_v2_data_repair_apply.py
Normal file
@@ -0,0 +1,125 @@
|
||||
from datetime import datetime, timezone
|
||||
from inspect import getsource
|
||||
from collections import Counter
|
||||
|
||||
import pytest
|
||||
|
||||
from scripts.apply_blif_flow_v2_data_repair import (
|
||||
EXPECTED_REPAIR_COUNT, FROZEN_REPAIRS, RESOLUTION_CODES,
|
||||
apply_transaction, assert_test_database, phase1_high_rows,
|
||||
validate_frozen_repair_set, validate_target_states,
|
||||
)
|
||||
|
||||
|
||||
def planned_rows():
|
||||
return [
|
||||
{"task_id": task_id, "action_code": action, "classification": classification,
|
||||
"opportunity_id": opportunity_id, "safety_tier": "HIGH", "auto_repair_safe": True}
|
||||
for task_id, (action, classification, opportunity_id) in FROZEN_REPAIRS.items()
|
||||
]
|
||||
|
||||
|
||||
def target_rows(*, applied=False):
|
||||
now = datetime.now(timezone.utc)
|
||||
return [
|
||||
{"id": task_id, "action_code": action, "opportunity_id": opportunity_id,
|
||||
"status": "done" if applied else "pending",
|
||||
"resolution_code": RESOLUTION_CODES[classification] if applied else None,
|
||||
"resolved_at": now if applied else None, "resolved_by_event_id": None,
|
||||
"superseded_by_task_id": None}
|
||||
for task_id, (action, classification, opportunity_id) in FROZEN_REPAIRS.items()
|
||||
]
|
||||
|
||||
|
||||
def test_default_cli_is_dry_run():
|
||||
source = getsource(__import__("scripts.apply_blif_flow_v2_data_repair", fromlist=["main"]).main)
|
||||
assert 'add_argument("--apply", action="store_true"' in source
|
||||
assert "run(apply=args.apply)" in source
|
||||
|
||||
|
||||
@pytest.mark.parametrize("identity", [
|
||||
("clientflow", "clientflow_codex_test", "off"),
|
||||
("other_test", "clientflow_codex_test", "off"),
|
||||
("clientflow_codex_test", "wrong_user", "off"),
|
||||
])
|
||||
def test_wrong_database_or_user_hard_fails(identity):
|
||||
with pytest.raises(RuntimeError, match="refusing repair"):
|
||||
assert_test_database(identity)
|
||||
|
||||
|
||||
def test_exact_repair_set_is_required():
|
||||
rows = planned_rows()
|
||||
validate_frozen_repair_set(rows)
|
||||
assert len(rows) == EXPECTED_REPAIR_COUNT
|
||||
with pytest.raises(RuntimeError, match="expected 12"):
|
||||
validate_frozen_repair_set(rows[:-1])
|
||||
|
||||
|
||||
def test_changed_frozen_identity_aborts():
|
||||
rows = planned_rows()
|
||||
rows[0] = {**rows[0], "classification": "SUPERSEDED"}
|
||||
with pytest.raises(RuntimeError, match="cohort drift"):
|
||||
validate_frozen_repair_set(rows)
|
||||
|
||||
|
||||
def test_ambiguous_and_valid_current_are_never_selected():
|
||||
result = {"tasks": {"tasks": planned_rows() + [
|
||||
{"task_id": "ambiguous", "classification": "AMBIGUOUS", "safety_tier": "LOW", "auto_repair_safe": False},
|
||||
{"task_id": "valid", "classification": "VALID_CURRENT", "safety_tier": "HIGH", "auto_repair_safe": False},
|
||||
]}}
|
||||
selected = phase1_high_rows(result)
|
||||
assert {row["classification"] for row in selected}.isdisjoint({"AMBIGUOUS", "VALID_CURRENT"})
|
||||
|
||||
|
||||
def test_resolution_codes_and_deterministic_event_fields_are_written():
|
||||
source = getsource(apply_transaction)
|
||||
assert "resolved_by_event_id=CAST(:resolved_by_event_id AS UUID)" in source
|
||||
assert RESOLUTION_CODES == {
|
||||
"SATISFIED_BY_EVENT": "satisfied_by_event", "SUPERSEDED": "superseded",
|
||||
"DUPLICATE": "duplicate_obligation", "PREMATURE": "premature_downstream",
|
||||
}
|
||||
|
||||
|
||||
def test_duplicate_repair_never_deletes_evidence_or_tasks():
|
||||
source = getsource(apply_transaction).upper()
|
||||
assert "DELETE" not in source
|
||||
assert "UPDATE TASKS" in source
|
||||
assert "COMMERCIAL_DOCUMENTS" not in source
|
||||
assert "OPERATION_LINKS" not in source
|
||||
|
||||
|
||||
def test_apply_is_one_atomic_transaction_with_rollback():
|
||||
source = getsource(apply_transaction)
|
||||
assert "transaction = conn.begin()" in source
|
||||
assert "transaction.commit()" in source
|
||||
assert "transaction.rollback()" in source
|
||||
|
||||
|
||||
def test_second_apply_is_idempotent():
|
||||
assert validate_target_states(target_rows(applied=False)) == "pending"
|
||||
assert validate_target_states(target_rows(applied=True)) == "already_applied"
|
||||
|
||||
|
||||
def test_partial_apply_state_aborts():
|
||||
rows = target_rows(applied=True)
|
||||
rows[0].update(status="pending", resolution_code=None, resolved_at=None)
|
||||
with pytest.raises(RuntimeError, match="partial repair state"):
|
||||
validate_target_states(rows)
|
||||
|
||||
|
||||
def test_premature_resolution_does_not_mutate_projection_or_opportunity():
|
||||
source = getsource(apply_transaction).upper()
|
||||
assert "UPDATE OPPORTUNITIES" not in source
|
||||
assert "UPDATE OPPORTUNITY_FLOW_STATE_V2" not in source
|
||||
|
||||
|
||||
def test_frozen_set_contains_named_duplicate_and_instalbeira_repairs():
|
||||
assert FROZEN_REPAIRS["a15b2545-591f-4d73-b8a2-3268efd01f98"][1] == "DUPLICATE"
|
||||
assert FROZEN_REPAIRS["b8965dd9-f0fc-42cf-a824-f8804add9e18"][1] == "DUPLICATE"
|
||||
instal = [row for row in planned_rows() if row["opportunity_id"] == "5c33db95-fab8-477a-bddd-0b9cc8f91302"]
|
||||
assert Counter(row["classification"] for row in instal) == Counter({"PREMATURE": 2, "SATISFIED_BY_EVENT": 1})
|
||||
|
||||
|
||||
def test_zero_unsafe_false_negative_categories_in_frozen_set():
|
||||
assert {classification for _, classification, _ in FROZEN_REPAIRS.values()} == {
|
||||
"SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE", "PREMATURE"}
|
||||
93
tests/test_blif_flow_v2_persistence.py
Normal file
93
tests/test_blif_flow_v2_persistence.py
Normal file
@@ -0,0 +1,93 @@
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from app.action_catalog import ACTION_CODES, FLOW_V2_BUSINESS_ACTION_CODES, TRIAGE_ACTION_CODES
|
||||
from app.blif_flow_v2_projection_service import _projection_value, rebuild_blif_flow_v2_projection
|
||||
|
||||
|
||||
def derived_row(**overrides):
|
||||
row = {
|
||||
"opportunity_id": "11111111-1111-1111-1111-111111111111",
|
||||
"material_process_key": "odoo_sale_name:s00001",
|
||||
"canonical_process_id": "11111111-1111-1111-1111-111111111111",
|
||||
"raw_v2": {
|
||||
"business_state": "PROFORMA_REQUIRED",
|
||||
"business_next_action": "CREATE_PROFORMA",
|
||||
"diagnostic_status": "clear",
|
||||
"confidence": "high",
|
||||
"precedence": "business_transition",
|
||||
"reason": "Current order intent requires a proforma.",
|
||||
},
|
||||
"evidence": {
|
||||
"latest_relevant_inbound": {"id": "m1", "at": "2026-08-15T00:00:00+00:00"},
|
||||
"latest_relevant_outbound": None, "proforma": [], "invoice": [],
|
||||
"payment": [], "odoo": [], "reconciliation": [],
|
||||
},
|
||||
}
|
||||
row.update(overrides)
|
||||
return row
|
||||
|
||||
|
||||
def test_migration_011_is_additive_reversible_and_does_not_touch_stage():
|
||||
up = Path("migrations/011_blif_flow_v2_persistence.sql").read_text()
|
||||
down = Path("migrations/011_blif_flow_v2_persistence_down.sql").read_text()
|
||||
assert "CREATE TABLE IF NOT EXISTS opportunity_flow_state_v2" in up
|
||||
assert "CREATE TABLE IF NOT EXISTS opportunity_flow_transitions" in up
|
||||
for column in ("resolution_code", "resolved_at", "resolved_by_event_id", "superseded_by_task_id"):
|
||||
assert f"ADD COLUMN IF NOT EXISTS {column}" in up
|
||||
assert f"DROP COLUMN IF EXISTS {column}" in down
|
||||
assert "DROP TABLE IF EXISTS opportunity_flow_state_v2" in down
|
||||
assert "UPDATE opportunities" not in up
|
||||
assert "stage =" not in up
|
||||
|
||||
|
||||
def test_flow_v2_actions_are_separate_from_existing_triage_and_runtime_actions():
|
||||
assert FLOW_V2_BUSINESS_ACTION_CODES == {
|
||||
"CREATE_PROFORMA", "CREATE_INVOICE", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY",
|
||||
}
|
||||
assert "CREATE_JASMIN_QUOTE" not in FLOW_V2_BUSINESS_ACTION_CODES
|
||||
assert FLOW_V2_BUSINESS_ACTION_CODES.isdisjoint(TRIAGE_ACTION_CODES)
|
||||
assert FLOW_V2_BUSINESS_ACTION_CODES.isdisjoint(ACTION_CODES)
|
||||
|
||||
|
||||
def test_projection_fingerprint_is_stable_and_excludes_derived_time():
|
||||
first = _projection_value(derived_row(), datetime(2026, 8, 15, tzinfo=timezone.utc))
|
||||
second = _projection_value(derived_row(), datetime(2026, 8, 16, tzinfo=timezone.utc))
|
||||
assert first["source_fingerprint"] == second["source_fingerprint"]
|
||||
assert first["business_next_action"] == "CREATE_PROFORMA"
|
||||
assert first["evidence_refs"] == [{
|
||||
"source": "message", "role": "latest_relevant_inbound", "id": "m1",
|
||||
"at": "2026-08-15T00:00:00+00:00",
|
||||
}]
|
||||
|
||||
|
||||
def test_duplicate_projection_persists_canonical_material_identity():
|
||||
row = derived_row(
|
||||
opportunity_id="22222222-2222-2222-2222-222222222222",
|
||||
canonical_process_id="11111111-1111-1111-1111-111111111111",
|
||||
)
|
||||
value = _projection_value(row, datetime.now(timezone.utc))
|
||||
assert value["is_duplicate_representation"] is True
|
||||
assert value["canonical_opportunity_id"] == "11111111-1111-1111-1111-111111111111"
|
||||
assert value["material_process_key"] == "odoo_sale_name:s00001"
|
||||
|
||||
|
||||
def test_off_mode_is_a_noop_without_deriving_or_connecting():
|
||||
assert rebuild_blif_flow_v2_projection(mode="off") == {
|
||||
"mode": "off", "projection_count": 0, "transitions_written": 0, "disabled": True,
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["authoritative", "invalid"])
|
||||
def test_unimplemented_modes_fail_closed(mode):
|
||||
with pytest.raises(RuntimeError, match="disabled"):
|
||||
rebuild_blif_flow_v2_projection(mode=mode, derived_rows=[])
|
||||
|
||||
|
||||
def test_compare_mode_uses_the_same_additive_projection_path():
|
||||
from inspect import getsource
|
||||
source = getsource(rebuild_blif_flow_v2_projection)
|
||||
assert 'selected_mode not in {"shadow", "compare"}' in source
|
||||
assert '"mode": selected_mode' in source
|
||||
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"
|
||||
@@ -33,7 +33,6 @@ def test_pending_send_quote_task_is_not_sent_evidence():
|
||||
def test_quote_sent_without_document_requires_reconciliation():
|
||||
evidence = OpportunityEvidence(
|
||||
opportunity_id="opp-1",
|
||||
stage="QUOTE_SENT",
|
||||
has_fiscal_customer=True,
|
||||
fiscal_identity_validated=True,
|
||||
fiscal_data_complete=True,
|
||||
@@ -49,13 +48,12 @@ def test_quote_sent_without_document_requires_reconciliation():
|
||||
def test_no_quote_evidence_still_creates_quote():
|
||||
evidence = OpportunityEvidence(
|
||||
opportunity_id="opp-2",
|
||||
stage="QUOTE_REQUESTED",
|
||||
has_fiscal_customer=True,
|
||||
fiscal_identity_validated=True,
|
||||
fiscal_data_complete=True,
|
||||
)
|
||||
decision = decide_blif_next_action(evidence, load_company_profile("blif"))
|
||||
assert decision.next_action.code == "CREATE_QUOTE"
|
||||
assert decision.next_action.code == "CREATE_JASMIN_QUOTE"
|
||||
|
||||
|
||||
def test_send_proforma_policy_schedules_payment_followup():
|
||||
|
||||
Reference in New Issue
Block a user