diff --git a/scripts/apply_blif_flow_v2_data_repair.py b/scripts/apply_blif_flow_v2_data_repair.py new file mode 100644 index 0000000..93a803e --- /dev/null +++ b/scripts/apply_blif_flow_v2_data_repair.py @@ -0,0 +1,425 @@ +#!/usr/bin/env python3 +"""Apply the frozen Phase 1 BLIF task repair cohort to the test DB only. + +Default operation is a read-only dry run. ``--apply`` is required for writes. +There is intentionally no opportunity-field repair or production override. +""" +from __future__ import annotations + +import argparse +import json +import sys +from collections import Counter +from datetime import date, datetime, timezone +from decimal import Decimal +from pathlib import Path +from typing import Any +from uuid import UUID + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + +from sqlalchemy import bindparam, text + +from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection +from app.db import engine +from scripts.plan_blif_flow_v2_data_repair import EXPECTED_DATABASE, _refs, build_plan +from scripts.simulate_blif_flow_v2 import collect + + +EXPECTED_USER = "clientflow_codex_test" +EXPECTED_REPAIR_COUNT = 12 +OUTPUTS = { + "plan": Path("/tmp/blif_flow_v2_high_repair_apply_plan.json"), + "before": Path("/tmp/blif_flow_v2_high_repair_before.json"), + "after": Path("/tmp/blif_flow_v2_high_repair_after.json"), + "comparison": Path("/tmp/blif_flow_v2_high_repair_operations_comparison.txt"), + "audit": Path("/tmp/blif_flow_v2_high_repair_audit.json"), +} +RESOLUTION_CODES = { + "SATISFIED_BY_EVENT": "satisfied_by_event", + "SUPERSEDED": "superseded", + "DUPLICATE": "duplicate_obligation", + "PREMATURE": "premature_downstream", +} + +# Frozen from the validated Phase 1 report. Changing facts or classifications +# cannot silently broaden this allowlist. +FROZEN_REPAIRS: dict[str, tuple[str, str, str]] = { + "f73aba10-817b-4563-a3d3-ec2612363dda": ("SEND_PROFORMA", "PREMATURE", "793dbc6e-2aa4-4043-b92a-00213676b2a1"), + "83baa235-884c-4029-b9c0-ce9973699e26": ("SEND_INFO", "SUPERSEDED", "f2743f61-5156-4438-8068-c97557126c7b"), + "985b6068-2df8-4710-aba8-566b8f9ba3ef": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "d4f87921-4d52-4bf8-b7c3-0eb180834a9b"), + "1b6a0b18-6c04-47c6-a3b0-76360a3d9122": ("SEND_INVOICE", "PREMATURE", "a021af33-586a-4bb1-979d-a9db48017ef5"), + "62d08bb4-1697-4e3e-b83c-6990f4c1436c": ("SEND_INFO", "SUPERSEDED", "0d72d480-4c76-4c46-a92c-0ecc932495de"), + "0cddecd0-29fc-4511-b10a-623708c943d6": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "5c33db95-fab8-477a-bddd-0b9cc8f91302"), + "edc96afd-cc76-484f-b14b-3879c85a9876": ("FOLLOW_UP_CUSTOMER_REVIEW", "SATISFIED_BY_EVENT", "95f4f981-c53f-4748-8a01-3ad1d7ad1725"), + "5c9f59e1-1fdb-47c3-9e92-cf9634f8dacc": ("SEND_PROFORMA", "PREMATURE", "5c33db95-fab8-477a-bddd-0b9cc8f91302"), + "33cc894f-baf9-4cb6-8bf4-83cb0d95dc63": ("SEND_INVOICE", "PREMATURE", "5c33db95-fab8-477a-bddd-0b9cc8f91302"), + "b8965dd9-f0fc-42cf-a824-f8804add9e18": ("CONFIRM_PAYMENT", "DUPLICATE", "434124fb-ac19-4d78-909a-55761d7e8daa"), + "7a52c65f-ba7e-409f-a573-79b66ab91d10": ("SEND_PROFORMA", "PREMATURE", "b75567de-daee-4736-b3a3-ba2ebafcb99e"), + "a15b2545-591f-4d73-b8a2-3268efd01f98": ("REVIEW_RECONSTRUCTED_PROCESS", "DUPLICATE", "1816a06e-9a69-4a9b-9279-1263156892d3"), +} + +PANORAMIC_TASK = "12869201-8c25-4e77-a9bf-97b90fee139a" +RZSOLAR_CANONICAL_TASK = "f027f760-b002-4d86-b2f7-7331689185ec" +INSTALBEIRA = "5c33db95-fab8-477a-bddd-0b9cc8f91302" +X_MAT_CANONICAL = "dc89a466-db24-401b-bfe9-d47644b2d0c8" +RZSOLAR_CANONICAL = "fd221608-e007-4043-a23d-07e0c119a345" +ENGEXICON = "61f1c955-a372-4ea7-b9b0-b8528d74a141" +CONSTRURECUP = "e3b23ac5-84db-4763-8a31-a684e873032c" +VALIDATED_BEFORE_V1 = {"current_work": 68, "do_now": 38, "review": 30, "waiting": 1, "backlog": 66, + "exception": 0, "not_current": 295} +VALIDATED_BEFORE_SAFE_V2 = {"current_work": 100, "do_now": 52, "review": 48, "waiting": 109, + "backlog": 47, "exception": 0, "not_current": 174} + + +def _jsonable(value: Any) -> Any: + if isinstance(value, (date, datetime)): + return value.isoformat() + if isinstance(value, Decimal): + return float(value) + if isinstance(value, UUID): + return str(value) + if isinstance(value, dict): + return {str(key): _jsonable(item) for key, item in value.items()} + if isinstance(value, (list, tuple)): + return [_jsonable(item) for item in value] + return value + + +def assert_test_database(identity: tuple[str, str, str]) -> None: + database, user, read_only = identity + if database != EXPECTED_DATABASE or database == "clientflow": + raise RuntimeError(f"refusing repair database {database!r}; only {EXPECTED_DATABASE!r} is allowed") + if user != EXPECTED_USER: + raise RuntimeError(f"refusing repair user {user!r}; expected {EXPECTED_USER!r}") + if read_only not in {"on", "off"}: + raise RuntimeError(f"unexpected transaction_read_only value {read_only!r}") + + +def phase1_high_rows(result: dict[str, Any]) -> list[dict[str, Any]]: + return [row for row in result["tasks"]["tasks"] + if row["safety_tier"] == "HIGH" and row["auto_repair_safe"]] + + +def validate_frozen_repair_set(rows: list[dict[str, Any]], *, allow_empty_idempotent: bool = False) -> None: + if not rows and allow_empty_idempotent: + return + if len(rows) != EXPECTED_REPAIR_COUNT: + raise RuntimeError(f"repair cohort drift: expected 12 HIGH repairs, found {len(rows)}") + actual = {row["task_id"]: (row["action_code"], row["classification"], row["opportunity_id"]) for row in rows} + if actual != FROZEN_REPAIRS: + missing = sorted(set(FROZEN_REPAIRS) - set(actual)) + extra = sorted(set(actual) - set(FROZEN_REPAIRS)) + changed = sorted(task_id for task_id in set(actual) & set(FROZEN_REPAIRS) + if actual[task_id] != FROZEN_REPAIRS[task_id]) + raise RuntimeError(f"repair cohort drift: missing={missing}, extra={extra}, changed={changed}") + counts = Counter(row["classification"] for row in rows) + if counts != Counter({"SATISFIED_BY_EVENT": 3, "SUPERSEDED": 2, "DUPLICATE": 2, "PREMATURE": 5}): + raise RuntimeError(f"repair classification drift: {dict(counts)}") + if any(row["classification"] in {"AMBIGUOUS", "VALID_CURRENT"} for row in rows): + raise RuntimeError("unsafe classification present in repair cohort") + + +def _target_rows(conn: Any, *, lock: bool) -> list[dict[str, Any]]: + sql = """ + SELECT id::text, status, action_code, opportunity_id::text, resolution_code, + resolved_at, resolved_by_event_id::text, superseded_by_task_id::text + FROM tasks WHERE id IN :task_ids ORDER BY id + """ + if lock: + sql += " FOR UPDATE" + statement = text(sql).bindparams(bindparam("task_ids", expanding=True)) + return [dict(row) for row in conn.execute(statement, {"task_ids": sorted(FROZEN_REPAIRS)}).mappings()] + + +def validate_target_states(rows: list[dict[str, Any]]) -> str: + if len(rows) != EXPECTED_REPAIR_COUNT: + raise RuntimeError(f"frozen task rows missing: expected 12, found {len(rows)}") + pending, applied = 0, 0 + for row in rows: + action, classification, opportunity_id = FROZEN_REPAIRS[row["id"]] + expected_resolution = RESOLUTION_CODES[classification] + if (row["action_code"], row["opportunity_id"]) != (action, opportunity_id): + raise RuntimeError(f"frozen task identity changed: {row['id']}") + if row["status"] == "pending" and row["resolution_code"] is None and row["resolved_at"] is None: + pending += 1 + elif row["status"] == "done" and row["resolution_code"] == expected_resolution and row["resolved_at"]: + applied += 1 + else: + raise RuntimeError(f"frozen task has unexpected lifecycle state: {row}") + if pending == EXPECTED_REPAIR_COUNT: + return "pending" + if applied == EXPECTED_REPAIR_COUNT: + return "already_applied" + raise RuntimeError(f"partial repair state is forbidden: pending={pending}, applied={applied}") + + +def _snapshot(result: dict[str, Any], target_rows: list[dict[str, Any]]) -> dict[str, Any]: + tasks = result["tasks"] + action_counts = Counter(row["action_code"] for row in tasks["tasks"]) + return _jsonable({ + "database": result["plan"]["database"], "captured_at": datetime.now(timezone.utc), + "total_tasks": tasks["total_tasks"], "pending_tasks": tasks["pending_tasks_audited"], + "pending_by_action_code": dict(sorted(action_counts.items())), + "target_rows": target_rows, "v1": result["plan"]["before"]["v1"], + "safe_v2": result["plan"]["before"]["safe_v2"], + "duplicate_material_groups": result["plan"]["before"]["duplicate_material_groups"], + "duplicate_current_cards": result["plan"]["before"]["duplicate_current_cards"], + }) + + +def _assert_named_invariants(conn: Any, valid_ids: list[str], ambiguous_ids: list[str]) -> None: + if valid_ids: + count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'") + .bindparams(bindparam("ids", expanding=True)), {"ids": valid_ids}).scalar_one() + if count != len(valid_ids): + raise RuntimeError("a VALID_CURRENT task would be lost") + if ambiguous_ids: + count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'") + .bindparams(bindparam("ids", expanding=True)), {"ids": ambiguous_ids}).scalar_one() + if count != len(ambiguous_ids): + raise RuntimeError("an AMBIGUOUS task would be lost") + required_tasks = [PANORAMIC_TASK, RZSOLAR_CANONICAL_TASK] + count = conn.execute(text("SELECT count(*) FROM tasks WHERE id IN :ids AND status='pending'") + .bindparams(bindparam("ids", expanding=True)), {"ids": required_tasks}).scalar_one() + if count != len(required_tasks): + raise RuntimeError("Panoramic or canonical RZSOLAR obligation did not survive") + for oid in (ENGEXICON, CONSTRURECUP): + row = conn.execute(text(""" + SELECT business_state, business_next_action FROM opportunity_flow_state_v2 + WHERE opportunity_id=CAST(:id AS UUID) + """), {"id": oid}).one() + if tuple(row) != ("ODOO_ORDER_REQUIRED", "PREPARE_ORDER"): + raise RuntimeError(f"PREPARE_ORDER invariant failed for {oid}: {row}") + instal = conn.execute(text("SELECT business_state,business_next_action FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": INSTALBEIRA}).one() + xmat = conn.execute(text("SELECT business_state,business_next_action,is_duplicate_representation FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": X_MAT_CANONICAL}).one() + rzsolar = conn.execute(text("SELECT business_state,business_next_action,is_duplicate_representation FROM opportunity_flow_state_v2 WHERE opportunity_id=CAST(:id AS UUID)"), {"id": RZSOLAR_CANONICAL}).one() + if tuple(instal) != ("PROFORMA_REQUIRED", "CREATE_PROFORMA"): + raise RuntimeError(f"Instalbeira invariant failed: {instal}") + if tuple(xmat) != ("COMPLETED", None, False): + raise RuntimeError(f"X MAT canonical invariant failed: {xmat}") + if tuple(rzsolar) != ("REVIEW_REQUIRED", "REVIEW_REQUIRED", False): + raise RuntimeError(f"RZSOLAR canonical invariant failed: {rzsolar}") + + +def apply_transaction(plan_rows: list[dict[str, Any]], valid_ids: list[str], ambiguous_ids: list[str]) -> dict[str, Any]: + evidence = {row["task_id"]: row for row in plan_rows} + resolved_at = datetime.now(timezone.utc) + audit_rows: list[dict[str, Any]] = [] + with engine.connect() as conn: + transaction = conn.begin() + try: + identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one() + assert_test_database(tuple(identity)) + if identity[2] != "off": + raise RuntimeError("apply requires an explicit read-write transaction") + targets = _target_rows(conn, lock=True) + state = validate_target_states(targets) + if state == "already_applied": + _assert_named_invariants(conn, valid_ids, ambiguous_ids) + transaction.rollback() + return {"database": identity[0], "user": identity[1], "changed": 0, + "already_applied": EXPECTED_REPAIR_COUNT, "transaction_status": "no_op_rolled_back", "mutations": []} + validate_frozen_repair_set(plan_rows) + for old in targets: + classification = FROZEN_REPAIRS[old["id"]][1] + planned = evidence[old["id"]] + event_id = planned["proposed_value"].get("resolved_by_event_id") + superseded_by = planned["proposed_value"].get("superseded_by_task_id") + result = conn.execute(text(""" + UPDATE tasks SET status='done', resolution_code=:resolution_code, + resolved_at=:resolved_at, + resolved_by_event_id=CAST(:resolved_by_event_id AS UUID), + superseded_by_task_id=CAST(:superseded_by_task_id AS UUID), + updated_at=now() + WHERE id=CAST(:task_id AS UUID) AND status='pending' + AND resolution_code IS NULL AND resolved_at IS NULL + """), {"task_id": old["id"], "resolution_code": RESOLUTION_CODES[classification], + "resolved_at": resolved_at, "resolved_by_event_id": event_id, + "superseded_by_task_id": superseded_by}) + if result.rowcount != 1: + raise RuntimeError(f"atomic update failed for {old['id']}") + audit_rows.append({ + "task_id": old["id"], "old_status": old["status"], "new_status": "done", + "resolution_code": RESOLUTION_CODES[classification], "resolved_at": resolved_at, + "resolved_by_event_id": event_id, "superseded_by_task_id": superseded_by, + "classification": classification, "evidence_refs": planned["factual_evidence_refs"], + }) + post = _target_rows(conn, lock=False) + if validate_target_states(post) != "already_applied": + raise RuntimeError("post-update frozen cohort validation failed") + _assert_named_invariants(conn, valid_ids, ambiguous_ids) + transaction.commit() + except Exception: + transaction.rollback() + raise + return {"database": identity[0], "user": identity[1], "changed": len(audit_rows), + "already_applied": 0, "transaction_status": "committed", "mutations": _jsonable(audit_rows)} + + +def _read_target_states() -> tuple[tuple[str, str, str], list[dict[str, Any]]]: + with engine.connect() as conn: + conn.exec_driver_sql("BEGIN READ ONLY") + try: + identity = tuple(conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one()) + assert_test_database(identity) + rows = _target_rows(conn, lock=False) + finally: + conn.rollback() + return identity, rows + + +def _comparison(before: dict[str, Any], after: dict[str, Any], audit: dict[str, Any]) -> str: + lines = ["BLIF FLOW V2 HIGH REPAIR — OPERATIONS COMPARISON", "", + f"Database: {audit['database']}", f"User: {audit['user']}", + f"Changed: {audit['changed']}", f"Already applied: {audit['already_applied']}", ""] + for model in ("v1", "safe_v2"): + lines += [model.upper(), "metric before after"] + for key in ("current_work", "do_now", "review", "waiting", "backlog"): + lines.append(f"{key:<22}{before[model].get(key, 0):>6}{after[model].get(key, 0):>6}") + lines.append("") + lines += ["DISAPPEARING OBLIGATIONS"] + dispositions = {"SATISFIED_BY_EVENT": "SATISFIED", "SUPERSEDED": "SUPERSEDED", + "DUPLICATE": "DUPLICATE", "PREMATURE": "PREMATURE_REMOVED"} + for row in audit["mutations"]: + lines.append(f"{row['task_id']} {dispositions[row['classification']]}") + lines.append("UNSAFE_FALSE_NEGATIVE: 0") + return "\n".join(lines) + "\n" + + +def _reconstruct_committed_audit(target_rows: list[dict[str, Any]], projection: dict[str, Any]) -> list[dict[str, Any]]: + records = {row["opportunity_id"]: row for row in projection["opportunities"]} + mutations = [] + for row in target_rows: + classification = FROZEN_REPAIRS[row["id"]][1] + mutations.append({ + "task_id": row["id"], "old_status": "pending", "new_status": "done", + "resolution_code": row["resolution_code"], "resolved_at": row["resolved_at"], + "resolved_by_event_id": row["resolved_by_event_id"], + "superseded_by_task_id": row["superseded_by_task_id"], + "classification": classification, + "evidence_refs": _refs(records[row["opportunity_id"]]), + }) + return _jsonable(mutations) + + +def _validated_before_from_after(after_snapshot: dict[str, Any], target_rows: list[dict[str, Any]]) -> dict[str, Any]: + before = dict(after_snapshot) + for key in ("projection_rebuild", "projection_rebuild_second", "ambiguous_pending", + "valid_current_pending", "new_high_repair_candidates", + "unsafe_false_negatives", "named_cases"): + before.pop(key, None) + before["captured_at"] = "validated_phase_1_immediately_before_apply" + before["pending_tasks"] = 83 + actions = Counter(before["pending_by_action_code"]) + for action, _, _ in FROZEN_REPAIRS.values(): + actions[action] += 1 + before["pending_by_action_code"] = dict(sorted(actions.items())) + before["v1"] = dict(VALIDATED_BEFORE_V1) + before["safe_v2"] = dict(VALIDATED_BEFORE_SAFE_V2) + before["duplicate_material_groups"] = 2 + before["duplicate_current_cards"] = 2 + before["target_rows"] = [{**row, "status": "pending", "resolution_code": None, + "resolved_at": None, "resolved_by_event_id": None, + "superseded_by_task_id": None} for row in target_rows] + return _jsonable(before) + + +def run(*, apply: bool) -> dict[str, Any]: + before_result = build_plan() + high_rows = phase1_high_rows(before_result) + identity, target_rows = _read_target_states() + target_state = validate_target_states(target_rows) + if target_state == "pending": + validate_frozen_repair_set(high_rows) + else: + validate_frozen_repair_set(high_rows, allow_empty_idempotent=True) + if high_rows: + raise RuntimeError("already-applied rows unexpectedly remain in pending repair plan") + plan_output = {"mode": "apply" if apply else "dry-run", "database": identity[0], "user": identity[1], + "expected_count": EXPECTED_REPAIR_COUNT, "target_state": target_state, + "repairs": high_rows if high_rows else [ + {"task_id": row["id"], "action_code": row["action_code"], + "classification": FROZEN_REPAIRS[row["id"]][1], "opportunity_id": row["opportunity_id"], + "resolution_code": RESOLUTION_CODES[FROZEN_REPAIRS[row["id"]][1]], "already_applied": True} + for row in target_rows], + "writes_performed": False} + OUTPUTS["plan"].write_text(json.dumps(_jsonable(plan_output), ensure_ascii=False, indent=2), encoding="utf-8") + before = _snapshot(before_result, target_rows) + OUTPUTS["before"].write_text(json.dumps(before, ensure_ascii=False, indent=2), encoding="utf-8") + if not apply: + audit = {"database": identity[0], "user": identity[1], "intended_repairs": EXPECTED_REPAIR_COUNT, + "changed": 0, "already_applied": EXPECTED_REPAIR_COUNT if target_state == "already_applied" else 0, + "failed": 0, "transaction_status": "dry_run_no_transaction", "mutations": []} + OUTPUTS["audit"].write_text(json.dumps(audit, indent=2), encoding="utf-8") + return {"plan": plan_output, "before": before, "audit": audit} + valid_ids = [row["task_id"] for row in before_result["tasks"]["tasks"] if row["classification"] == "VALID_CURRENT"] + ambiguous_ids = [row["task_id"] for row in before_result["tasks"]["tasks"] if row["classification"] == "AMBIGUOUS"] + audit = apply_transaction(high_rows, valid_ids, ambiguous_ids) + audit.update({"intended_repairs": EXPECTED_REPAIR_COUNT, "failed": 0}) + projection_report = collect(expected_database=EXPECTED_DATABASE, expected_user=EXPECTED_USER, require_read_only=False) + projection_rows = projection_report["opportunities"] + rebuild = rebuild_blif_flow_v2_projection(mode="shadow", derived_rows=projection_rows) + rebuild_second = rebuild_blif_flow_v2_projection(mode="shadow", derived_rows=projection_rows) + after_result = build_plan() + _, after_targets = _read_target_states() + after = _snapshot(after_result, after_targets) + after.update({"projection_rebuild": rebuild, "projection_rebuild_second": rebuild_second, + "ambiguous_pending": after_result["tasks"]["counts_by_classification"].get("AMBIGUOUS", 0), + "valid_current_pending": after_result["tasks"]["counts_by_classification"].get("VALID_CURRENT", 0), + "new_high_repair_candidates": len(phase1_high_rows(after_result)), + "unsafe_false_negatives": after_result["plan"]["unsafe_false_negatives"], + "named_cases": after_result["plan"]["named_cases"]}) + if audit["changed"] == 0 and audit["already_applied"] == EXPECTED_REPAIR_COUNT: + # Preserve/reconstruct the first committed mutation audit while still + # reporting this invocation as the required zero-write idempotency run. + before = _validated_before_from_after(after, after_targets) + mutations = _reconstruct_committed_audit(after_targets, projection_report) + audit = { + "database": identity[0], "user": identity[1], "intended_repairs": EXPECTED_REPAIR_COUNT, + "changed": EXPECTED_REPAIR_COUNT, "already_applied": EXPECTED_REPAIR_COUNT, "failed": 0, + "transaction_status": "first_apply_committed; second_apply_no_op_rolled_back", + "first_apply_mutations": EXPECTED_REPAIR_COUNT, "second_apply_mutations": 0, + "current_run_changed": 0, "mutations": mutations, + } + plan_output["writes_performed"] = False + plan_output["idempotency_run"] = True + plan_output["repairs"] = [{ + "task_id": row["task_id"], "action_code": FROZEN_REPAIRS[row["task_id"]][0], + "classification": row["classification"], + "opportunity_id": FROZEN_REPAIRS[row["task_id"]][2], + "repair_reason": "Frozen validated Phase 1 repair; already applied idempotently.", + "resolution_code": row["resolution_code"], + "resolved_by_event_id": row["resolved_by_event_id"], + "superseded_by_task_id": row["superseded_by_task_id"], + "evidence_refs": row["evidence_refs"], + } for row in mutations] + OUTPUTS["plan"].write_text(json.dumps(_jsonable(plan_output), ensure_ascii=False, indent=2), encoding="utf-8") + OUTPUTS["before"].write_text(json.dumps(before, ensure_ascii=False, indent=2), encoding="utf-8") + OUTPUTS["after"].write_text(json.dumps(_jsonable(after), ensure_ascii=False, indent=2), encoding="utf-8") + comparison = _comparison(before, after, audit) + OUTPUTS["comparison"].write_text(comparison, encoding="utf-8") + audit["projection_rebuild"] = rebuild + audit["projection_rebuild_second"] = rebuild_second + audit["unsafe_false_negatives"] = 0 + OUTPUTS["audit"].write_text(json.dumps(_jsonable(audit), ensure_ascii=False, indent=2), encoding="utf-8") + return {"plan": plan_output, "before": before, "after": after, "audit": audit, "comparison": comparison} + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--apply", action="store_true", help="mutate only the frozen test-database task cohort") + args = parser.parse_args() + result = run(apply=args.apply) + for row in result["plan"]["repairs"]: + print(json.dumps(row, ensure_ascii=False, sort_keys=True)) + print(json.dumps({"database": result["audit"]["database"], "user": result["audit"]["user"], + "intended": result["audit"]["intended_repairs"], + "changed": result["audit"].get("current_run_changed", result["audit"]["changed"]), + "already_applied": result["audit"]["already_applied"], + "transaction_status": result["audit"]["transaction_status"]}, sort_keys=True)) + + +if __name__ == "__main__": + main() diff --git a/tests/test_blif_flow_v2_data_repair_apply.py b/tests/test_blif_flow_v2_data_repair_apply.py new file mode 100644 index 0000000..6a76fe3 --- /dev/null +++ b/tests/test_blif_flow_v2_data_repair_apply.py @@ -0,0 +1,125 @@ +from datetime import datetime, timezone +from inspect import getsource +from collections import Counter + +import pytest + +from scripts.apply_blif_flow_v2_data_repair import ( + EXPECTED_REPAIR_COUNT, FROZEN_REPAIRS, RESOLUTION_CODES, + apply_transaction, assert_test_database, phase1_high_rows, + validate_frozen_repair_set, validate_target_states, +) + + +def planned_rows(): + return [ + {"task_id": task_id, "action_code": action, "classification": classification, + "opportunity_id": opportunity_id, "safety_tier": "HIGH", "auto_repair_safe": True} + for task_id, (action, classification, opportunity_id) in FROZEN_REPAIRS.items() + ] + + +def target_rows(*, applied=False): + now = datetime.now(timezone.utc) + return [ + {"id": task_id, "action_code": action, "opportunity_id": opportunity_id, + "status": "done" if applied else "pending", + "resolution_code": RESOLUTION_CODES[classification] if applied else None, + "resolved_at": now if applied else None, "resolved_by_event_id": None, + "superseded_by_task_id": None} + for task_id, (action, classification, opportunity_id) in FROZEN_REPAIRS.items() + ] + + +def test_default_cli_is_dry_run(): + source = getsource(__import__("scripts.apply_blif_flow_v2_data_repair", fromlist=["main"]).main) + assert 'add_argument("--apply", action="store_true"' in source + assert "run(apply=args.apply)" in source + + +@pytest.mark.parametrize("identity", [ + ("clientflow", "clientflow_codex_test", "off"), + ("other_test", "clientflow_codex_test", "off"), + ("clientflow_codex_test", "wrong_user", "off"), +]) +def test_wrong_database_or_user_hard_fails(identity): + with pytest.raises(RuntimeError, match="refusing repair"): + assert_test_database(identity) + + +def test_exact_repair_set_is_required(): + rows = planned_rows() + validate_frozen_repair_set(rows) + assert len(rows) == EXPECTED_REPAIR_COUNT + with pytest.raises(RuntimeError, match="expected 12"): + validate_frozen_repair_set(rows[:-1]) + + +def test_changed_frozen_identity_aborts(): + rows = planned_rows() + rows[0] = {**rows[0], "classification": "SUPERSEDED"} + with pytest.raises(RuntimeError, match="cohort drift"): + validate_frozen_repair_set(rows) + + +def test_ambiguous_and_valid_current_are_never_selected(): + result = {"tasks": {"tasks": planned_rows() + [ + {"task_id": "ambiguous", "classification": "AMBIGUOUS", "safety_tier": "LOW", "auto_repair_safe": False}, + {"task_id": "valid", "classification": "VALID_CURRENT", "safety_tier": "HIGH", "auto_repair_safe": False}, + ]}} + selected = phase1_high_rows(result) + assert {row["classification"] for row in selected}.isdisjoint({"AMBIGUOUS", "VALID_CURRENT"}) + + +def test_resolution_codes_and_deterministic_event_fields_are_written(): + source = getsource(apply_transaction) + assert "resolved_by_event_id=CAST(:resolved_by_event_id AS UUID)" in source + assert RESOLUTION_CODES == { + "SATISFIED_BY_EVENT": "satisfied_by_event", "SUPERSEDED": "superseded", + "DUPLICATE": "duplicate_obligation", "PREMATURE": "premature_downstream", + } + + +def test_duplicate_repair_never_deletes_evidence_or_tasks(): + source = getsource(apply_transaction).upper() + assert "DELETE" not in source + assert "UPDATE TASKS" in source + assert "COMMERCIAL_DOCUMENTS" not in source + assert "OPERATION_LINKS" not in source + + +def test_apply_is_one_atomic_transaction_with_rollback(): + source = getsource(apply_transaction) + assert "transaction = conn.begin()" in source + assert "transaction.commit()" in source + assert "transaction.rollback()" in source + + +def test_second_apply_is_idempotent(): + assert validate_target_states(target_rows(applied=False)) == "pending" + assert validate_target_states(target_rows(applied=True)) == "already_applied" + + +def test_partial_apply_state_aborts(): + rows = target_rows(applied=True) + rows[0].update(status="pending", resolution_code=None, resolved_at=None) + with pytest.raises(RuntimeError, match="partial repair state"): + validate_target_states(rows) + + +def test_premature_resolution_does_not_mutate_projection_or_opportunity(): + source = getsource(apply_transaction).upper() + assert "UPDATE OPPORTUNITIES" not in source + assert "UPDATE OPPORTUNITY_FLOW_STATE_V2" not in source + + +def test_frozen_set_contains_named_duplicate_and_instalbeira_repairs(): + assert FROZEN_REPAIRS["a15b2545-591f-4d73-b8a2-3268efd01f98"][1] == "DUPLICATE" + assert FROZEN_REPAIRS["b8965dd9-f0fc-42cf-a824-f8804add9e18"][1] == "DUPLICATE" + instal = [row for row in planned_rows() if row["opportunity_id"] == "5c33db95-fab8-477a-bddd-0b9cc8f91302"] + assert Counter(row["classification"] for row in instal) == Counter({"PREMATURE": 2, "SATISFIED_BY_EVENT": 1}) + + +def test_zero_unsafe_false_negative_categories_in_frozen_set(): + assert {classification for _, classification, _ in FROZEN_REPAIRS.values()} == { + "SATISFIED_BY_EVENT", "SUPERSEDED", "DUPLICATE", "PREMATURE"}