fix: allow guarded production Flow v2 read-only audit

This commit is contained in:
plx
2026-08-16 00:53:29 +00:00
parent e3b8750ebb
commit e456cbc0e4
2 changed files with 400 additions and 190 deletions

View File

@@ -1,228 +1,328 @@
#!/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
import argparse
import json
import os
import sys
from collections import Counter
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from typing import Any, Sequence
ROOT = Path(__file__).resolve().parents[1]
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 (
_load_v2_comparison_rows, compare_v1_v2_decisions, get_opportunity_next_actions,
def build_parser() -> argparse.ArgumentParser:
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"
CUTOVER = Path("/tmp/blif_flow_v2_cutover_report.json")
INVENTORY = Path("/tmp/blif_flow_v2_legacy_field_inventory.json")
COMPARE = Path("/tmp/blif_flow_v2_compare_report.json")
def validate_execution(
*, production_readonly_audit: bool, database: str, transaction_read_only: str,
mode: str,
) -> 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 = [
{
"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.",
},
]
def _code(value: Any) -> str:
return str(value or "").strip().upper()
SOURCE_OF_TRUTH = [
["business process state", "opportunity_flow_state_v2.business_state"],
["business next action", "opportunity_flow_state_v2.business_next_action"],
["time-sensitive operational queue", "OperationalEligibility"],
["scheduled follow-up", "pending task action_code + due_at"],
["formal document state", "commercial_documents + active opportunity_document_links"],
["payment", "confirmed factual payment evidence/operation link"],
["Odoo execution", "operation_links / factual Odoo evidence"],
["messages/customer response", "messages and communications chronology"],
["legacy stage", "compatibility only"],
["tasks", "operator obligations/history; never business fact proof"],
]
def classify_semantic_difference(
*, v1_state: str | None, v1_action: str | None,
v2_state: str | None, v2_action: str | None,
v2_projection_present: bool = True,
is_duplicate_representation: bool = False,
operational_action: str | None = None,
operational_precedence: str | None = None,
operational_queue: str | None = None,
v2_confidence: str | None = None,
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:
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:
conn.exec_driver_sql("BEGIN READ ONLY")
try:
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!r}")
ids = [str(value) for value in conn.execute(text("SELECT opportunity_id FROM opportunity_flow_state_v2 ORDER BY opportunity_id")).scalars()]
historical = conn.execute(text("""
SELECT count(*) FROM opportunities o
WHERE o.next_follow_up_at IS NOT NULL
AND NOT EXISTS (
SELECT 1 FROM tasks t WHERE t.opportunity_id=o.id AND t.status='pending'
AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN
('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
)
""")).scalar_one()
identity = conn.execute(text(
"SELECT current_database(),current_user,current_setting('transaction_read_only')"
)).one()
opportunity_ids = [str(value) for value in conn.execute(text(
"SELECT id FROM opportunities ORDER BY id"
)).scalars()]
rows = 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, flow_version
FROM opportunity_flow_state_v2 ORDER BY opportunity_id
""")).mappings().all()
finally:
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:
identity, ids, historical_timestamps = _identity_and_ids()
if len(ids) != 328:
raise RuntimeError(f"expected 328 persisted Flow v2 opportunities, found {len(ids)}")
v1 = get_opportunity_next_actions(ids)
v2 = _load_v2_comparison_rows(ids)
comparisons = compare_v1_v2_decisions(v1, v2)
counts = Counter("agree" if row["action_agrees"] else "different" for row in comparisons)
counts["missing_projection"] = sum(not row["projection_present"] for row in comparisons)
compare_report = {
"generated_at": datetime.now(timezone.utc), "database": identity,
"mode_semantics": "V1 returned; persisted V2 observed; no business mutation",
"opportunity_count": len(comparisons), "counts": dict(counts), "comparisons": comparisons,
def run_audit(*, production_readonly_audit: bool) -> dict[str, Any]:
# Imports occur only after argparse, so --help cannot initialize DB code.
from app.config import settings
from app.db import engine
from scripts.simulate_blif_flow_v2 import collect
identity, opportunity_ids, persisted = _database_snapshot(engine)
mode = str(settings.blif_flow_v2_mode or "").strip().lower()
validate_execution(
production_readonly_audit=production_readonly_audit,
database=identity["database"], transaction_read_only=identity["transaction_read_only"],
mode=mode,
)
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 = {}
for name in ("INSTALBEIRA", "PANORAMIC SUCCESS", "ENGEXICON", "CONSTRURECUP", "X MAT", "RZSOLAR"):
matches = [row for row in projection["opportunities"]
if name.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()]
named[name] = [{"opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"],
"canonical_process_id": row["canonical_process_id"], "business_state": row["safe_v2"]["business_state"],
"effective_action": row["safe_v2"]["effective_operational_action"],
"queue": row["safe_v2"]["effective_operational_queue"],
"duplicate_suppressed": row["safe_v2"]["precedence"] == "duplicate_representation"}
for row in matches]
cutover = {
"generated_at": datetime.now(timezone.utc), "database": identity,
"source_of_truth": [{"concern": concern, "source": source} for concern, source in SOURCE_OF_TRUTH],
"data_migration": {"opportunities_requiring_mutation_before_shadow": 0,
"opportunities_requiring_no_mutation_before_shadow": len(ids),
"broad_repair_required": False,
"stage_write_migration_required": False,
"lifecycle_write_migration_required": False,
"followup_timestamp_cleanup_required": False},
"mode_contract": {
"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.",
by_id = {row["opportunity_id"]: row for row in comparisons}
for name, opportunity_id in NAMED_IDS.items():
row = by_id[opportunity_id]
named[name] = row
# Use the complete canonical Operations candidate universe, including
# preserved standalone obligations, rather than opportunity cards alone.
v1_metrics = _operation_totals(projection["v1_totals"])
safe_metrics = _operation_totals(projection["safe_v2_totals"])
audit = {
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
"read_only": True, "business_writes": 0, "projection_writes": 0,
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
"canonical_count": sum(not row.get("is_duplicate_representation") for row in persisted.values()),
"duplicate_representation_count": sum(bool(row.get("is_duplicate_representation")) for row in persisted.values()),
"operations_metrics": {"v1": v1_metrics, "safe_v2": safe_metrics},
"semantic_classification_counts": dict(sorted(counts.items())),
"authoritative_cutover_blockers": {
"real_conflicts": len(real_conflicts), "missing_projections": len(missing),
"blocked": bool(real_conflicts or missing),
},
"production_shadow_blockers": [
"Production migration 011 presence was not and must not be checked from this environment.",
"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,
"real_conflicts": real_conflicts, "missing_projections": missing,
"named_cases": named,
}
_write(CUTOVER, cutover)
print(json.dumps({"database": identity, "opportunities": len(ids),
"compare": compare_report["counts"], "historical_timestamps": historical_timestamps,
"mutation_required": 0, "outputs": [str(CUTOVER), str(INVENTORY), str(COMPARE)]}, indent=2))
_write(AUDIT, audit)
lines = [
"BLIF FLOW V2 PRODUCTION SHADOW READ-ONLY AUDIT", "",
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__":
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",
}