diff --git a/app/blif_flow_v2_projection_service.py b/app/blif_flow_v2_projection_service.py index b6d80ac..cfc807f 100644 --- a/app/blif_flow_v2_projection_service.py +++ b/app/blif_flow_v2_projection_service.py @@ -19,6 +19,7 @@ from app.config import settings FLOW_VERSION = "blif-flow-v2-shadow-3" WRITE_DATABASE_ALLOWLIST = frozenset({"clientflow_codex_test"}) +PROJECTION_WRITE_TABLES = frozenset({"opportunity_flow_state_v2", "opportunity_flow_transitions"}) def _jsonable(value: Any) -> Any: @@ -78,16 +79,24 @@ def _projection_value(row: dict[str, Any], derived_at: datetime) -> dict[str, An return stable | {"source_fingerprint": fingerprint, "derived_at": derived_at} -def _derive_all() -> list[dict[str, Any]]: +def _derive_all( + *, + expected_database: str = "clientflow_codex_test", + expected_user: str | None = "clientflow_codex_test", + expected_opportunity_count: int | None = 328, + require_opportunities: bool = False, +) -> list[dict[str, Any]]: # Reuse the validated shadow evidence adapter without making it authoritative. from scripts.simulate_blif_flow_v2 import collect report = collect( - expected_database="clientflow_codex_test", - expected_user="clientflow_codex_test", + expected_database=expected_database, + expected_user=expected_user, # collect() still opens its factual read phase with BEGIN READ ONLY; # the session default may be read-write in the isolated test database. require_read_only=False, + expected_opportunity_count=expected_opportunity_count, + require_opportunities=require_opportunities, ) return list(report["opportunities"]) @@ -99,6 +108,10 @@ def rebuild_blif_flow_v2_projection( allowed_databases: frozenset[str] = WRITE_DATABASE_ALLOWLIST, target_schema: str = "public", connection: Any | None = None, + derive_expected_database: str = "clientflow_codex_test", + derive_expected_user: str | None = "clientflow_codex_test", + expected_opportunity_count: int | None = 328, + require_opportunities: bool = False, ) -> dict[str, Any]: """Idempotently rebuild projection rows and state-change transitions. @@ -114,7 +127,14 @@ def rebuild_blif_flow_v2_projection( if not re.fullmatch(r"[a-z_][a-z0-9_]*", target_schema): raise ValueError("invalid target_schema") - rows = list(derived_rows) if derived_rows is not None else _derive_all() + rows = list(derived_rows) if derived_rows is not None else _derive_all( + expected_database=derive_expected_database, + expected_user=derive_expected_user, + expected_opportunity_count=expected_opportunity_count, + require_opportunities=require_opportunities, + ) + if require_opportunities and not rows: + raise RuntimeError("Flow v2 projection requires at least one opportunity") derived_at = datetime.now(timezone.utc) values = [_projection_value(row, derived_at) for row in rows] if len({value["opportunity_id"] for value in values}) != len(values): diff --git a/scripts/rebuild_blif_flow_v2_projection.py b/scripts/rebuild_blif_flow_v2_projection.py index f030163..4d67ca1 100644 --- a/scripts/rebuild_blif_flow_v2_projection.py +++ b/scripts/rebuild_blif_flow_v2_projection.py @@ -1,18 +1,84 @@ #!/usr/bin/env python3 -"""Rebuild additive BLIF Flow v2 projection tables in safe shadow mode.""" +"""Rebuild additive BLIF Flow v2 projection tables with explicit safeguards.""" from __future__ import annotations +import argparse import json import os import sys from pathlib import Path +from typing import Sequence + ROOT = Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) os.chdir(ROOT) -from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument( + "--production-shadow", action="store_true", + help="explicitly authorize projection-only shadow writes to database clientflow", + ) + return parser + + +def validate_execution(*, production_shadow: bool, mode: str, database: str) -> None: + """Validate CLI intent independently of DATABASE_URL inference.""" + normalized_mode = str(mode or "").strip().lower() + if production_shadow: + if normalized_mode != "shadow": + raise RuntimeError("--production-shadow requires BLIF_FLOW_V2_MODE=shadow") + if database != "clientflow": + raise RuntimeError( + f"--production-shadow requires database 'clientflow', found {database!r}" + ) + return + if database == "clientflow": + raise RuntimeError("production database requires explicit --production-shadow opt-in") + + +def main(argv: Sequence[str] | None = None) -> int: + # argparse handles --help and exits before any application/DB import below. + args = build_parser().parse_args(argv) + + from sqlalchemy import text + from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection + from app.config import settings + from app.db import engine + + mode = str(settings.blif_flow_v2_mode or "").strip().lower() + with engine.connect() as conn: + identity = conn.execute(text( + "SELECT current_database(), current_user, current_setting('transaction_read_only')" + )).one() + database, user, transaction_read_only = identity + validate_execution( + production_shadow=args.production_shadow, mode=mode, database=database, + ) + print(json.dumps({ + "database": database, "user": user, "blif_flow_v2_mode": mode, + "transaction_read_only": transaction_read_only, + "production_shadow_opt_in": args.production_shadow, + }, sort_keys=True), flush=True) + + if args.production_shadow: + result = rebuild_blif_flow_v2_projection( + mode="shadow", + # Production is deliberately scoped to this invocation; the + # module-level default allowlist remains test-only. + allowed_databases=frozenset({"clientflow"}), + derive_expected_database="clientflow", + derive_expected_user=user, + expected_opportunity_count=None, + require_opportunities=True, + ) + else: + result = rebuild_blif_flow_v2_projection() + print(json.dumps(result, indent=2, sort_keys=True)) + return 0 if __name__ == "__main__": - print(json.dumps(rebuild_blif_flow_v2_projection(), indent=2, sort_keys=True)) + raise SystemExit(main()) diff --git a/scripts/simulate_blif_flow_v2.py b/scripts/simulate_blif_flow_v2.py index 1fdca4d..302fa7b 100644 --- a/scripts/simulate_blif_flow_v2.py +++ b/scripts/simulate_blif_flow_v2.py @@ -530,6 +530,8 @@ def collect( *, expected_database: str = "clientflow_codex_shadow", expected_user: str | None = "clientflow_codex", require_read_only: bool = True, + expected_opportunity_count: int | None = None, + require_opportunities: bool = False, ) -> dict[str, Any]: data = _load( expected_database=expected_database, @@ -537,8 +539,12 @@ def collect( require_read_only=require_read_only, ) opportunities = data["opportunities"] - if len(opportunities) != 328: - raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}") + if expected_opportunity_count is not None and len(opportunities) != expected_opportunity_count: + raise RuntimeError( + f"expected {expected_opportunity_count} opportunities, found {len(opportunities)}" + ) + if require_opportunities and not opportunities: + raise RuntimeError("Flow v2 projection requires at least one opportunity") ids = [_s(opp["id"]) for opp in opportunities] v1_decisions = get_opportunity_next_actions(ids) operations = get_operations_summary(limit=200) diff --git a/tests/test_blif_flow_v2_production_shadow_rebuild.py b/tests/test_blif_flow_v2_production_shadow_rebuild.py new file mode 100644 index 0000000..0982ae2 --- /dev/null +++ b/tests/test_blif_flow_v2_production_shadow_rebuild.py @@ -0,0 +1,104 @@ +from inspect import getsource + +import pytest + +import app.blif_flow_v2_projection_service as projection +import scripts.rebuild_blif_flow_v2_projection as cli +import scripts.simulate_blif_flow_v2 as simulator + + +def test_help_performs_no_rebuild(monkeypatch, capsys): + monkeypatch.setattr(cli, "validate_execution", lambda **kwargs: pytest.fail("help reached execution")) + with pytest.raises(SystemExit) as exc: + cli.main(["--help"]) + assert exc.value.code == 0 + assert "--production-shadow" in capsys.readouterr().out + + +def test_default_execution_does_not_authorize_production(): + with pytest.raises(RuntimeError, match="explicit --production-shadow"): + cli.validate_execution(production_shadow=False, mode="shadow", database="clientflow") + + +@pytest.mark.parametrize("mode", ["", "off", "compare", "authoritative"]) +def test_production_shadow_requires_exact_shadow_mode(mode): + with pytest.raises(RuntimeError, match="requires BLIF_FLOW_V2_MODE=shadow"): + cli.validate_execution(production_shadow=True, mode=mode, database="clientflow") + + +@pytest.mark.parametrize("database", ["clientflow_codex_test", "clientflow_codex_shadow", "other"]) +def test_production_shadow_requires_exact_production_database(database): + with pytest.raises(RuntimeError, match="requires database 'clientflow'"): + cli.validate_execution(production_shadow=True, mode="shadow", database=database) + + +def test_production_shadow_explicit_proof_is_accepted(): + cli.validate_execution(production_shadow=True, mode="shadow", database="clientflow") + + +def test_global_write_allowlist_remains_test_only(): + assert projection.WRITE_DATABASE_ALLOWLIST == frozenset({"clientflow_codex_test"}) + assert "clientflow" not in projection.WRITE_DATABASE_ALLOWLIST + + +def test_factual_read_transaction_remains_read_only(): + source = getsource(simulator._load) + assert 'conn.execute(text("BEGIN READ ONLY"))' in source + assert 'conn.execute(text("ROLLBACK"))' in source + + +def test_projection_writer_only_mutates_two_additive_tables(): + source = getsource(projection.rebuild_blif_flow_v2_projection).upper() + assert projection.PROJECTION_WRITE_TABLES == { + "opportunity_flow_state_v2", "opportunity_flow_transitions"} + assert "INSERT INTO OPPORTUNITY_FLOW_TRANSITIONS" in source + assert "INSERT INTO OPPORTUNITY_FLOW_STATE_V2" in source + for forbidden in ("OPPORTUNITIES", "TASKS", "MESSAGES", "COMMUNICATIONS", + "COMMERCIAL_DOCUMENTS", "CUSTOMERS", "PAYMENTS", "OPERATION_LINKS"): + assert f"INSERT INTO {forbidden}" not in source + assert f"UPDATE {forbidden}" not in source + assert f"DELETE FROM {forbidden}" not in source + + +def test_production_derivation_has_no_fixed_328_requirement(monkeypatch): + captured = {} + + def fake_collect(**kwargs): + captured.update(kwargs) + return {"opportunities": [{"opportunity_id": "one"}]} + + monkeypatch.setattr(simulator, "collect", fake_collect) + rows = projection._derive_all( + expected_database="clientflow", expected_user="runtime-role", + expected_opportunity_count=None, require_opportunities=True, + ) + assert rows == [{"opportunity_id": "one"}] + assert captured["expected_opportunity_count"] is None + assert captured["require_opportunities"] is True + + +def test_snapshot_expected_count_validation_is_still_available(monkeypatch): + monkeypatch.setattr(simulator, "_load", lambda **kwargs: {"opportunities": [object()]}) + with pytest.raises(RuntimeError, match="expected 328 opportunities, found 1"): + simulator.collect(expected_opportunity_count=328) + + +def test_empty_production_universe_is_rejected(monkeypatch): + monkeypatch.setattr(simulator, "collect", lambda **kwargs: {"opportunities": []}) + assert projection._derive_all( + expected_database="clientflow", expected_user=None, + expected_opportunity_count=None, require_opportunities=True, + ) == [] + with pytest.raises(RuntimeError, match="at least one opportunity"): + projection.rebuild_blif_flow_v2_projection( + mode="shadow", derived_rows=[], require_opportunities=True, + ) + + +def test_existing_test_defaults_remain_guarded(monkeypatch): + captured = {} + monkeypatch.setattr(simulator, "collect", lambda **kwargs: captured.update(kwargs) or {"opportunities": []}) + projection._derive_all() + assert captured["expected_database"] == "clientflow_codex_test" + assert captured["expected_user"] == "clientflow_codex_test" + assert captured["expected_opportunity_count"] == 328