fix: allow guarded production Flow v2 shadow rebuild

This commit is contained in:
plx
2026-08-16 00:39:23 +00:00
parent 32e7773957
commit e3b8750ebb
4 changed files with 205 additions and 9 deletions

View File

@@ -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):

View File

@@ -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)
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())

View File

@@ -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)

View File

@@ -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