217 lines
10 KiB
Python
217 lines
10 KiB
Python
"""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"})
|
|
|
|
|
|
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() -> 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",
|
|
# 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,
|
|
)
|
|
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,
|
|
) -> 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()
|
|
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,
|
|
}
|