feat: add BLIF Flow v2 data repair planner

This commit is contained in:
plx
2026-08-15 23:50:54 +00:00
parent e99da64d8b
commit acbbd1a84b
3 changed files with 718 additions and 0 deletions

View 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}

View 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()

View 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"])