diff --git a/app/domain/opportunity_flow/repair.py b/app/domain/opportunity_flow/repair.py new file mode 100644 index 0000000..4637e7b --- /dev/null +++ b/app/domain/opportunity_flow/repair.py @@ -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} diff --git a/scripts/plan_blif_flow_v2_data_repair.py b/scripts/plan_blif_flow_v2_data_repair.py new file mode 100644 index 0000000..37045de --- /dev/null +++ b/scripts/plan_blif_flow_v2_data_repair.py @@ -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() diff --git a/tests/domain/opportunity_flow/test_blif_flow_v2_data_repair.py b/tests/domain/opportunity_flow/test_blif_flow_v2_data_repair.py new file mode 100644 index 0000000..5e92f66 --- /dev/null +++ b/tests/domain/opportunity_flow/test_blif_flow_v2_data_repair.py @@ -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"])