#!/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()