"""Persistence scaffolding for the rebuildable BLIF Flow v2 projection. The factual sources remain authoritative. This module writes only the additive projection/audit tables introduced by migration 011 and never updates stages or tasks. """ from __future__ import annotations import hashlib import json import re from datetime import datetime, timezone from typing import Any, Iterable from sqlalchemy import text 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: if isinstance(value, datetime): return value.isoformat() if isinstance(value, dict): return {str(key): _jsonable(item) for key, item in value.items()} if isinstance(value, (list, tuple)): return [_jsonable(item) for item in value] return value def _evidence_refs(row: dict[str, Any]) -> list[dict[str, Any]]: evidence = row.get("evidence") or {} refs: list[dict[str, Any]] = [] for name in ("latest_relevant_inbound", "latest_relevant_outbound"): item = evidence.get(name) if item and item.get("id"): refs.append({"source": "message", "role": name, "id": item["id"], "at": item.get("at")}) for name, source in (("proforma", "jasmin_proforma"), ("invoice", "jasmin_invoice"), ("payment", "payment"), ("odoo", "odoo"), ("reconciliation", "reconciliation")): for item in evidence.get(name) or []: if item.get("id"): refs.append({ "source": source, "id": item["id"], "external_id": item.get("external_id"), "document_number": item.get("document_number"), "external_type": item.get("external_type"), "status": item.get("status"), }) return refs def _projection_value(row: dict[str, Any], derived_at: datetime) -> dict[str, Any]: raw = row["raw_v2"] duplicate = row.get("canonical_process_id") != row.get("opportunity_id") evidence_refs = _evidence_refs(row) reason_code = str(raw.get("precedence") or "business_transition").upper() stable = { "opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"], "canonical_opportunity_id": row["canonical_process_id"], "is_duplicate_representation": duplicate, "business_state": raw["business_state"], "business_next_action": raw.get("business_next_action"), "diagnostic_status": raw.get("diagnostic_status") or row.get("evidence", {}).get("diagnostic_status") or "clear", "confidence": raw.get("confidence") or "low", "reason_code": reason_code, "reason_text": raw.get("reason") or "", "evidence_refs": evidence_refs, "flow_version": FLOW_VERSION, } fingerprint = hashlib.sha256( json.dumps(_jsonable(stable), ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8") ).hexdigest() return stable | {"source_fingerprint": fingerprint, "derived_at": derived_at} 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=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"]) def rebuild_blif_flow_v2_projection( *, mode: str | None = None, derived_rows: Iterable[dict[str, Any]] | None = None, 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. ``off`` is a no-op. ``shadow`` and ``compare`` write only the same additive projection tables. Compare affects observation at the read boundary, never the persisted business rows. Authoritative remains deliberately disabled. """ selected_mode = str(mode or settings.blif_flow_v2_mode or "off").strip().lower() if selected_mode == "off": return {"mode": "off", "projection_count": 0, "transitions_written": 0, "disabled": True} if selected_mode not in {"shadow", "compare"}: raise RuntimeError(f"BLIF Flow v2 mode {selected_mode!r} is disabled; only off/shadow/compare are safe") 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( 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): raise RuntimeError("Flow v2 projection requires exactly one derived row per opportunity") owns_connection = connection is None if connection is None: from app.db import engine conn = engine.connect().execution_options(isolation_level="AUTOCOMMIT") else: conn = connection try: identity = conn.execute(text( "SELECT current_database(), current_user, current_setting('transaction_read_only')" )).one() if identity[0] not in allowed_databases: raise RuntimeError(f"refusing Flow v2 projection write to database {identity[0]!r}") if owns_connection: conn.exec_driver_sql("BEGIN READ WRITE") try: if conn.execute(text("SELECT current_setting('transaction_read_only')")).scalar_one() != "off": raise RuntimeError("Flow v2 projection rebuild requires an explicit READ WRITE transaction") conn.exec_driver_sql(f'SET LOCAL search_path TO "{target_schema}"') existing = { str(row["opportunity_id"]): dict(row) for row in conn.execute(text(""" SELECT opportunity_id::text, business_state, source_fingerprint FROM opportunity_flow_state_v2 """)).mappings() } transitions_written = 0 for value in values: previous = existing.get(value["opportunity_id"]) if previous is None or previous["business_state"] != value["business_state"]: result = conn.execute(text(""" INSERT INTO opportunity_flow_transitions ( opportunity_id, from_state, to_state, reason_code, reason_text, evidence_refs, flow_version, source_fingerprint ) VALUES ( CAST(:opportunity_id AS UUID), :from_state, :to_state, :reason_code, :reason_text, CAST(:evidence_refs AS JSONB), :flow_version, :source_fingerprint ) ON CONFLICT (opportunity_id, from_state, to_state, flow_version, source_fingerprint) DO NOTHING """), { **value, "from_state": previous["business_state"] if previous else None, "to_state": value["business_state"], "evidence_refs": json.dumps(_jsonable(value["evidence_refs"]), ensure_ascii=False), }) transitions_written += result.rowcount conn.execute(text(""" INSERT INTO opportunity_flow_state_v2 ( opportunity_id, material_process_key, canonical_opportunity_id, is_duplicate_representation, business_state, business_next_action, diagnostic_status, confidence, reason_code, reason_text, evidence_refs, flow_version, source_fingerprint, derived_at, updated_at ) VALUES ( CAST(:opportunity_id AS UUID), :material_process_key, CAST(:canonical_opportunity_id AS UUID), :is_duplicate_representation, :business_state, :business_next_action, :diagnostic_status, :confidence, :reason_code, :reason_text, CAST(:evidence_refs_json AS JSONB), :flow_version, :source_fingerprint, :derived_at, now() ) ON CONFLICT (opportunity_id) DO UPDATE SET material_process_key=EXCLUDED.material_process_key, canonical_opportunity_id=EXCLUDED.canonical_opportunity_id, is_duplicate_representation=EXCLUDED.is_duplicate_representation, business_state=EXCLUDED.business_state, business_next_action=EXCLUDED.business_next_action, diagnostic_status=EXCLUDED.diagnostic_status, confidence=EXCLUDED.confidence, reason_code=EXCLUDED.reason_code, reason_text=EXCLUDED.reason_text, evidence_refs=EXCLUDED.evidence_refs, flow_version=EXCLUDED.flow_version, source_fingerprint=EXCLUDED.source_fingerprint, derived_at=EXCLUDED.derived_at, updated_at=CASE WHEN opportunity_flow_state_v2.source_fingerprint IS DISTINCT FROM EXCLUDED.source_fingerprint THEN now() ELSE opportunity_flow_state_v2.updated_at END """), value | { "evidence_refs_json": json.dumps(_jsonable(value["evidence_refs"]), ensure_ascii=False) }) if owns_connection: conn.exec_driver_sql("COMMIT") except Exception: if owns_connection: conn.exec_driver_sql("ROLLBACK") raise finally: if owns_connection: conn.close() return { "mode": selected_mode, "projection_count": len(values), "canonical_count": sum(not value["is_duplicate_representation"] for value in values), "duplicate_count": sum(value["is_duplicate_representation"] for value in values), "transitions_written": transitions_written, "disabled": False, }