1 Commits

Author SHA1 Message Date
plx
e456cbc0e4 fix: allow guarded production Flow v2 read-only audit 2026-08-16 00:53:29 +00:00
2 changed files with 400 additions and 190 deletions

View File

@@ -1,228 +1,328 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""Read-only compatibility/cutover audit for BLIF Flow v2.""" """Run a strictly read-only BLIF Flow v2 shadow/cutover audit."""
from __future__ import annotations from __future__ import annotations
import argparse
import json import json
import os
import sys import sys
from collections import Counter from collections import Counter
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any, Sequence
ROOT = Path(__file__).resolve().parents[1] ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT)) sys.path.insert(0, str(ROOT))
os.chdir(ROOT)
from sqlalchemy import text PRODUCTION_DATABASE = "clientflow"
TEST_DATABASE = "clientflow_codex_test"
AUDIT = Path("/tmp/blif_flow_v2_production_shadow_audit.json")
SEMANTIC = Path("/tmp/blif_flow_v2_production_semantic_compare.json")
SUMMARY = Path("/tmp/blif_flow_v2_production_shadow_summary.txt")
CLASSIFICATIONS = {
"SEMANTICALLY_EQUIVALENT", "V1_OPERATIONAL_OVERRIDE", "V2_CORRECTS_V1",
"LEGACY_ONLY", "REAL_CONFLICT", "MISSING_PROJECTION",
}
CURRENT_QUEUES = {"do_now", "review", "exception"}
OVERRIDE_PRECEDENCE = {
"scheduled_call", "due_followup", "future_followup", "integration_exception",
"document_prerequisite", "fiscal_prerequisite", "safe_preserve_v1",
}
LEGACY_ACTIONS = {
"CREATE_JASMIN_QUOTE", "NO_ACTION", "WAIT_CUSTOMER", "WAIT_PAYMENT",
"WAIT_PRODUCTION", "WAIT_LOGISTICS", "WAIT_SUPPLIER", "WAIT_SCHEDULED_DATE",
}
REVIEW_ACTIONS = {
"REVIEW", "REVIEW_REQUIRED", "REVIEW_MANUALLY", "REVIEW_RECONSTRUCTED_PROCESS",
"RECONCILE_DOCUMENTS", "VALIDATE_FISCAL_CUSTOMER", "REVIEW_EXCEPTION",
}
NAMED_IDS = {
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
"PANORAMIC": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
"ENGEXICON": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
"CONSTRURECUP": "e3b23ac5-84db-4763-8a31-a684e873032c",
"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",
}
from app.db import engine
from app.opportunity_next_action_service import ( def build_parser() -> argparse.ArgumentParser:
_load_v2_comparison_rows, compare_v1_v2_decisions, get_opportunity_next_actions, parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument(
"--production-readonly-audit", action="store_true",
help="explicitly authorize a read-only audit of database clientflow",
) )
from scripts.simulate_blif_flow_v2 import collect return parser
EXPECTED_DATABASE = "clientflow_codex_test" def validate_execution(
CUTOVER = Path("/tmp/blif_flow_v2_cutover_report.json") *, production_readonly_audit: bool, database: str, transaction_read_only: str,
INVENTORY = Path("/tmp/blif_flow_v2_legacy_field_inventory.json") mode: str,
COMPARE = Path("/tmp/blif_flow_v2_compare_report.json") ) -> None:
normalized_mode = str(mode or "").strip().lower()
if production_readonly_audit:
if database != PRODUCTION_DATABASE:
raise RuntimeError(
f"--production-readonly-audit requires database {PRODUCTION_DATABASE!r}, found {database!r}"
)
if transaction_read_only != "on":
raise RuntimeError("production audit requires transaction_read_only=on")
if normalized_mode not in {"shadow", "compare"}:
raise RuntimeError("production audit requires BLIF_FLOW_V2_MODE=shadow or compare")
return
if database == PRODUCTION_DATABASE:
raise RuntimeError("production database requires explicit --production-readonly-audit opt-in")
if database != TEST_DATABASE:
raise RuntimeError(f"default audit requires database {TEST_DATABASE!r}, found {database!r}")
if transaction_read_only != "on":
raise RuntimeError("cutover audit requires an explicit READ ONLY transaction")
FIELD_INVENTORY = [ def _code(value: Any) -> str:
{ return str(value or "").strip().upper()
"field": "stage", "kind": "derived_compatibility", "factual": False,
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
"app/reconciliation_service.py", "app/odoo_service.py", "app/external_reconciliation_sync.py"],
"readers": ["app/domain/opportunity_flow/evidence.py (V1 only)", "app/workflow_guard.py",
"app/admin_ui/pages/opportunities.py", "app/admin_ui/pages/orders.py",
"app/admin_dashboard.py", "app/revenue_forecast_service.py"],
"flow_v2_replacement": "opportunity_flow_state_v2.business_state",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
"notes": "Keep stored value initially for V1 boards, reports, forecasts and action guards; never use as factual V2 input. Consumers need code migration, not historical mass rewrite.",
},
{
"field": "lifecycle_state", "kind": "operational_compatibility", "factual": False,
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
"readers": ["app/operational_eligibility.py", "app/operations_service.py",
"app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "business state plus active operational obligations",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
"notes": "Awaiting/recovery/nurture presentation remains V1 compatibility. It cannot create a scheduled obligation without an active task.",
},
{
"field": "next_follow_up_at", "kind": "denormalized_followup_compatibility", "factual": False,
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
"readers": ["app/operations_service.py", "app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "pending follow-up task action_code/due_at",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
"notes": "Historical timestamp is inert without a pending follow-up task; retain all 44 values for now.",
},
{
"field": "nurture_until", "kind": "denormalized_schedule_compatibility", "factual": False,
"writers": ["app/opportunity_service.py"],
"readers": ["app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "explicit pending REVIEW_NURTURE/follow-up task due_at",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
},
{
"field": "follow_up_attempts", "kind": "historical_counter", "factual": False,
"writers": ["app/opportunity_service.py", "app/followup_service.py"],
"readers": ["app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "event/task history aggregation", "compatibility_required": True,
"one_time_migration_required": False, "eventually_deprecatable": True,
},
{
"field": "last_action_code", "kind": "last-known-action_compatibility", "factual": False,
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
"app/reconciliation_service.py", "app/odoo_service.py"],
"readers": ["app/workflow_guard.py", "app/admin_dashboard.py",
"app/admin_ui/pages/opportunities.py", "app/action_prompt.py"],
"flow_v2_replacement": "opportunity_flow_state_v2.business_next_action plus OperationalEligibility",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": True,
},
{
"field": "last_task_id", "kind": "historical_pointer", "factual": False,
"writers": ["app/opportunity_service.py", "app/reconciliation_service.py",
"app/company_opportunity_linking.py"],
"readers": ["app/workflow_guard.py"],
"flow_v2_replacement": "active canonical task lookup", "compatibility_required": True,
"one_time_migration_required": False, "eventually_deprecatable": True,
},
{
"field": "pending_primary_* / pending_follow_up_*", "kind": "runtime_read_model", "factual": False,
"writers": ["none (SQL projections in app/opportunity_service.py)"],
"readers": ["app/admin_ui/pages/opportunities.py"],
"flow_v2_replacement": "active task lookup remains an explicit operational overlay",
"compatibility_required": True, "one_time_migration_required": False,
"eventually_deprecatable": False,
"notes": "These virtual fields correctly derive active obligations and are not stored opportunity state.",
},
]
SOURCE_OF_TRUTH = [
["business process state", "opportunity_flow_state_v2.business_state"], def classify_semantic_difference(
["business next action", "opportunity_flow_state_v2.business_next_action"], *, v1_state: str | None, v1_action: str | None,
["time-sensitive operational queue", "OperationalEligibility"], v2_state: str | None, v2_action: str | None,
["scheduled follow-up", "pending task action_code + due_at"], v2_projection_present: bool = True,
["formal document state", "commercial_documents + active opportunity_document_links"], is_duplicate_representation: bool = False,
["payment", "confirmed factual payment evidence/operation link"], operational_action: str | None = None,
["Odoo execution", "operation_links / factual Odoo evidence"], operational_precedence: str | None = None,
["messages/customer response", "messages and communications chronology"], operational_queue: str | None = None,
["legacy stage", "compatibility only"], v2_confidence: str | None = None,
["tasks", "operator obligations/history; never business fact proof"], v2_diagnostic_status: str | None = None,
] ) -> tuple[str, str]:
"""Conservatively compare meanings rather than raw action vocabulary."""
if not v2_projection_present:
return "MISSING_PROJECTION", "No persisted Flow v2 projection exists for this opportunity."
v1, v2, effective = _code(v1_action), _code(v2_action), _code(operational_action)
precedence = str(operational_precedence or "").strip().lower()
queue = str(operational_queue or "").strip().lower()
if is_duplicate_representation:
return "V2_CORRECTS_V1", "Material identity suppresses a duplicate representation without deleting evidence."
if precedence in OVERRIDE_PRECEDENCE and effective and (v1 == effective or queue in CURRENT_QUEUES | {"waiting"}):
return "V1_OPERATIONAL_OVERRIDE", f"Explicit operational precedence {precedence} validly overlays the V2 business transition."
if v1 == v2 and v1:
return "SEMANTICALLY_EQUIVALENT", "V1 and V2 select the same action."
if not v1 and not v2:
return "SEMANTICALLY_EQUIVALENT", "Neither model has a current business action."
if v1 in REVIEW_ACTIONS and (v2 in REVIEW_ACTIONS or _code(v2_state) in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED"}):
return "SEMANTICALLY_EQUIVALENT", "Both decisions require review/blocker handling."
if v1 in {"NO_ACTION", ""} and not v2 and _code(v2_state) in {"COMPLETED", "LOST", "NO_INTEREST"}:
return "SEMANTICALLY_EQUIVALENT", "Both decisions represent a terminal/non-current process."
if v1.startswith("FOLLOW_UP_") and _code(v2_state) in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}:
return "V1_OPERATIONAL_OVERRIDE", "A current follow-up obligation overlays a waiting V2 business state."
if v1 == "CONFIRM_PAYMENT" and _code(v2_state) == "AWAITING_PAYMENT" and not v2:
return "SEMANTICALLY_EQUIVALENT", "Both decisions mean payment remains outstanding; V1 names the compatibility action."
if v1 in {"VALIDATE_FISCAL_CUSTOMER", "RECONCILE_DOCUMENTS"} and not effective and (
queue in {"not_current", "backlog"} or _code(v2_state) in {"INQUIRY", "AWAITING_CUSTOMER", "AWAITING_PAYMENT"}
):
return "LEGACY_ONLY", "Legacy data-hygiene/blocker vocabulary is not a current factual V2 obligation."
if precedence in {"safe_diagnostic_only", "safe_ambiguous_review"}:
return "V2_CORRECTS_V1", "SAFE V2 prevents ambiguous historical compatibility state from creating current work."
if _code(v2_state) in {"REVIEW_REQUIRED", "FISCAL_BLOCKED", "DOCUMENT_RECONCILIATION_REQUIRED"}:
return "V2_CORRECTS_V1", "V2 converts conflicting factual history into an explicit protected review state."
if effective and effective == v2 and v1 != v2:
return "V2_CORRECTS_V1", "The safe operational action follows the factual V2 transition rather than the legacy action."
if v1 in LEGACY_ACTIONS:
if v2 and _code(v2_state) not in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"}:
return "V2_CORRECTS_V1", "V2 replaces a legacy compatibility action with a factual business transition."
return "LEGACY_ONLY", "V1 action is compatibility vocabulary with no native V2 business transition."
if v2 and str(v2_confidence or "").lower() == "high" and str(v2_diagnostic_status or "").lower() == "clear":
return "V2_CORRECTS_V1", "High-confidence factual V2 transition corrects a different legacy action."
if effective and v1 == effective:
return "V1_OPERATIONAL_OVERRIDE", "V1 matches the safe operational overlay rather than the business action."
return "REAL_CONFLICT", "The available evidence does not establish equivalence, a valid override, or a safe V2 correction."
def _write(path: Path, value: Any) -> None: def _write(path: Path, value: Any) -> None:
path.write_text(json.dumps(value, ensure_ascii=False, indent=2, default=str), encoding="utf-8") path.write_text(json.dumps(value, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
def _identity_and_ids() -> tuple[dict[str, str], list[str], int]: def _queue_totals(rows: list[dict[str, Any]], key: str) -> dict[str, int]:
counts = Counter(str(row[key].get("operational_queue") or row[key].get("effective_operational_queue") or "not_current") for row in rows)
return {
"current": sum(counts[name] for name in CURRENT_QUEUES),
"do_now": counts["do_now"], "review": counts["review"],
"waiting": counts["waiting"], "backlog": counts["backlog"],
"not_current": counts["not_current"],
}
def _operation_totals(source: dict[str, Any]) -> dict[str, int]:
return {"current": int(source.get("current_work") or 0),
"do_now": int(source.get("do_now") or 0),
"review": int(source.get("review") or 0),
"waiting": int(source.get("waiting") or 0),
"backlog": int(source.get("backlog") or 0),
"not_current": int(source.get("not_current") or 0)}
def _database_snapshot(engine: Any) -> tuple[dict[str, str], list[str], dict[str, dict[str, Any]]]:
from sqlalchemy import text
with engine.connect() as conn: with engine.connect() as conn:
conn.exec_driver_sql("BEGIN READ ONLY") conn.exec_driver_sql("BEGIN READ ONLY")
try: try:
identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one() identity = conn.execute(text(
if identity[0] != EXPECTED_DATABASE or identity[2] != "on": "SELECT current_database(),current_user,current_setting('transaction_read_only')"
raise RuntimeError(f"refusing unexpected/non-read-only database: {identity!r}") )).one()
ids = [str(value) for value in conn.execute(text("SELECT opportunity_id FROM opportunity_flow_state_v2 ORDER BY opportunity_id")).scalars()] opportunity_ids = [str(value) for value in conn.execute(text(
historical = conn.execute(text(""" "SELECT id FROM opportunities ORDER BY id"
SELECT count(*) FROM opportunities o )).scalars()]
WHERE o.next_follow_up_at IS NOT NULL rows = conn.execute(text("""
AND NOT EXISTS ( SELECT opportunity_id::text, material_process_key,
SELECT 1 FROM tasks t WHERE t.opportunity_id=o.id AND t.status='pending' canonical_opportunity_id::text, is_duplicate_representation,
AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN business_state, business_next_action, diagnostic_status,
('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE')) confidence, reason_code, reason_text, evidence_refs, flow_version
) FROM opportunity_flow_state_v2 ORDER BY opportunity_id
""")).scalar_one() """)).mappings().all()
finally: finally:
conn.rollback() conn.rollback()
return {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]}, ids, int(historical) return (
{"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
opportunity_ids, {str(row["opportunity_id"]): dict(row) for row in rows},
)
def main() -> None: def run_audit(*, production_readonly_audit: bool) -> dict[str, Any]:
identity, ids, historical_timestamps = _identity_and_ids() # Imports occur only after argparse, so --help cannot initialize DB code.
if len(ids) != 328: from app.config import settings
raise RuntimeError(f"expected 328 persisted Flow v2 opportunities, found {len(ids)}") from app.db import engine
v1 = get_opportunity_next_actions(ids) from scripts.simulate_blif_flow_v2 import collect
v2 = _load_v2_comparison_rows(ids)
comparisons = compare_v1_v2_decisions(v1, v2) identity, opportunity_ids, persisted = _database_snapshot(engine)
counts = Counter("agree" if row["action_agrees"] else "different" for row in comparisons) mode = str(settings.blif_flow_v2_mode or "").strip().lower()
counts["missing_projection"] = sum(not row["projection_present"] for row in comparisons) validate_execution(
compare_report = { production_readonly_audit=production_readonly_audit,
"generated_at": datetime.now(timezone.utc), "database": identity, database=identity["database"], transaction_read_only=identity["transaction_read_only"],
"mode_semantics": "V1 returned; persisted V2 observed; no business mutation", mode=mode,
"opportunity_count": len(comparisons), "counts": dict(counts), "comparisons": comparisons, )
preflight = {**identity, "BLIF_FLOW_V2_MODE": mode,
"production_readonly_opt_in": production_readonly_audit}
print(json.dumps(preflight, sort_keys=True), flush=True)
expected_count = None if production_readonly_audit else 328
projection = collect(
expected_database=identity["database"], expected_user=identity["user"],
require_read_only=production_readonly_audit,
expected_opportunity_count=expected_count, require_opportunities=True,
)
records = {row["opportunity_id"]: row for row in projection["opportunities"]}
if set(records) != set(opportunity_ids):
raise RuntimeError("factual collector did not return the complete opportunity universe")
comparisons = []
for opportunity_id in opportunity_ids:
record = records[opportunity_id]
v1, operational = record["v1"], record["safe_v2"]
v2 = persisted.get(opportunity_id)
classification, reason = classify_semantic_difference(
v1_state=v1.get("commercial_stage"), v1_action=v1.get("current_action"),
v2_state=(v2 or {}).get("business_state"), v2_action=(v2 or {}).get("business_next_action"),
v2_projection_present=v2 is not None,
is_duplicate_representation=bool((v2 or {}).get("is_duplicate_representation")),
operational_action=operational.get("effective_operational_action"),
operational_precedence=operational.get("precedence"),
operational_queue=operational.get("effective_operational_queue"),
v2_confidence=(v2 or {}).get("confidence"),
v2_diagnostic_status=(v2 or {}).get("diagnostic_status"),
)
comparisons.append({
"opportunity_id": opportunity_id, "title": record.get("title"),
"customer_name": record.get("customer"), "classification": classification,
"classification_reason": reason,
"v1": {"state": v1.get("commercial_stage"), "action": v1.get("current_action"),
"queue": v1.get("operational_queue"), "reason_code": v1.get("reason")},
"v2": {"business_state": (v2 or {}).get("business_state"),
"business_next_action": (v2 or {}).get("business_next_action"),
"reason_code": (v2 or {}).get("reason_code"),
"diagnostic_status": (v2 or {}).get("diagnostic_status"),
"confidence": (v2 or {}).get("confidence")},
"operational_override": {"action": operational.get("effective_operational_action"),
"queue": operational.get("effective_operational_queue"),
"precedence": operational.get("precedence")},
"material_process_key": (v2 or {}).get("material_process_key"),
"canonical_opportunity_id": (v2 or {}).get("canonical_opportunity_id"),
"is_duplicate_representation": bool((v2 or {}).get("is_duplicate_representation")),
"evidence_refs": (v2 or {}).get("evidence_refs", []),
})
counts = Counter(row["classification"] for row in comparisons)
for name in CLASSIFICATIONS:
counts.setdefault(name, 0)
real_conflicts = [row for row in comparisons if row["classification"] == "REAL_CONFLICT"]
missing = [row for row in comparisons if row["classification"] == "MISSING_PROJECTION"]
semantic_report = {
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
"classification_counts": dict(sorted(counts.items())),
"real_conflicts": real_conflicts, "missing_projections": missing,
"comparisons": comparisons,
} }
_write(COMPARE, compare_report) _write(SEMANTIC, semantic_report)
inventory = {
"generated_at": datetime.now(timezone.utc), "database": identity,
"fields": FIELD_INVENTORY,
"summary": {
"fields_requiring_data_migration": [],
"compatibility_only_fields": [row["field"] for row in FIELD_INVENTORY if row["compatibility_required"]],
"safe_to_deprecate_after_consumer_migration": [row["field"] for row in FIELD_INVENTORY if row["eventually_deprecatable"]],
"historical_followup_timestamps_preserved": historical_timestamps,
},
}
_write(INVENTORY, inventory)
projection = collect(expected_database=EXPECTED_DATABASE, expected_user=identity["user"], require_read_only=False)
named = {} named = {}
for name in ("INSTALBEIRA", "PANORAMIC SUCCESS", "ENGEXICON", "CONSTRURECUP", "X MAT", "RZSOLAR"): by_id = {row["opportunity_id"]: row for row in comparisons}
matches = [row for row in projection["opportunities"] for name, opportunity_id in NAMED_IDS.items():
if name.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()] row = by_id[opportunity_id]
named[name] = [{"opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"], named[name] = row
"canonical_process_id": row["canonical_process_id"], "business_state": row["safe_v2"]["business_state"], # Use the complete canonical Operations candidate universe, including
"effective_action": row["safe_v2"]["effective_operational_action"], # preserved standalone obligations, rather than opportunity cards alone.
"queue": row["safe_v2"]["effective_operational_queue"], v1_metrics = _operation_totals(projection["v1_totals"])
"duplicate_suppressed": row["safe_v2"]["precedence"] == "duplicate_representation"} safe_metrics = _operation_totals(projection["safe_v2_totals"])
for row in matches] audit = {
cutover = { "generated_at": datetime.now(timezone.utc), "preflight": preflight,
"generated_at": datetime.now(timezone.utc), "database": identity, "read_only": True, "business_writes": 0, "projection_writes": 0,
"source_of_truth": [{"concern": concern, "source": source} for concern, source in SOURCE_OF_TRUTH], "opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
"data_migration": {"opportunities_requiring_mutation_before_shadow": 0, "canonical_count": sum(not row.get("is_duplicate_representation") for row in persisted.values()),
"opportunities_requiring_no_mutation_before_shadow": len(ids), "duplicate_representation_count": sum(bool(row.get("is_duplicate_representation")) for row in persisted.values()),
"broad_repair_required": False, "operations_metrics": {"v1": v1_metrics, "safe_v2": safe_metrics},
"stage_write_migration_required": False, "semantic_classification_counts": dict(sorted(counts.items())),
"lifecycle_write_migration_required": False, "authoritative_cutover_blockers": {
"followup_timestamp_cleanup_required": False}, "real_conflicts": len(real_conflicts), "missing_projections": len(missing),
"mode_contract": { "blocked": bool(real_conflicts or missing),
"off": "No Flow v2 derivation or projection writes are triggered by runtime reads.",
"shadow": "V1 returned; explicit projection rebuild is additive/idempotent; no tasks, stage, or UI behavior changed.",
"compare": "V1 returned; V2 projection read and structured comparison logged; no disagreement writes.",
"authoritative": "Disabled and fail-closed. No activation performed.",
}, },
"production_shadow_blockers": [ "real_conflicts": real_conflicts, "missing_projections": missing,
"Production migration 011 presence was not and must not be checked from this environment.", "named_cases": named,
"Deployment must provide an explicit projection rebuild cadence and comparison-log monitoring/retention.",
],
"mode_off_deployment_blockers": [],
"shadow_compare_activation_prerequisites": [
"Install migration 011 through the normal controlled production migration process.",
"Run mode=off first, then explicitly rebuild projection in shadow.",
"Monitor structured disagreement rates and missing projections before compare enablement.",
],
"authoritative_activation_blockers": [
"Authoritative switch intentionally raises and has no enabled code path.",
"Operations/UI must consume V2 business state/action followed by OperationalEligibility.",
"Legacy stage consumers in orders, forecasts, dashboard, workflow guards and opportunity columns require migration or explicit compatibility adapters.",
"Explicit operational override precedence must be implemented for CALL_CUSTOMER, due follow-up, SUPPORT, SEND_INFO, manual review and blockers.",
"Production shadow/compare observation, rollback criteria and zero-false-negative acceptance must be completed.",
],
"legacy_findings_are": "field-level compatibility findings, not required row mutations",
"historical_followup_timestamps_preserved": historical_timestamps,
"compare_summary": compare_report["counts"], "named_cases": named,
} }
_write(CUTOVER, cutover) _write(AUDIT, audit)
print(json.dumps({"database": identity, "opportunities": len(ids), lines = [
"compare": compare_report["counts"], "historical_timestamps": historical_timestamps, "BLIF FLOW V2 PRODUCTION SHADOW READ-ONLY AUDIT", "",
"mutation_required": 0, "outputs": [str(CUTOVER), str(INVENTORY), str(COMPARE)]}, indent=2)) f"database: {identity['database']}", f"user: {identity['user']}",
f"transaction_read_only: {identity['transaction_read_only']}",
f"BLIF_FLOW_V2_MODE: {mode}",
f"production_readonly_opt_in: {production_readonly_audit}", "",
f"opportunities: {len(opportunity_ids)}", f"projections: {len(persisted)}",
f"canonical: {audit['canonical_count']}",
f"duplicate representations: {audit['duplicate_representation_count']}", "",
f"V1 metrics: {json.dumps(v1_metrics, sort_keys=True)}",
f"SAFE V2 metrics: {json.dumps(safe_metrics, sort_keys=True)}", "",
f"semantic classifications: {json.dumps(dict(sorted(counts.items())), sort_keys=True)}",
f"REAL_CONFLICT blockers: {len(real_conflicts)}",
f"MISSING_PROJECTION blockers: {len(missing)}",
"business writes: 0", "projection writes: 0",
]
SUMMARY.write_text("\n".join(lines) + "\n", encoding="utf-8")
return audit
def main(argv: Sequence[str] | None = None) -> int:
# --help exits here before application/database imports or connections.
args = build_parser().parse_args(argv)
result = run_audit(production_readonly_audit=args.production_readonly_audit)
print(json.dumps({
"database": result["preflight"]["database"],
"opportunities": result["opportunity_count"], "projections": result["projection_count"],
"semantic_classifications": result["semantic_classification_counts"],
"authoritative_cutover_blockers": result["authoritative_cutover_blockers"],
"outputs": [str(AUDIT), str(SEMANTIC), str(SUMMARY)],
}, indent=2, sort_keys=True))
return 0
if __name__ == "__main__": if __name__ == "__main__":
main() raise SystemExit(main())

View File

@@ -0,0 +1,110 @@
from inspect import getsource
import pytest
import scripts.audit_blif_flow_v2_cutover as audit
def classify(**overrides):
values = dict(
v1_state="INFO_SENT", v1_action="SEND_INFO",
v2_state="INQUIRY", v2_action="SEND_INFO",
v2_projection_present=True, is_duplicate_representation=False,
operational_action="SEND_INFO", operational_precedence="business_transition",
operational_queue="do_now", v2_confidence="high", v2_diagnostic_status="clear",
)
values.update(overrides)
return audit.classify_semantic_difference(**values)[0]
def test_help_performs_zero_database_work(monkeypatch, capsys):
monkeypatch.setattr(audit, "run_audit", lambda **kwargs: pytest.fail("help reached DB audit"))
with pytest.raises(SystemExit) as exc:
audit.main(["--help"])
assert exc.value.code == 0
assert "--production-readonly-audit" in capsys.readouterr().out
def test_production_requires_explicit_opt_in():
with pytest.raises(RuntimeError, match="explicit --production-readonly-audit"):
audit.validate_execution(production_readonly_audit=False, database="clientflow",
transaction_read_only="on", mode="shadow")
def test_production_database_must_be_clientflow():
with pytest.raises(RuntimeError, match="requires database 'clientflow'"):
audit.validate_execution(production_readonly_audit=True, database="clientflow_codex_test",
transaction_read_only="on", mode="shadow")
def test_production_transaction_must_be_read_only():
with pytest.raises(RuntimeError, match="transaction_read_only=on"):
audit.validate_execution(production_readonly_audit=True, database="clientflow",
transaction_read_only="off", mode="shadow")
def test_production_audit_contains_no_sql_writes():
source = (getsource(audit._database_snapshot) + getsource(audit.run_audit)).upper()
for verb in ("INSERT ", "UPDATE ", "DELETE ", "CREATE ", "ALTER ", "DROP ", "TRUNCATE "):
assert verb not in source
assert 'BEGIN READ ONLY' in source
@pytest.mark.parametrize("mode", ["shadow", "compare"])
def test_shadow_and_compare_are_allowed(mode):
audit.validate_execution(production_readonly_audit=True, database="clientflow",
transaction_read_only="on", mode=mode)
@pytest.mark.parametrize("mode", ["", "off", "authoritative"])
def test_unsafe_production_modes_are_refused(mode):
with pytest.raises(RuntimeError, match="requires BLIF_FLOW_V2_MODE=shadow or compare"):
audit.validate_execution(production_readonly_audit=True, database="clientflow",
transaction_read_only="on", mode=mode)
def test_missing_projection_is_reported():
assert classify(v2_projection_present=False, v2_state=None, v2_action=None) == "MISSING_PROJECTION"
def test_semantically_equivalent_is_classified():
assert classify() == "SEMANTICALLY_EQUIVALENT"
assert classify(v1_action="NO_ACTION", v2_state="COMPLETED", v2_action=None,
operational_action=None, operational_queue="not_current") == "SEMANTICALLY_EQUIVALENT"
def test_operational_override_is_classified():
assert classify(v1_action="FOLLOW_UP_CUSTOMER_REVIEW", v2_state="AWAITING_CUSTOMER",
v2_action=None, operational_action="FOLLOW_UP_CUSTOMER_REVIEW",
operational_precedence="due_followup") == "V1_OPERATIONAL_OVERRIDE"
def test_v2_correction_is_classified():
assert classify(v1_action="SEND_INVOICE", v2_state="PROFORMA_REQUIRED",
v2_action="CREATE_PROFORMA") == "V2_CORRECTS_V1"
assert classify(v1_action="REVIEW_RECONSTRUCTED_PROCESS", v2_state="REVIEW_REQUIRED",
v2_action="REVIEW_REQUIRED", is_duplicate_representation=True) == "V2_CORRECTS_V1"
def test_legacy_only_is_classified():
assert classify(v1_action="CREATE_JASMIN_QUOTE", v2_state="AWAITING_CUSTOMER",
v2_action=None, v2_confidence="medium") == "LEGACY_ONLY"
def test_real_conflict_is_classified():
assert classify(v1_action="UNMAPPED_ACTION", v2_state="INQUIRY",
v2_action=None, v2_confidence="medium",
v2_diagnostic_status="ambiguous", operational_action=None) == "REAL_CONFLICT"
def test_named_case_ids_and_expected_semantics_are_frozen():
assert audit.NAMED_IDS == {
"INSTALBEIRA": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
"PANORAMIC": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
"ENGEXICON": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
"CONSTRURECUP": "e3b23ac5-84db-4763-8a31-a684e873032c",
"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",
}