From e99da64d8b87bd6304a8be24e484e9392028724a Mon Sep 17 00:00:00 2001 From: plx Date: Sat, 15 Aug 2026 23:30:37 +0000 Subject: [PATCH] feat: add BLIF Flow v2 persistence scaffolding --- app/action_catalog.py | 10 + app/blif_flow_v2_projection_service.py | 215 ++++++++++++++++++ app/config.py | 3 + migrations/011_blif_flow_v2_persistence.sql | 52 +++++ .../011_blif_flow_v2_persistence_down.sql | 12 + scripts/rebuild_blif_flow_v2_projection.py | 18 ++ scripts/simulate_blif_flow_v2.py | 22 +- scripts/validate_blif_flow_v2_persistence.py | 119 ++++++++++ tests/test_blif_flow_v2_persistence.py | 86 +++++++ 9 files changed, 533 insertions(+), 4 deletions(-) create mode 100644 app/blif_flow_v2_projection_service.py create mode 100644 migrations/011_blif_flow_v2_persistence.sql create mode 100644 migrations/011_blif_flow_v2_persistence_down.sql create mode 100644 scripts/rebuild_blif_flow_v2_projection.py create mode 100644 scripts/validate_blif_flow_v2_persistence.py create mode 100644 tests/test_blif_flow_v2_persistence.py diff --git a/app/action_catalog.py b/app/action_catalog.py index 49cfccf..a99e03d 100644 --- a/app/action_catalog.py +++ b/app/action_catalog.py @@ -40,6 +40,16 @@ 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 = { "CALL_CUSTOMER": { "route": "vendas", diff --git a/app/blif_flow_v2_projection_service.py b/app/blif_flow_v2_projection_service.py new file mode 100644 index 0000000..1eb051a --- /dev/null +++ b/app/blif_flow_v2_projection_service.py @@ -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, + } diff --git a/app/config.py b/app/config.py index 2702166..4386255 100644 --- a/app/config.py +++ b/app/config.py @@ -21,6 +21,9 @@ class Settings(BaseSettings): # TIMESTAMPTZ nem os índices usados pelo schema core. database_url: str 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_admin_auth_mode: Literal["proxy", "token", "local"] = "proxy" diff --git a/migrations/011_blif_flow_v2_persistence.sql b/migrations/011_blif_flow_v2_persistence.sql new file mode 100644 index 0000000..001c1e2 --- /dev/null +++ b/migrations/011_blif_flow_v2_persistence.sql @@ -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); diff --git a/migrations/011_blif_flow_v2_persistence_down.sql b/migrations/011_blif_flow_v2_persistence_down.sql new file mode 100644 index 0000000..022d42c --- /dev/null +++ b/migrations/011_blif_flow_v2_persistence_down.sql @@ -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; diff --git a/scripts/rebuild_blif_flow_v2_projection.py b/scripts/rebuild_blif_flow_v2_projection.py new file mode 100644 index 0000000..f030163 --- /dev/null +++ b/scripts/rebuild_blif_flow_v2_projection.py @@ -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)) diff --git a/scripts/simulate_blif_flow_v2.py b/scripts/simulate_blif_flow_v2.py index 19f808b..1fdca4d 100644 --- a/scripts/simulate_blif_flow_v2.py +++ b/scripts/simulate_blif_flow_v2.py @@ -92,14 +92,20 @@ def _group(rows: Iterable[dict[str, Any]], key: str = "opportunity_id") -> dict[ 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: conn = conn.execution_options(isolation_level="AUTOCOMMIT") identity = conn.execute(text( "SELECT current_database(), current_user, current_setting('transaction_read_only')" )).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}") + 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")) try: 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]: - data = _load() +def collect( + *, 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"] if len(opportunities) != 328: raise RuntimeError(f"expected 328 opportunities, found {len(opportunities)}") diff --git a/scripts/validate_blif_flow_v2_persistence.py b/scripts/validate_blif_flow_v2_persistence.py new file mode 100644 index 0000000..ab0b839 --- /dev/null +++ b/scripts/validate_blif_flow_v2_persistence.py @@ -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() diff --git a/tests/test_blif_flow_v2_persistence.py b/tests/test_blif_flow_v2_persistence.py new file mode 100644 index 0000000..7eac097 --- /dev/null +++ b/tests/test_blif_flow_v2_persistence.py @@ -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=[])