Compare commits
1 Commits
feat/blif-
...
fix/blif-f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e3b8750ebb |
@@ -19,6 +19,7 @@ from app.config import settings
|
|||||||
|
|
||||||
FLOW_VERSION = "blif-flow-v2-shadow-3"
|
FLOW_VERSION = "blif-flow-v2-shadow-3"
|
||||||
WRITE_DATABASE_ALLOWLIST = frozenset({"clientflow_codex_test"})
|
WRITE_DATABASE_ALLOWLIST = frozenset({"clientflow_codex_test"})
|
||||||
|
PROJECTION_WRITE_TABLES = frozenset({"opportunity_flow_state_v2", "opportunity_flow_transitions"})
|
||||||
|
|
||||||
|
|
||||||
def _jsonable(value: Any) -> Any:
|
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}
|
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.
|
# Reuse the validated shadow evidence adapter without making it authoritative.
|
||||||
from scripts.simulate_blif_flow_v2 import collect
|
from scripts.simulate_blif_flow_v2 import collect
|
||||||
|
|
||||||
report = collect(
|
report = collect(
|
||||||
expected_database="clientflow_codex_test",
|
expected_database=expected_database,
|
||||||
expected_user="clientflow_codex_test",
|
expected_user=expected_user,
|
||||||
# collect() still opens its factual read phase with BEGIN READ ONLY;
|
# collect() still opens its factual read phase with BEGIN READ ONLY;
|
||||||
# the session default may be read-write in the isolated test database.
|
# the session default may be read-write in the isolated test database.
|
||||||
require_read_only=False,
|
require_read_only=False,
|
||||||
|
expected_opportunity_count=expected_opportunity_count,
|
||||||
|
require_opportunities=require_opportunities,
|
||||||
)
|
)
|
||||||
return list(report["opportunities"])
|
return list(report["opportunities"])
|
||||||
|
|
||||||
@@ -99,6 +108,10 @@ def rebuild_blif_flow_v2_projection(
|
|||||||
allowed_databases: frozenset[str] = WRITE_DATABASE_ALLOWLIST,
|
allowed_databases: frozenset[str] = WRITE_DATABASE_ALLOWLIST,
|
||||||
target_schema: str = "public",
|
target_schema: str = "public",
|
||||||
connection: Any | None = None,
|
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]:
|
) -> dict[str, Any]:
|
||||||
"""Idempotently rebuild projection rows and state-change transitions.
|
"""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):
|
if not re.fullmatch(r"[a-z_][a-z0-9_]*", target_schema):
|
||||||
raise ValueError("invalid 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)
|
derived_at = datetime.now(timezone.utc)
|
||||||
values = [_projection_value(row, derived_at) for row in rows]
|
values = [_projection_value(row, derived_at) for row in rows]
|
||||||
if len({value["opportunity_id"] for value in values}) != len(values):
|
if len({value["opportunity_id"] for value in values}) != len(values):
|
||||||
|
|||||||
@@ -1,18 +1,84 @@
|
|||||||
#!/usr/bin/env python3
|
#!/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
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import argparse
|
||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
from typing import 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)
|
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.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__":
|
if __name__ == "__main__":
|
||||||
print(json.dumps(rebuild_blif_flow_v2_projection(), indent=2, sort_keys=True))
|
raise SystemExit(main())
|
||||||
|
|||||||
@@ -530,6 +530,8 @@ def collect(
|
|||||||
*, expected_database: str = "clientflow_codex_shadow",
|
*, expected_database: str = "clientflow_codex_shadow",
|
||||||
expected_user: str | None = "clientflow_codex",
|
expected_user: str | None = "clientflow_codex",
|
||||||
require_read_only: bool = True,
|
require_read_only: bool = True,
|
||||||
|
expected_opportunity_count: int | None = None,
|
||||||
|
require_opportunities: bool = False,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
data = _load(
|
data = _load(
|
||||||
expected_database=expected_database,
|
expected_database=expected_database,
|
||||||
@@ -537,8 +539,12 @@ def collect(
|
|||||||
require_read_only=require_read_only,
|
require_read_only=require_read_only,
|
||||||
)
|
)
|
||||||
opportunities = data["opportunities"]
|
opportunities = data["opportunities"]
|
||||||
if len(opportunities) != 328:
|
if expected_opportunity_count is not None and len(opportunities) != expected_opportunity_count:
|
||||||
raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}")
|
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]
|
ids = [_s(opp["id"]) for opp in opportunities]
|
||||||
v1_decisions = get_opportunity_next_actions(ids)
|
v1_decisions = get_opportunity_next_actions(ids)
|
||||||
operations = get_operations_summary(limit=200)
|
operations = get_operations_summary(limit=200)
|
||||||
|
|||||||
104
tests/test_blif_flow_v2_production_shadow_rebuild.py
Normal file
104
tests/test_blif_flow_v2_production_shadow_rebuild.py
Normal 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
|
||||||
Reference in New Issue
Block a user