feat: apply validated BLIF Flow v2 task repairs
This commit is contained in:
425
scripts/apply_blif_flow_v2_data_repair.py
Normal file
425
scripts/apply_blif_flow_v2_data_repair.py
Normal file
@@ -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()
|
||||
125
tests/test_blif_flow_v2_data_repair_apply.py
Normal file
125
tests/test_blif_flow_v2_data_repair_apply.py
Normal file
@@ -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"}
|
||||
Reference in New Issue
Block a user