diff --git a/app/operations_service.py b/app/operations_service.py index 6a109a6..d496750 100644 --- a/app/operations_service.py +++ b/app/operations_service.py @@ -8,6 +8,7 @@ from __future__ import annotations import os import re import subprocess +from contextlib import nullcontext from typing import Any, Dict, List, Optional from sqlalchemy import text @@ -26,7 +27,7 @@ from app.canonical_operations import ( opportunity_ids_from_work_seeds, partition_canonical_items, ) -from app.opportunity_next_action_service import get_opportunity_next_actions +from app.opportunity_next_action_service import get_opportunity_next_actions, get_opportunity_next_actions_for_mode from app.document_reconciliation_service import active_document_link_exclusion_sql @@ -104,7 +105,7 @@ def _norm_identity(value: Any) -> str: def _identity_tokens(value: Any) -> set[str]: - return { + result = { token for token in _norm_identity(value).split() if len(token) >= 3 and token not in _IDENTITY_STOPWORDS @@ -302,7 +303,9 @@ def _attach_operation_urls(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]: return _sanitize_operation_identities(cleaned) -def get_operations_summary(limit: int = 24) -> Dict[str, Any]: +def get_operations_summary( + limit: int = 24, *, flow_mode: str | None = None, connection: Any = None, +) -> Dict[str, Any]: """Build the /operations work queue summary. /operations is intentionally not a mini-dashboard. It returns a compact @@ -315,7 +318,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: # after canonicalization. This prevents duplicate rows from consuming the # display limit while still avoiding an unbounded history scan. candidate_limit = max(200, min(display_limit * 20, 1000)) - with engine.begin() as conn: + active_flow_mode = str(flow_mode if flow_mode is not None else settings.blif_flow_v2_mode or "off").strip().lower() + with (nullcontext(connection) if connection is not None else engine.begin()) as conn: counts = conn.execute(text(""" SELECT (SELECT COUNT(*) FROM opportunities WHERE status = 'open')::int AS open_opportunities, @@ -675,9 +679,9 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: # In authoritative mode the factual projection, not a legacy task/stage, # selects material processes that have current business work. Synthetic # rows are read-model seeds only and are never persisted. - if str(settings.blif_flow_v2_mode or "off").strip().lower() == "authoritative": + if active_flow_mode == "authoritative": seeded = {str(row.get("opportunity_id") or "") for row in seed_rows} - with engine.connect() as conn: + with (nullcontext(connection) if connection is not None else engine.connect()) as conn: v2_seeds = conn.execute(text(""" SELECT p.opportunity_id::text, p.business_next_action, p.derived_at, o.title, o.value_amount, o.currency, o.lifecycle_state, @@ -713,7 +717,9 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")}) normalized_rows = _normalise_work_item_intent(seed_rows) opportunity_ids = opportunity_ids_from_work_seeds(normalized_rows) - decisions = get_opportunity_next_actions(opportunity_ids) if opportunity_ids else {} + decisions = (get_opportunity_next_actions_for_mode( + opportunity_ids, flow_mode=active_flow_mode, connection=connection, + ) if opportunity_ids else {}) projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows]) canonical_items = _attach_operation_urls(list(projection["items"])) partition = partition_canonical_items(canonical_items, display_limit=display_limit) @@ -758,6 +764,11 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: **diagnostics, "projection_metrics": diagnostics, } + if flow_mode is not None: + # Complete, untruncated read-model population for offline simulation; + # normal API/GET callers do not receive this audit-only field. + result["simulation_all_items"] = canonical_items + return result def get_system_health_summary() -> Dict[str, Any]: diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index c041c55..3a6f39b 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -376,13 +376,24 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di if not ids: return {} - if _flow_v2_mode() == "authoritative": - return _get_authoritative_opportunity_decisions(ids) + return get_opportunity_next_actions_for_mode(ids, flow_mode=_flow_v2_mode()) + + +def get_opportunity_next_actions_for_mode( + opportunity_ids: Iterable[str], *, flow_mode: str, connection: Any = None, +) -> Dict[str, Dict[str, Any]]: + """Explicit-mode read boundary used by the read-only cutover simulation.""" + ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip())) + if not ids: + return {} + if flow_mode == "authoritative": + return _get_authoritative_opportunity_decisions(ids, connection=connection) # Read path only: schema creation belongs to startup/migrations. In # particular, GET /operations must never perform DDL while calculating its # canonical projection. - with engine.begin() as conn: + from contextlib import nullcontext + with (nullcontext(connection) if connection is not None else engine.begin()) as conn: opportunities = _bulk_rows(conn, """ SELECT id::text, stage, status, title, local_customer_id::text AS fiscal_customer_id, @@ -453,13 +464,15 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di company_profile="blif", ) decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict() - _observe_flow_v2(decisions) + if flow_mode == "compare": + _observe_flow_v2(decisions) return decisions -def _get_authoritative_opportunity_decisions(ids: list[str]) -> Dict[str, Dict[str, Any]]: +def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any = None) -> Dict[str, Dict[str, Any]]: """Read-only set-oriented adapter input loader; never derives or persists V2.""" - with engine.connect() as conn: + from contextlib import nullcontext + with (nullcontext(connection) if connection is not None else engine.connect()) as conn: projections = _bulk_rows(conn, """ SELECT opportunity_id::text, material_process_key, canonical_opportunity_id::text, is_duplicate_representation, diff --git a/scripts/simulate_blif_flow_v2_authoritative.py b/scripts/simulate_blif_flow_v2_authoritative.py index 79102b5..8562479 100644 --- a/scripts/simulate_blif_flow_v2_authoritative.py +++ b/scripts/simulate_blif_flow_v2_authoritative.py @@ -1,95 +1,226 @@ #!/usr/bin/env python3 -"""Read-only audit of the exact Operations boundary in V1 and authoritative V2.""" +"""Simulate authoritative BLIF Flow v2 without changing application mode.""" from __future__ import annotations +import argparse import json +import os from pathlib import Path import sys - -ROOT = Path(__file__).resolve().parents[1] -if str(ROOT) not in sys.path: - sys.path.insert(0, str(ROOT)) - -from app.config import settings -from app.operations_service import get_operations_summary +from typing import Any JSON_PATH = Path("/tmp/blif_flow_v2_authoritative_simulation.json") TEXT_PATH = Path("/tmp/blif_flow_v2_authoritative_operations.txt") CURRENT = {"do_now", "review", "blocked", "exception"} +CLASSIFICATIONS = { + "SATISFIED", "SUPERSEDED", "DUPLICATE_SUPPRESSED", "PREMATURE_REMOVED", + "REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE", +} +NAMED = { + "Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "do_now"), + "Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"), + "ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "do_now"), + "CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "do_now"), + "X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, "not_current"), + "X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, "suppressed"), + "RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "review"), + "RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", None, "suppressed"), +} -def _all(summary): - return list(summary.get("work_items") or []) + list(summary.get("waiting_items") or []) + list(summary.get("backlog_items") or []) + list(summary.get("not_current_items") or []) +def parse_args(argv: list[str] | None = None) -> argparse.Namespace: + parser = argparse.ArgumentParser( + description="Read-only V1 versus simulated authoritative BLIF Flow v2 Operations audit." + ) + parser.add_argument( + "--production-readonly-simulation", action="store_true", + help="Required opt-in when current_database() is clientflow; requires compare mode and a read-only transaction.", + ) + return parser.parse_args(argv) -def _metrics(summary): - items = _all(summary) - queues = [str(item.get("operational_queue") or "") for item in items] +def _load_runtime(): + """Application imports are intentionally delayed until after argparse.""" + root = Path(__file__).resolve().parents[1] + if str(root) not in sys.path: + sys.path.insert(0, str(root)) + from app.db import engine + from app.operations_service import get_operations_summary + from app.opportunity_next_action_service import get_opportunity_next_actions_for_mode + return engine, get_operations_summary, get_opportunity_next_actions_for_mode + + +def _all(summary: dict[str, Any]) -> list[dict[str, Any]]: + if "simulation_all_items" in summary: + return list(summary.get("simulation_all_items") or []) + return sum((list(summary.get(key) or []) for key in + ("work_items", "waiting_items", "backlog_items", "not_current_items")), []) + + +def _metrics(summary: dict[str, Any]) -> dict[str, int]: + queues = [str(item.get("operational_queue") or "") for item in _all(summary)] + return {"current": sum(q in CURRENT for q in queues), "do_now": queues.count("do_now"), + "review": queues.count("review"), "waiting": queues.count("waiting"), + "backlog": queues.count("backlog"), "blocked": queues.count("blocked"), + "not_current": queues.count("not_current")} + + +def _key(item: dict[str, Any]) -> str: + return str(item.get("work_item_key") or item.get("process_key") or item.get("opportunity_id") or item.get("id")) + + +def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[str, str]: + if new is not None: + return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item." + decision = old.get("decision") if isinstance(old.get("decision"), dict) else {} + code = str(decision.get("reason_code") or old.get("eligibility_reason_code") or "").upper() + if "DUPLICATE" in code: + return "DUPLICATE_SUPPRESSED", "The local opportunity is a duplicate material representation." + if "SATISF" in code or "ANSWERED" in code: + return "SATISFIED", "Later factual evidence satisfies the old obligation." + if "SUPERSE" in code or "TERMINAL" in code: + return "SUPERSEDED", "A later factual state supersedes the legacy card." + if "PREMATURE" in code: + return "PREMATURE_REMOVED", "The legacy action is premature for the factual state." + if old.get("source") in {"task", "opportunity"}: + return "LEGACY_ONLY", "The card is supported only by legacy operational representation." + return "UNSAFE_FALSE_NEGATIVE", "No safe factual explanation was found for removing this current card." + + +def build_report( + *, identity: dict[str, Any], configured_mode: str, v1: dict[str, Any], v2: dict[str, Any], + named_decisions: dict[str, dict[str, Any]], missing_projections: int, +) -> dict[str, Any]: + v1_current = {_key(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT} + v2_current = {_key(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT} + removed, changed = [], [] + for key, old in v1_current.items(): + new = v2_current.get(key) + if new is not None and new.get("action_code") == old.get("action_code"): + continue + classification, reason = _classification(old, new) + row = {"work_item_key": key, "opportunity_id": old.get("opportunity_id"), + "v1_action": old.get("action_code"), "simulated_v2_action": (new or {}).get("action_code"), + "classification": classification, "reason": reason, + "source_refs": old.get("source_refs") or []} + (changed if new is not None else removed).append(row) + added = [{"work_item_key": key, "opportunity_id": item.get("opportunity_id"), + "action": item.get("action_code"), "source_refs": item.get("source_refs") or []} + for key, item in v2_current.items() if key not in v1_current] + + material_counts: dict[str, int] = {} + for item in v2_current.values(): + material = str(item.get("canonical_opportunity_id") or item.get("opportunity_id") or item.get("process_key")) + material_counts[material] = material_counts.get(material, 0) + 1 + duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1] + + named_cases = {} + for name, (oid, expected_action, expected_queue) in NAMED.items(): + decision = named_decisions.get(oid) or {} + suppressed = decision.get("suppress_current_card") is True + actual_queue = "suppressed" if suppressed else decision.get("operational_queue") + actual_action = decision.get("effective_action") + passed = actual_action == expected_action and actual_queue == expected_queue + named_cases[name] = {"opportunity_id": oid, "expected_action": expected_action, + "expected_queue": expected_queue, "actual_action": actual_action, + "actual_queue": actual_queue, "passed": passed} + + fiscal = [] + for item in v2_current.values(): + decision = item.get("decision") if isinstance(item.get("decision"), dict) else {} + if item.get("action_code") == "VALIDATE_FISCAL_CUSTOMER": + fiscal.append({"opportunity_id": item.get("opportunity_id"), + "business_action": decision.get("business_next_action"), + "reason_code": decision.get("reason_code"), + "reason": decision.get("description")}) + unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed) return { - "current": sum(queue in CURRENT for queue in queues), - "do_now": queues.count("do_now"), "review": queues.count("review"), - "waiting": queues.count("waiting"), "backlog": queues.count("backlog"), - "not_current": queues.count("not_current"), + "identity": identity, "configured_mode": configured_mode, + "v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2), + "semantic_changes": len(removed) + len(changed) + len(added), + "removed_cards": removed, "added_cards": added, "changed_actions": changed, + "duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates), + "missing_projections": missing_projections, "unsafe_false_negatives": unsafe, + "named_cases": named_cases, "fiscal_validation_cases": fiscal, } -def _identity(item): - return str(item.get("canonical_opportunity_id") or item.get("opportunity_id") or item.get("process_key") or item.get("work_item_key")) - - -def _classification(old, new): - if new is not None: - return "REPLACED_BY_CORRECT_ACTION" - decision = old.get("decision") or {} - code = str(decision.get("reason_code") or old.get("eligibility_reason_code") or "").upper() - if "DUPLICATE" in code: - return "DUPLICATE_SUPPRESSED" - if "SATISF" in code or "ANSWERED" in code: - return "SATISFIED" - if "SUPERSE" in code or "TERMINAL" in code: - return "SUPERSEDED" - if "PREMATURE" in code: - return "PREMATURE_REMOVED" - if old.get("source") in {"task", "opportunity"}: - return "LEGACY_ONLY" - return "UNSAFE_FALSE_NEGATIVE" - - -def simulate(): - original = settings.blif_flow_v2_mode - try: - settings.blif_flow_v2_mode = "off" - v1 = get_operations_summary(limit=200) - settings.blif_flow_v2_mode = "authoritative" - v2 = get_operations_summary(limit=200) - finally: - settings.blif_flow_v2_mode = original - v1_current = {_identity(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT} - v2_current = {_identity(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT} - removed = [] - for key, old in v1_current.items(): - new = v2_current.get(key) - if new is None or new.get("action_code") != old.get("action_code"): - removed.append({"identity": key, "v1_action": old.get("action_code"), - "v2_action": (new or {}).get("action_code"), - "classification": _classification(old, new)}) - added = [{"identity": key, "action": item.get("action_code")} - for key, item in v2_current.items() - if key not in v1_current or v1_current[key].get("action_code") != item.get("action_code")] - result = {"v1": _metrics(v1), "authoritative_v2": _metrics(v2), - "cards_removed": removed, "cards_added_or_replaced": added, - "unsafe_false_negatives": sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed)} - JSON_PATH.write_text(json.dumps(result, indent=2, ensure_ascii=False, default=str) + "\n") - lines = ["BLIF Flow v2 authoritative Operations simulation", "", f"V1: {result['v1']}", - f"Authoritative V2: {result['authoritative_v2']}", - f"Cards removed: {len(removed)}", f"Cards added/replaced: {len(added)}", - f"UNSAFE_FALSE_NEGATIVE: {result['unsafe_false_negatives']}"] +def _write_reports(report: dict[str, Any]) -> None: + JSON_PATH.write_text(json.dumps(report, indent=2, ensure_ascii=False, default=str) + "\n") + lines = ["BLIF Flow v2 authoritative Operations simulation", "", + f"Identity: {report['identity']}", f"Configured mode: {report['configured_mode']}", + f"V1: {report['v1_metrics']}", f"Simulated authoritative V2: {report['authoritative_metrics']}", + f"Removed: {len(report['removed_cards'])}", f"Added: {len(report['added_cards'])}", + f"Changed actions: {len(report['changed_actions'])}", + f"Duplicate current cards: {report['duplicate_current_cards']}", + f"Missing projections: {report['missing_projections']}", + f"UNSAFE_FALSE_NEGATIVE: {report['unsafe_false_negatives']}", "", "Named cases:"] + lines.extend(f"- {name}: {case}" for name, case in report["named_cases"].items()) + lines.extend(["", f"Fiscal validation cases: {len(report['fiscal_validation_cases'])}"]) TEXT_PATH.write_text("\n".join(lines) + "\n") - return result + + +def exit_code(report: dict[str, Any]) -> int: + failed_named = any(not case.get("passed") for case in report.get("named_cases", {}).values()) + unsafe = int(report.get("unsafe_false_negatives") or 0) + missing = int(report.get("missing_projections") or 0) + duplicates = int(report.get("duplicate_current_cards") or 0) + return 2 if unsafe or missing or duplicates or failed_named else 0 + + +def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, Any]: + engine, get_operations, get_decisions = runtime_loader() + configured_mode = os.environ.get("BLIF_FLOW_V2_MODE") + conn = engine.connect() + try: + # One explicit transaction and one factual snapshot for identity, V1, + # V2, projection completeness, and named-case decisions. + conn.exec_driver_sql("BEGIN READ ONLY") + identity = dict(conn.exec_driver_sql(""" + SELECT current_database() AS database, current_user AS user, + current_setting('transaction_read_only') AS transaction_read_only + """).mappings().one()) + production = identity.get("database") == "clientflow" + if production and not args.production_readonly_simulation: + raise RuntimeError("production simulation requires --production-readonly-simulation") + if args.production_readonly_simulation: + if identity.get("database") != "clientflow": + raise RuntimeError("production simulation requires current_database() = clientflow") + if identity.get("transaction_read_only") != "on": + raise RuntimeError("production simulation requires transaction_read_only = on") + if configured_mode is None: + raise RuntimeError("production simulation requires explicit BLIF_FLOW_V2_MODE=compare") + if configured_mode.strip().lower() != "compare": + raise RuntimeError("production simulation requires BLIF_FLOW_V2_MODE=compare") + + v1 = get_operations(limit=200, flow_mode="off", connection=conn) + v2 = get_operations(limit=200, flow_mode="authoritative", connection=conn) + projection_counts = conn.exec_driver_sql(""" + SELECT (SELECT count(*) FROM opportunities) AS opportunities, + (SELECT count(*) FROM opportunity_flow_state_v2) AS projections + """).mappings().one() + named_ids = [value[0] for value in NAMED.values()] + named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn) + if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode: + raise RuntimeError("configured BLIF_FLOW_V2_MODE changed during simulation") + report = build_report( + identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2, + named_decisions=named_decisions, + missing_projections=max(0, int(projection_counts["opportunities"]) - int(projection_counts["projections"])), + ) + _write_reports(report) + return report + finally: + conn.exec_driver_sql("ROLLBACK") + conn.close() + + +def main(argv: list[str] | None = None, *, runtime_loader=_load_runtime) -> int: + args = parse_args(argv) # --help exits before runtime_loader/application imports. + report = run(args, runtime_loader=runtime_loader) + print(json.dumps(report, indent=2, ensure_ascii=False, default=str)) + return exit_code(report) if __name__ == "__main__": - audit = simulate() - print(json.dumps(audit, indent=2, ensure_ascii=False, default=str)) - if audit["unsafe_false_negatives"]: - raise SystemExit(2) + raise SystemExit(main()) diff --git a/tests/test_blif_flow_v2_authoritative_simulation.py b/tests/test_blif_flow_v2_authoritative_simulation.py new file mode 100644 index 0000000..f4c841e --- /dev/null +++ b/tests/test_blif_flow_v2_authoritative_simulation.py @@ -0,0 +1,138 @@ +from argparse import Namespace +import importlib.util +from pathlib import Path + +import pytest + + +PATH = Path("scripts/simulate_blif_flow_v2_authoritative.py") +SPEC = importlib.util.spec_from_file_location("authoritative_simulation", PATH) +sim = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(sim) + + +class Result: + def __init__(self, row): self.row = row + def mappings(self): return self + def one(self): return self.row + + +class Connection: + def __init__(self, database="clientflow", readonly="on"): + self.database, self.readonly, self.sql = database, readonly, [] + def exec_driver_sql(self, sql): + self.sql.append(sql.strip()) + if "current_database" in sql: + return Result({"database": self.database, "user": "audit", "transaction_read_only": self.readonly}) + if "count(*) FROM opportunities" in sql: + return Result({"opportunities": 328, "projections": 328}) + return Result({}) + def close(self): pass + + +class Engine: + def __init__(self, connection): self.connection = connection + def connect(self): return self.connection + + +def named_decisions(fail=False): + result = {} + for name, (oid, action, queue) in sim.NAMED.items(): + result[oid] = {"effective_action": action, + "operational_queue": "not_current" if queue == "suppressed" else queue, + "suppress_current_card": queue == "suppressed"} + if fail: + result[sim.NAMED["Instalbeira"][0]]["effective_action"] = "SEND_INVOICE" + return result + + +def runtime(connection, *, fail_named=False, modes=None): + empty = {"work_items": [], "waiting_items": [], "backlog_items": [], "not_current_items": []} + def operations(**kwargs): + assert kwargs["connection"] is connection + if modes is not None: + modes.append(kwargs["flow_mode"]) + return empty + def decisions(ids, **kwargs): + assert kwargs == {"flow_mode": "authoritative", "connection": connection} + return named_decisions(fail_named) + return lambda: (Engine(connection), operations, decisions) + + +def test_help_performs_zero_db_work(): + called = False + def loader(): + nonlocal called + called = True + raise AssertionError("DB/application runtime must not load") + with pytest.raises(SystemExit) as exc: + sim.main(["--help"], runtime_loader=loader) + assert exc.value.code == 0 and called is False + + +def test_production_requires_explicit_flag(monkeypatch): + monkeypatch.setenv("BLIF_FLOW_V2_MODE", "compare") + with pytest.raises(RuntimeError, match="requires --production"): + sim.run(Namespace(production_readonly_simulation=False), runtime_loader=runtime(Connection())) + + +@pytest.mark.parametrize("database,readonly,mode,message", [ + ("test_db", "on", "compare", "current_database"), + ("clientflow", "off", "compare", "transaction_read_only"), + ("clientflow", "on", None, "explicit BLIF_FLOW"), + ("clientflow", "on", "off", "BLIF_FLOW_V2_MODE=compare"), + ("clientflow", "on", "shadow", "BLIF_FLOW_V2_MODE=compare"), + ("clientflow", "on", "authoritative", "BLIF_FLOW_V2_MODE=compare"), +]) +def test_production_guards(monkeypatch, database, readonly, mode, message): + if mode is None: + monkeypatch.delenv("BLIF_FLOW_V2_MODE", raising=False) + else: + monkeypatch.setenv("BLIF_FLOW_V2_MODE", mode) + with pytest.raises(RuntimeError, match=message): + sim.run(Namespace(production_readonly_simulation=True), + runtime_loader=runtime(Connection(database, readonly))) + + +def test_safe_simulation_uses_explicit_modes_without_changing_config_and_has_no_db_writes(monkeypatch, tmp_path): + monkeypatch.setenv("BLIF_FLOW_V2_MODE", "compare") + monkeypatch.setattr(sim, "JSON_PATH", tmp_path / "report.json") + monkeypatch.setattr(sim, "TEXT_PATH", tmp_path / "report.txt") + connection = Connection() + modes = [] + report = sim.run(Namespace(production_readonly_simulation=True), runtime_loader=runtime(connection, modes=modes)) + assert report["configured_mode"] == "compare" + assert sim.exit_code(report) == 0 + assert [sql for sql in connection.sql if sql.split(None, 1)[0].upper() in {"INSERT", "UPDATE", "DELETE", "MERGE"}] == [] + assert connection.sql[0] == "BEGIN READ ONLY" and connection.sql[-1] == "ROLLBACK" + assert modes == ["off", "authoritative"] + + +def safe_report(): + return sim.build_report(identity={}, configured_mode="compare", + v1={"work_items": []}, v2={"work_items": []}, + named_decisions=named_decisions(), missing_projections=0) + + +def test_unsafe_false_negative_is_nonzero(): + report = safe_report() + report["unsafe_false_negatives"] = 1 + assert sim.exit_code(report) != 0 + + +def test_missing_projection_is_nonzero(): + report = safe_report() + report["missing_projections"] = 1 + assert sim.exit_code(report) != 0 + + +def test_duplicate_current_card_is_nonzero(): + report = safe_report() + report["duplicate_current_cards"] = 1 + assert sim.exit_code(report) != 0 + + +def test_named_case_failure_is_nonzero(): + report = sim.build_report(identity={}, configured_mode="compare", v1={"work_items": []}, + v2={"work_items": []}, named_decisions=named_decisions(True), missing_projections=0) + assert sim.exit_code(report) != 0