feat: add BLIF Flow v2 persistence scaffolding
This commit is contained in:
@@ -40,6 +40,16 @@ INTERNAL_ACTION_CODES = {
|
|||||||
|
|
||||||
ACTION_CODES = TRIAGE_ACTION_CODES | INTERNAL_ACTION_CODES
|
ACTION_CODES = TRIAGE_ACTION_CODES | INTERNAL_ACTION_CODES
|
||||||
|
|
||||||
|
# Persistence vocabulary for the shadow-only Flow v2 projection. These are not
|
||||||
|
# added to TRIAGE_ACTION_CODES or ACTION_CODES, so no existing task/LLM/runtime
|
||||||
|
# behavior changes.
|
||||||
|
FLOW_V2_BUSINESS_ACTION_CODES = {
|
||||||
|
"CREATE_PROFORMA",
|
||||||
|
"CREATE_INVOICE",
|
||||||
|
"VALIDATE_ODOO_ORDER",
|
||||||
|
"COMPLETE_OPPORTUNITY",
|
||||||
|
}
|
||||||
|
|
||||||
ACTION_MAP = {
|
ACTION_MAP = {
|
||||||
"CALL_CUSTOMER": {
|
"CALL_CUSTOMER": {
|
||||||
"route": "vendas",
|
"route": "vendas",
|
||||||
|
|||||||
215
app/blif_flow_v2_projection_service.py
Normal file
215
app/blif_flow_v2_projection_service.py
Normal file
@@ -0,0 +1,215 @@
|
|||||||
|
"""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`` writes only additive projection tables.
|
||||||
|
Compare/authoritative behavior is deliberately not implemented.
|
||||||
|
"""
|
||||||
|
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 != "shadow":
|
||||||
|
raise RuntimeError(f"BLIF Flow v2 mode {selected_mode!r} is not implemented; only off/shadow 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": "shadow", "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,
|
||||||
|
}
|
||||||
@@ -21,6 +21,9 @@ class Settings(BaseSettings):
|
|||||||
# TIMESTAMPTZ nem os índices usados pelo schema core.
|
# TIMESTAMPTZ nem os índices usados pelo schema core.
|
||||||
database_url: str
|
database_url: str
|
||||||
clientflow_persist: bool = True
|
clientflow_persist: bool = True
|
||||||
|
# Flow v2 is additive shadow scaffolding only. Compare/authoritative values
|
||||||
|
# are reserved for future work and are not activated by this implementation.
|
||||||
|
blif_flow_v2_mode: Literal["off", "shadow", "compare", "authoritative"] = "off"
|
||||||
|
|
||||||
clientflow_webhook_secret: str = ""
|
clientflow_webhook_secret: str = ""
|
||||||
clientflow_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy"
|
clientflow_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy"
|
||||||
|
|||||||
52
migrations/011_blif_flow_v2_persistence.sql
Normal file
52
migrations/011_blif_flow_v2_persistence.sql
Normal file
@@ -0,0 +1,52 @@
|
|||||||
|
-- Additive, rebuildable BLIF Flow v2 projection scaffolding.
|
||||||
|
-- This migration does not alter opportunities.stage or activate Flow v2.
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS opportunity_flow_state_v2 (
|
||||||
|
opportunity_id UUID PRIMARY KEY REFERENCES opportunities(id) ON DELETE CASCADE,
|
||||||
|
material_process_key TEXT NOT NULL,
|
||||||
|
canonical_opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE CASCADE,
|
||||||
|
is_duplicate_representation BOOLEAN NOT NULL DEFAULT FALSE,
|
||||||
|
business_state TEXT NOT NULL,
|
||||||
|
business_next_action TEXT,
|
||||||
|
diagnostic_status TEXT NOT NULL DEFAULT 'clear',
|
||||||
|
confidence TEXT NOT NULL DEFAULT 'low',
|
||||||
|
reason_code TEXT NOT NULL,
|
||||||
|
reason_text TEXT NOT NULL,
|
||||||
|
evidence_refs JSONB NOT NULL DEFAULT '[]'::jsonb,
|
||||||
|
flow_version TEXT NOT NULL,
|
||||||
|
source_fingerprint TEXT NOT NULL,
|
||||||
|
derived_at TIMESTAMPTZ NOT NULL,
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_material_process_key
|
||||||
|
ON opportunity_flow_state_v2(material_process_key);
|
||||||
|
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_canonical
|
||||||
|
ON opportunity_flow_state_v2(canonical_opportunity_id);
|
||||||
|
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_state_v2_business_state
|
||||||
|
ON opportunity_flow_state_v2(business_state);
|
||||||
|
|
||||||
|
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolution_code TEXT;
|
||||||
|
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolved_at TIMESTAMPTZ;
|
||||||
|
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS resolved_by_event_id UUID REFERENCES opportunity_events(id) ON DELETE SET NULL;
|
||||||
|
ALTER TABLE tasks ADD COLUMN IF NOT EXISTS superseded_by_task_id UUID REFERENCES tasks(id) ON DELETE SET NULL;
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS ix_tasks_resolved_by_event_id ON tasks(resolved_by_event_id);
|
||||||
|
CREATE INDEX IF NOT EXISTS ix_tasks_superseded_by_task_id ON tasks(superseded_by_task_id);
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS opportunity_flow_transitions (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
opportunity_id UUID NOT NULL REFERENCES opportunities(id) ON DELETE CASCADE,
|
||||||
|
from_state TEXT,
|
||||||
|
to_state TEXT NOT NULL,
|
||||||
|
reason_code TEXT NOT NULL,
|
||||||
|
reason_text TEXT NOT NULL,
|
||||||
|
evidence_refs JSONB NOT NULL DEFAULT '[]'::jsonb,
|
||||||
|
flow_version TEXT NOT NULL,
|
||||||
|
source_fingerprint TEXT NOT NULL,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
UNIQUE (opportunity_id, from_state, to_state, flow_version, source_fingerprint)
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS ix_opportunity_flow_transitions_opportunity_created
|
||||||
|
ON opportunity_flow_transitions(opportunity_id, created_at DESC);
|
||||||
12
migrations/011_blif_flow_v2_persistence_down.sql
Normal file
12
migrations/011_blif_flow_v2_persistence_down.sql
Normal file
@@ -0,0 +1,12 @@
|
|||||||
|
-- Reversible rollback for BLIF Flow v2 persistence scaffolding.
|
||||||
|
|
||||||
|
DROP TABLE IF EXISTS opportunity_flow_transitions;
|
||||||
|
|
||||||
|
DROP INDEX IF EXISTS ix_tasks_superseded_by_task_id;
|
||||||
|
DROP INDEX IF EXISTS ix_tasks_resolved_by_event_id;
|
||||||
|
ALTER TABLE tasks DROP COLUMN IF EXISTS superseded_by_task_id;
|
||||||
|
ALTER TABLE tasks DROP COLUMN IF EXISTS resolved_by_event_id;
|
||||||
|
ALTER TABLE tasks DROP COLUMN IF EXISTS resolved_at;
|
||||||
|
ALTER TABLE tasks DROP COLUMN IF EXISTS resolution_code;
|
||||||
|
|
||||||
|
DROP TABLE IF EXISTS opportunity_flow_state_v2;
|
||||||
18
scripts/rebuild_blif_flow_v2_projection.py
Normal file
18
scripts/rebuild_blif_flow_v2_projection.py
Normal file
@@ -0,0 +1,18 @@
|
|||||||
|
#!/usr/bin/env python3
|
||||||
|
"""Rebuild additive BLIF Flow v2 projection tables in safe shadow mode."""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
ROOT = Path(__file__).resolve().parents[1]
|
||||||
|
sys.path.insert(0, str(ROOT))
|
||||||
|
os.chdir(ROOT)
|
||||||
|
|
||||||
|
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
print(json.dumps(rebuild_blif_flow_v2_projection(), indent=2, sort_keys=True))
|
||||||
@@ -92,14 +92,20 @@ def _group(rows: Iterable[dict[str, Any]], key: str = "opportunity_id") -> dict[
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
||||||
def _load() -> dict[str, Any]:
|
def _load(
|
||||||
|
*, expected_database: str = "clientflow_codex_shadow",
|
||||||
|
expected_user: str | None = "clientflow_codex",
|
||||||
|
require_read_only: bool = True,
|
||||||
|
) -> dict[str, Any]:
|
||||||
with engine.connect() as conn:
|
with engine.connect() as conn:
|
||||||
conn = conn.execution_options(isolation_level="AUTOCOMMIT")
|
conn = conn.execution_options(isolation_level="AUTOCOMMIT")
|
||||||
identity = conn.execute(text(
|
identity = conn.execute(text(
|
||||||
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
||||||
)).one()
|
)).one()
|
||||||
if tuple(identity) != ("clientflow_codex_shadow", "clientflow_codex", "on"):
|
if identity[0] != expected_database or (expected_user and identity[1] != expected_user):
|
||||||
raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
|
raise RuntimeError(f"refusing unexpected database identity: {identity!r}")
|
||||||
|
if require_read_only and identity[2] != "on":
|
||||||
|
raise RuntimeError(f"read-only simulation requires transaction_read_only=on: {identity!r}")
|
||||||
conn.execute(text("BEGIN READ ONLY"))
|
conn.execute(text("BEGIN READ ONLY"))
|
||||||
try:
|
try:
|
||||||
opportunities = [dict(row) for row in conn.execute(text("""
|
opportunities = [dict(row) for row in conn.execute(text("""
|
||||||
@@ -520,8 +526,16 @@ def _totals(items: Iterable[dict[str, Any]], queue_key: str) -> dict[str, int]:
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
def collect() -> dict[str, Any]:
|
def collect(
|
||||||
data = _load()
|
*, expected_database: str = "clientflow_codex_shadow",
|
||||||
|
expected_user: str | None = "clientflow_codex",
|
||||||
|
require_read_only: bool = True,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
data = _load(
|
||||||
|
expected_database=expected_database,
|
||||||
|
expected_user=expected_user,
|
||||||
|
require_read_only=require_read_only,
|
||||||
|
)
|
||||||
opportunities = data["opportunities"]
|
opportunities = data["opportunities"]
|
||||||
if len(opportunities) != 328:
|
if len(opportunities) != 328:
|
||||||
raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}")
|
raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}")
|
||||||
|
|||||||
119
scripts/validate_blif_flow_v2_persistence.py
Normal file
119
scripts/validate_blif_flow_v2_persistence.py
Normal file
@@ -0,0 +1,119 @@
|
|||||||
|
#!/usr/bin/env python3
|
||||||
|
"""Validate migration 011 and two rebuilds using temp tables in the test DB."""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
|
ROOT = Path(__file__).resolve().parents[1]
|
||||||
|
sys.path.insert(0, str(ROOT))
|
||||||
|
os.chdir(ROOT)
|
||||||
|
|
||||||
|
from app.blif_flow_v2_projection_service import rebuild_blif_flow_v2_projection
|
||||||
|
from app.db import engine
|
||||||
|
|
||||||
|
|
||||||
|
NAMED_IDS = {
|
||||||
|
"instalbeira": "5c33db95-fab8-477a-bddd-0b9cc8f91302",
|
||||||
|
"engexicon": "61f1c955-a372-4ea7-b9b0-b8528d74a141",
|
||||||
|
"construrecup": "e3b23ac5-84db-4763-8a31-a684e873032c",
|
||||||
|
"panoramic": "fd79b9a1-07e6-4f61-95e8-09eab89c155e",
|
||||||
|
"x_mat_canonical": "dc89a466-db24-401b-bfe9-d47644b2d0c8",
|
||||||
|
"x_mat_reconstructed": "1816a06e-9a69-4a9b-9279-1263156892d3",
|
||||||
|
"rzsolar_reconstructed": "fd221608-e007-4043-a23d-07e0c119a345",
|
||||||
|
"rzsolar_synthetic": "434124fb-ac19-4d78-909a-55761d7e8daa",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def main() -> None:
|
||||||
|
conn = engine.connect().execution_options(isolation_level="AUTOCOMMIT")
|
||||||
|
try:
|
||||||
|
identity = conn.execute(text(
|
||||||
|
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
|
||||||
|
)).one()
|
||||||
|
if tuple(identity[:2]) != ("clientflow_codex_test", "clientflow_codex_test"):
|
||||||
|
raise RuntimeError(f"refusing persistence validation on {identity!r}")
|
||||||
|
conn.exec_driver_sql("BEGIN READ WRITE")
|
||||||
|
try:
|
||||||
|
stage_fingerprint_before = conn.execute(text("""
|
||||||
|
SELECT md5(string_agg(id::text || ':' || COALESCE(stage,''), ',' ORDER BY id))
|
||||||
|
FROM public.opportunities
|
||||||
|
""")).scalar_one()
|
||||||
|
conn.exec_driver_sql("SET LOCAL search_path TO pg_temp, public")
|
||||||
|
conn.exec_driver_sql("CREATE TEMP TABLE opportunities (id UUID PRIMARY KEY)")
|
||||||
|
conn.exec_driver_sql("INSERT INTO opportunities SELECT id FROM public.opportunities")
|
||||||
|
conn.exec_driver_sql("CREATE TEMP TABLE opportunity_events (LIKE public.opportunity_events INCLUDING DEFAULTS INCLUDING CONSTRAINTS)")
|
||||||
|
conn.exec_driver_sql("ALTER TABLE opportunity_events ADD PRIMARY KEY (id)")
|
||||||
|
conn.exec_driver_sql("CREATE TEMP TABLE tasks (LIKE public.tasks INCLUDING DEFAULTS INCLUDING CONSTRAINTS)")
|
||||||
|
conn.exec_driver_sql("ALTER TABLE tasks ADD PRIMARY KEY (id)")
|
||||||
|
conn.exec_driver_sql(Path("migrations/011_blif_flow_v2_persistence.sql").read_text())
|
||||||
|
|
||||||
|
first = rebuild_blif_flow_v2_projection(
|
||||||
|
mode="shadow", target_schema="pg_temp", connection=conn,
|
||||||
|
)
|
||||||
|
transition_count_first = conn.execute(text(
|
||||||
|
"SELECT count(*) FROM pg_temp.opportunity_flow_transitions"
|
||||||
|
)).scalar_one()
|
||||||
|
first_fingerprints = dict(conn.execute(text("""
|
||||||
|
SELECT opportunity_id::text, source_fingerprint
|
||||||
|
FROM pg_temp.opportunity_flow_state_v2
|
||||||
|
""")).all())
|
||||||
|
second = rebuild_blif_flow_v2_projection(
|
||||||
|
mode="shadow", target_schema="pg_temp", connection=conn,
|
||||||
|
)
|
||||||
|
transition_count_second = conn.execute(text(
|
||||||
|
"SELECT count(*) FROM pg_temp.opportunity_flow_transitions"
|
||||||
|
)).scalar_one()
|
||||||
|
second_fingerprints = dict(conn.execute(text("""
|
||||||
|
SELECT opportunity_id::text, source_fingerprint
|
||||||
|
FROM pg_temp.opportunity_flow_state_v2
|
||||||
|
""")).all())
|
||||||
|
named = {}
|
||||||
|
for name, oid in NAMED_IDS.items():
|
||||||
|
row = conn.execute(text("""
|
||||||
|
SELECT opportunity_id::text, material_process_key,
|
||||||
|
canonical_opportunity_id::text, is_duplicate_representation,
|
||||||
|
business_state, business_next_action, diagnostic_status, confidence
|
||||||
|
FROM pg_temp.opportunity_flow_state_v2
|
||||||
|
WHERE opportunity_id=CAST(:oid AS UUID)
|
||||||
|
"""), {"oid": oid}).mappings().one()
|
||||||
|
named[name] = dict(row)
|
||||||
|
stage_fingerprint_after = conn.execute(text("""
|
||||||
|
SELECT md5(string_agg(id::text || ':' || COALESCE(stage,''), ',' ORDER BY id))
|
||||||
|
FROM public.opportunities
|
||||||
|
""")).scalar_one()
|
||||||
|
result = {
|
||||||
|
"database": {"name": identity[0], "user": identity[1]},
|
||||||
|
"migration_scope": "transaction-scoped pg_temp (public tasks is postgres-owned)",
|
||||||
|
"first_rebuild": first, "second_rebuild": second,
|
||||||
|
"transition_count_first": transition_count_first,
|
||||||
|
"transition_count_second": transition_count_second,
|
||||||
|
"idempotent": transition_count_first == transition_count_second
|
||||||
|
and first_fingerprints == second_fingerprints,
|
||||||
|
"opportunity_stage_unchanged": stage_fingerprint_before == stage_fingerprint_after,
|
||||||
|
"named": named,
|
||||||
|
}
|
||||||
|
conn.exec_driver_sql(Path("migrations/011_blif_flow_v2_persistence_down.sql").read_text())
|
||||||
|
remaining_task_columns = conn.execute(text("""
|
||||||
|
SELECT count(*) FROM information_schema.columns
|
||||||
|
WHERE table_schema LIKE 'pg_temp_%' AND table_name='tasks'
|
||||||
|
AND column_name IN ('resolution_code','resolved_at','resolved_by_event_id','superseded_by_task_id')
|
||||||
|
""")).scalar_one()
|
||||||
|
result["down_migration_reversible"] = (
|
||||||
|
conn.execute(text("SELECT to_regclass('pg_temp.opportunity_flow_state_v2')")).scalar_one() is None
|
||||||
|
and conn.execute(text("SELECT to_regclass('pg_temp.opportunity_flow_transitions')")).scalar_one() is None
|
||||||
|
and remaining_task_columns == 0
|
||||||
|
)
|
||||||
|
print(json.dumps(result, ensure_ascii=False, indent=2, sort_keys=True))
|
||||||
|
finally:
|
||||||
|
conn.exec_driver_sql("ROLLBACK")
|
||||||
|
finally:
|
||||||
|
conn.close()
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
main()
|
||||||
86
tests/test_blif_flow_v2_persistence.py
Normal file
86
tests/test_blif_flow_v2_persistence.py
Normal file
@@ -0,0 +1,86 @@
|
|||||||
|
from datetime import datetime, timezone
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from app.action_catalog import ACTION_CODES, FLOW_V2_BUSINESS_ACTION_CODES, TRIAGE_ACTION_CODES
|
||||||
|
from app.blif_flow_v2_projection_service import _projection_value, rebuild_blif_flow_v2_projection
|
||||||
|
|
||||||
|
|
||||||
|
def derived_row(**overrides):
|
||||||
|
row = {
|
||||||
|
"opportunity_id": "11111111-1111-1111-1111-111111111111",
|
||||||
|
"material_process_key": "odoo_sale_name:s00001",
|
||||||
|
"canonical_process_id": "11111111-1111-1111-1111-111111111111",
|
||||||
|
"raw_v2": {
|
||||||
|
"business_state": "PROFORMA_REQUIRED",
|
||||||
|
"business_next_action": "CREATE_PROFORMA",
|
||||||
|
"diagnostic_status": "clear",
|
||||||
|
"confidence": "high",
|
||||||
|
"precedence": "business_transition",
|
||||||
|
"reason": "Current order intent requires a proforma.",
|
||||||
|
},
|
||||||
|
"evidence": {
|
||||||
|
"latest_relevant_inbound": {"id": "m1", "at": "2026-08-15T00:00:00+00:00"},
|
||||||
|
"latest_relevant_outbound": None, "proforma": [], "invoice": [],
|
||||||
|
"payment": [], "odoo": [], "reconciliation": [],
|
||||||
|
},
|
||||||
|
}
|
||||||
|
row.update(overrides)
|
||||||
|
return row
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_011_is_additive_reversible_and_does_not_touch_stage():
|
||||||
|
up = Path("migrations/011_blif_flow_v2_persistence.sql").read_text()
|
||||||
|
down = Path("migrations/011_blif_flow_v2_persistence_down.sql").read_text()
|
||||||
|
assert "CREATE TABLE IF NOT EXISTS opportunity_flow_state_v2" in up
|
||||||
|
assert "CREATE TABLE IF NOT EXISTS opportunity_flow_transitions" in up
|
||||||
|
for column in ("resolution_code", "resolved_at", "resolved_by_event_id", "superseded_by_task_id"):
|
||||||
|
assert f"ADD COLUMN IF NOT EXISTS {column}" in up
|
||||||
|
assert f"DROP COLUMN IF EXISTS {column}" in down
|
||||||
|
assert "DROP TABLE IF EXISTS opportunity_flow_state_v2" in down
|
||||||
|
assert "UPDATE opportunities" not in up
|
||||||
|
assert "stage =" not in up
|
||||||
|
|
||||||
|
|
||||||
|
def test_flow_v2_actions_are_separate_from_existing_triage_and_runtime_actions():
|
||||||
|
assert FLOW_V2_BUSINESS_ACTION_CODES == {
|
||||||
|
"CREATE_PROFORMA", "CREATE_INVOICE", "VALIDATE_ODOO_ORDER", "COMPLETE_OPPORTUNITY",
|
||||||
|
}
|
||||||
|
assert "CREATE_JASMIN_QUOTE" not in FLOW_V2_BUSINESS_ACTION_CODES
|
||||||
|
assert FLOW_V2_BUSINESS_ACTION_CODES.isdisjoint(TRIAGE_ACTION_CODES)
|
||||||
|
assert FLOW_V2_BUSINESS_ACTION_CODES.isdisjoint(ACTION_CODES)
|
||||||
|
|
||||||
|
|
||||||
|
def test_projection_fingerprint_is_stable_and_excludes_derived_time():
|
||||||
|
first = _projection_value(derived_row(), datetime(2026, 8, 15, tzinfo=timezone.utc))
|
||||||
|
second = _projection_value(derived_row(), datetime(2026, 8, 16, tzinfo=timezone.utc))
|
||||||
|
assert first["source_fingerprint"] == second["source_fingerprint"]
|
||||||
|
assert first["business_next_action"] == "CREATE_PROFORMA"
|
||||||
|
assert first["evidence_refs"] == [{
|
||||||
|
"source": "message", "role": "latest_relevant_inbound", "id": "m1",
|
||||||
|
"at": "2026-08-15T00:00:00+00:00",
|
||||||
|
}]
|
||||||
|
|
||||||
|
|
||||||
|
def test_duplicate_projection_persists_canonical_material_identity():
|
||||||
|
row = derived_row(
|
||||||
|
opportunity_id="22222222-2222-2222-2222-222222222222",
|
||||||
|
canonical_process_id="11111111-1111-1111-1111-111111111111",
|
||||||
|
)
|
||||||
|
value = _projection_value(row, datetime.now(timezone.utc))
|
||||||
|
assert value["is_duplicate_representation"] is True
|
||||||
|
assert value["canonical_opportunity_id"] == "11111111-1111-1111-1111-111111111111"
|
||||||
|
assert value["material_process_key"] == "odoo_sale_name:s00001"
|
||||||
|
|
||||||
|
|
||||||
|
def test_off_mode_is_a_noop_without_deriving_or_connecting():
|
||||||
|
assert rebuild_blif_flow_v2_projection(mode="off") == {
|
||||||
|
"mode": "off", "projection_count": 0, "transitions_written": 0, "disabled": True,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("mode", ["compare", "authoritative", "invalid"])
|
||||||
|
def test_unimplemented_modes_fail_closed(mode):
|
||||||
|
with pytest.raises(RuntimeError, match="not implemented"):
|
||||||
|
rebuild_blif_flow_v2_projection(mode=mode, derived_rows=[])
|
||||||
Reference in New Issue
Block a user