Compare commits
2 Commits
feat/blif-
...
fix/blif-f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e456cbc0e4 | ||
|
|
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,228 +1,328 @@
|
|||||||
#!/usr/bin/env python3
|
#!/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
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import argparse
|
||||||
import json
|
import json
|
||||||
|
import os
|
||||||
import sys
|
import sys
|
||||||
from collections import Counter
|
from collections import Counter
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any
|
from typing import Any, 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)
|
||||||
|
|
||||||
from sqlalchemy import text
|
PRODUCTION_DATABASE = "clientflow"
|
||||||
|
TEST_DATABASE = "clientflow_codex_test"
|
||||||
from app.db import engine
|
AUDIT = Path("/tmp/blif_flow_v2_production_shadow_audit.json")
|
||||||
from app.opportunity_next_action_service import (
|
SEMANTIC = Path("/tmp/blif_flow_v2_production_semantic_compare.json")
|
||||||
_load_v2_comparison_rows, compare_v1_v2_decisions, get_opportunity_next_actions,
|
SUMMARY = Path("/tmp/blif_flow_v2_production_shadow_summary.txt")
|
||||||
)
|
CLASSIFICATIONS = {
|
||||||
from scripts.simulate_blif_flow_v2 import collect
|
"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",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
EXPECTED_DATABASE = "clientflow_codex_test"
|
def build_parser() -> argparse.ArgumentParser:
|
||||||
CUTOVER = Path("/tmp/blif_flow_v2_cutover_report.json")
|
parser = argparse.ArgumentParser(description=__doc__)
|
||||||
INVENTORY = Path("/tmp/blif_flow_v2_legacy_field_inventory.json")
|
parser.add_argument(
|
||||||
COMPARE = Path("/tmp/blif_flow_v2_compare_report.json")
|
"--production-readonly-audit", action="store_true",
|
||||||
|
help="explicitly authorize a read-only audit of database clientflow",
|
||||||
|
)
|
||||||
|
return parser
|
||||||
|
|
||||||
|
|
||||||
FIELD_INVENTORY = [
|
def validate_execution(
|
||||||
{
|
*, production_readonly_audit: bool, database: str, transaction_read_only: str,
|
||||||
"field": "stage", "kind": "derived_compatibility", "factual": False,
|
mode: str,
|
||||||
"writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py",
|
) -> None:
|
||||||
"app/reconciliation_service.py", "app/odoo_service.py", "app/external_reconciliation_sync.py"],
|
normalized_mode = str(mode or "").strip().lower()
|
||||||
"readers": ["app/domain/opportunity_flow/evidence.py (V1 only)", "app/workflow_guard.py",
|
if production_readonly_audit:
|
||||||
"app/admin_ui/pages/opportunities.py", "app/admin_ui/pages/orders.py",
|
if database != PRODUCTION_DATABASE:
|
||||||
"app/admin_dashboard.py", "app/revenue_forecast_service.py"],
|
raise RuntimeError(
|
||||||
"flow_v2_replacement": "opportunity_flow_state_v2.business_state",
|
f"--production-readonly-audit requires database {PRODUCTION_DATABASE!r}, found {database!r}"
|
||||||
"compatibility_required": True, "one_time_migration_required": False,
|
)
|
||||||
"eventually_deprecatable": True,
|
if transaction_read_only != "on":
|
||||||
"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.",
|
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")
|
||||||
"field": "lifecycle_state", "kind": "operational_compatibility", "factual": False,
|
return
|
||||||
"writers": ["app/followup_service.py", "app/opportunity_service.py"],
|
if database == PRODUCTION_DATABASE:
|
||||||
"readers": ["app/operational_eligibility.py", "app/operations_service.py",
|
raise RuntimeError("production database requires explicit --production-readonly-audit opt-in")
|
||||||
"app/admin_ui/pages/opportunities.py"],
|
if database != TEST_DATABASE:
|
||||||
"flow_v2_replacement": "business state plus active operational obligations",
|
raise RuntimeError(f"default audit requires database {TEST_DATABASE!r}, found {database!r}")
|
||||||
"compatibility_required": True, "one_time_migration_required": False,
|
if transaction_read_only != "on":
|
||||||
"eventually_deprecatable": True,
|
raise RuntimeError("cutover audit requires an explicit READ ONLY transaction")
|
||||||
"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.",
|
|
||||||
},
|
|
||||||
]
|
|
||||||
|
|
||||||
SOURCE_OF_TRUTH = [
|
|
||||||
["business process state", "opportunity_flow_state_v2.business_state"],
|
def _code(value: Any) -> str:
|
||||||
["business next action", "opportunity_flow_state_v2.business_next_action"],
|
return str(value or "").strip().upper()
|
||||||
["time-sensitive operational queue", "OperationalEligibility"],
|
|
||||||
["scheduled follow-up", "pending task action_code + due_at"],
|
|
||||||
["formal document state", "commercial_documents + active opportunity_document_links"],
|
def classify_semantic_difference(
|
||||||
["payment", "confirmed factual payment evidence/operation link"],
|
*, v1_state: str | None, v1_action: str | None,
|
||||||
["Odoo execution", "operation_links / factual Odoo evidence"],
|
v2_state: str | None, v2_action: str | None,
|
||||||
["messages/customer response", "messages and communications chronology"],
|
v2_projection_present: bool = True,
|
||||||
["legacy stage", "compatibility only"],
|
is_duplicate_representation: bool = False,
|
||||||
["tasks", "operator obligations/history; never business fact proof"],
|
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:
|
def _write(path: Path, value: Any) -> None:
|
||||||
path.write_text(json.dumps(value, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
|
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:
|
with engine.connect() as conn:
|
||||||
conn.exec_driver_sql("BEGIN READ ONLY")
|
conn.exec_driver_sql("BEGIN READ ONLY")
|
||||||
try:
|
try:
|
||||||
identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one()
|
identity = conn.execute(text(
|
||||||
if identity[0] != EXPECTED_DATABASE or identity[2] != "on":
|
"SELECT current_database(),current_user,current_setting('transaction_read_only')"
|
||||||
raise RuntimeError(f"refusing unexpected/non-read-only database: {identity!r}")
|
)).one()
|
||||||
ids = [str(value) for value in conn.execute(text("SELECT opportunity_id FROM opportunity_flow_state_v2 ORDER BY opportunity_id")).scalars()]
|
opportunity_ids = [str(value) for value in conn.execute(text(
|
||||||
historical = conn.execute(text("""
|
"SELECT id FROM opportunities ORDER BY id"
|
||||||
SELECT count(*) FROM opportunities o
|
)).scalars()]
|
||||||
WHERE o.next_follow_up_at IS NOT NULL
|
rows = conn.execute(text("""
|
||||||
AND NOT EXISTS (
|
SELECT opportunity_id::text, material_process_key,
|
||||||
SELECT 1 FROM tasks t WHERE t.opportunity_id=o.id AND t.status='pending'
|
canonical_opportunity_id::text, is_duplicate_representation,
|
||||||
AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN
|
business_state, business_next_action, diagnostic_status,
|
||||||
('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
|
confidence, reason_code, reason_text, evidence_refs, flow_version
|
||||||
)
|
FROM opportunity_flow_state_v2 ORDER BY opportunity_id
|
||||||
""")).scalar_one()
|
""")).mappings().all()
|
||||||
finally:
|
finally:
|
||||||
conn.rollback()
|
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:
|
def run_audit(*, production_readonly_audit: bool) -> dict[str, Any]:
|
||||||
identity, ids, historical_timestamps = _identity_and_ids()
|
# Imports occur only after argparse, so --help cannot initialize DB code.
|
||||||
if len(ids) != 328:
|
from app.config import settings
|
||||||
raise RuntimeError(f"expected 328 persisted Flow v2 opportunities, found {len(ids)}")
|
from app.db import engine
|
||||||
v1 = get_opportunity_next_actions(ids)
|
from scripts.simulate_blif_flow_v2 import collect
|
||||||
v2 = _load_v2_comparison_rows(ids)
|
|
||||||
comparisons = compare_v1_v2_decisions(v1, v2)
|
identity, opportunity_ids, persisted = _database_snapshot(engine)
|
||||||
counts = Counter("agree" if row["action_agrees"] else "different" for row in comparisons)
|
mode = str(settings.blif_flow_v2_mode or "").strip().lower()
|
||||||
counts["missing_projection"] = sum(not row["projection_present"] for row in comparisons)
|
validate_execution(
|
||||||
compare_report = {
|
production_readonly_audit=production_readonly_audit,
|
||||||
"generated_at": datetime.now(timezone.utc), "database": identity,
|
database=identity["database"], transaction_read_only=identity["transaction_read_only"],
|
||||||
"mode_semantics": "V1 returned; persisted V2 observed; no business mutation",
|
mode=mode,
|
||||||
"opportunity_count": len(comparisons), "counts": dict(counts), "comparisons": comparisons,
|
)
|
||||||
|
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 = {}
|
named = {}
|
||||||
for name in ("INSTALBEIRA", "PANORAMIC SUCCESS", "ENGEXICON", "CONSTRURECUP", "X MAT", "RZSOLAR"):
|
by_id = {row["opportunity_id"]: row for row in comparisons}
|
||||||
matches = [row for row in projection["opportunities"]
|
for name, opportunity_id in NAMED_IDS.items():
|
||||||
if name.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()]
|
row = by_id[opportunity_id]
|
||||||
named[name] = [{"opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"],
|
named[name] = row
|
||||||
"canonical_process_id": row["canonical_process_id"], "business_state": row["safe_v2"]["business_state"],
|
# Use the complete canonical Operations candidate universe, including
|
||||||
"effective_action": row["safe_v2"]["effective_operational_action"],
|
# preserved standalone obligations, rather than opportunity cards alone.
|
||||||
"queue": row["safe_v2"]["effective_operational_queue"],
|
v1_metrics = _operation_totals(projection["v1_totals"])
|
||||||
"duplicate_suppressed": row["safe_v2"]["precedence"] == "duplicate_representation"}
|
safe_metrics = _operation_totals(projection["safe_v2_totals"])
|
||||||
for row in matches]
|
audit = {
|
||||||
cutover = {
|
"generated_at": datetime.now(timezone.utc), "preflight": preflight,
|
||||||
"generated_at": datetime.now(timezone.utc), "database": identity,
|
"read_only": True, "business_writes": 0, "projection_writes": 0,
|
||||||
"source_of_truth": [{"concern": concern, "source": source} for concern, source in SOURCE_OF_TRUTH],
|
"opportunity_count": len(opportunity_ids), "projection_count": len(persisted),
|
||||||
"data_migration": {"opportunities_requiring_mutation_before_shadow": 0,
|
"canonical_count": sum(not row.get("is_duplicate_representation") for row in persisted.values()),
|
||||||
"opportunities_requiring_no_mutation_before_shadow": len(ids),
|
"duplicate_representation_count": sum(bool(row.get("is_duplicate_representation")) for row in persisted.values()),
|
||||||
"broad_repair_required": False,
|
"operations_metrics": {"v1": v1_metrics, "safe_v2": safe_metrics},
|
||||||
"stage_write_migration_required": False,
|
"semantic_classification_counts": dict(sorted(counts.items())),
|
||||||
"lifecycle_write_migration_required": False,
|
"authoritative_cutover_blockers": {
|
||||||
"followup_timestamp_cleanup_required": False},
|
"real_conflicts": len(real_conflicts), "missing_projections": len(missing),
|
||||||
"mode_contract": {
|
"blocked": bool(real_conflicts or missing),
|
||||||
"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.",
|
|
||||||
},
|
},
|
||||||
"production_shadow_blockers": [
|
"real_conflicts": real_conflicts, "missing_projections": missing,
|
||||||
"Production migration 011 presence was not and must not be checked from this environment.",
|
"named_cases": named,
|
||||||
"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,
|
|
||||||
}
|
}
|
||||||
_write(CUTOVER, cutover)
|
_write(AUDIT, audit)
|
||||||
print(json.dumps({"database": identity, "opportunities": len(ids),
|
lines = [
|
||||||
"compare": compare_report["counts"], "historical_timestamps": historical_timestamps,
|
"BLIF FLOW V2 PRODUCTION SHADOW READ-ONLY AUDIT", "",
|
||||||
"mutation_required": 0, "outputs": [str(CUTOVER), str(INVENTORY), str(COMPARE)]}, indent=2))
|
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__":
|
if __name__ == "__main__":
|
||||||
main()
|
raise SystemExit(main())
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|
||||||
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__":
|
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)
|
||||||
|
|||||||
110
tests/test_blif_flow_v2_production_readonly_audit.py
Normal file
110
tests/test_blif_flow_v2_production_readonly_audit.py
Normal 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",
|
||||||
|
}
|
||||||
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