447 lines
27 KiB
Python
447 lines
27 KiB
Python
#!/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()
|