From 32e77739575d2b0a03f0c2349407e3ea9c30b293 Mon Sep 17 00:00:00 2001 From: plx Date: Sun, 16 Aug 2026 00:12:20 +0000 Subject: [PATCH] feat: add BLIF Flow v2 shadow compare cutover layer --- app/admin_ui/pages/opportunities.py | 20 +-- app/blif_flow_v2_projection_service.py | 11 +- app/config.py | 5 +- app/opportunity_next_action_service.py | 77 ++++++++- scripts/audit_blif_flow_v2_cutover.py | 228 +++++++++++++++++++++++++ tests/test_blif_flow_v2_cutover.py | 89 ++++++++++ tests/test_blif_flow_v2_persistence.py | 11 +- 7 files changed, 417 insertions(+), 24 deletions(-) create mode 100644 scripts/audit_blif_flow_v2_cutover.py create mode 100644 tests/test_blif_flow_v2_cutover.py diff --git a/app/admin_ui/pages/opportunities.py b/app/admin_ui/pages/opportunities.py index 32609f2..4605f08 100644 --- a/app/admin_ui/pages/opportunities.py +++ b/app/admin_ui/pages/opportunities.py @@ -1544,20 +1544,14 @@ def _parse_opportunity_dt(value: object): def _opportunity_lifecycle_state(opp: dict) -> str: state = str(opp.get("lifecycle_state") or "active").strip().lower() or "active" now = datetime.now(timezone.utc) - pending_call = str(opp.get("pending_follow_up_action_code") or "").strip().upper() == "CALL_CUSTOMER" - pending_call_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at")) - if pending_call: - return "follow_up_due" if pending_call_due and pending_call_due <= now else "scheduled_follow_up" - # Compatibility only: a stale denormalized lifecycle marker cannot invent - # a scheduled call when no active CALL_CUSTOMER task exists. - if state == "scheduled_follow_up": + pending_followup = str(opp.get("pending_follow_up_action_code") or "").strip().upper() + pending_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at")) + if pending_followup: + return "follow_up_due" if pending_due and pending_due <= now else "scheduled_follow_up" + # Compatibility-only timestamps/lifecycle markers cannot create work. An + # active pending follow-up task and its due_at are the operational source. + if state in {"scheduled_follow_up", "follow_up_due"}: return "active" - nurture_until = _parse_opportunity_dt(opp.get("nurture_until")) - next_follow_up = _parse_opportunity_dt(opp.get("next_follow_up_at")) - if state == "nurture" and nurture_until and nurture_until <= now: - return "follow_up_due" - if state in {"awaiting_customer", "active", "scheduled_follow_up"} and next_follow_up and next_follow_up <= now: - return "follow_up_due" return state diff --git a/app/blif_flow_v2_projection_service.py b/app/blif_flow_v2_projection_service.py index 1eb051a..b6d80ac 100644 --- a/app/blif_flow_v2_projection_service.py +++ b/app/blif_flow_v2_projection_service.py @@ -102,14 +102,15 @@ def rebuild_blif_flow_v2_projection( ) -> 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. + ``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 != "shadow": - raise RuntimeError(f"BLIF Flow v2 mode {selected_mode!r} is not implemented; only off/shadow are safe") + 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") @@ -208,7 +209,7 @@ def rebuild_blif_flow_v2_projection( conn.close() return { - "mode": "shadow", "projection_count": len(values), + "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, diff --git a/app/config.py b/app/config.py index 4386255..4a603ce 100644 --- a/app/config.py +++ b/app/config.py @@ -21,8 +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. + # Flow v2 remains non-authoritative. Compare observes V1/V2 differences + # while returning V1; authoritative is accepted by config only to fail + # closed at the switching boundary. blif_flow_v2_mode: Literal["off", "shadow", "compare", "authoritative"] = "off" clientflow_webhook_secret: str = "" diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index e512250..80fa3cf 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -7,11 +7,14 @@ from __future__ import annotations from dataclasses import asdict, dataclass from collections import defaultdict +import json +import logging from typing import Any, Dict, Iterable, Optional from sqlalchemy import bindparam, text from app.db import engine +from app.config import settings # Backward-compatible static anchors from v1.5.59: quotation_doc, invoice_doc, confirmar pagamento antes de emitir fatura, Pagamento confirmado com base em, Criar/enviar fatura. from app.domain.opportunity_flow import ( OpportunityEvidence, @@ -21,6 +24,9 @@ from app.domain.opportunity_flow import ( ) +logger = logging.getLogger(__name__) + + @dataclass class OpportunityNextAction: action_code: str @@ -37,6 +43,67 @@ class OpportunityNextAction: return asdict(self) +def _flow_v2_mode() -> str: + return str(settings.blif_flow_v2_mode or "off").strip().lower() + + +def _load_v2_comparison_rows(opportunity_ids: list[str]) -> dict[str, dict[str, Any]]: + """Read the persisted projection only; comparison must never derive writes.""" + if not opportunity_ids: + return {} + with engine.connect() as conn: + rows = 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, reason_code, reason_text, flow_version + FROM opportunity_flow_state_v2 + WHERE opportunity_id::text IN :opportunity_ids + """).bindparams(bindparam("opportunity_ids", expanding=True)), + {"opportunity_ids": opportunity_ids}).mappings().all() + return {str(row["opportunity_id"]): dict(row) for row in rows} + + +def compare_v1_v2_decisions( + v1_decisions: Dict[str, Dict[str, Any]], + v2_rows: dict[str, dict[str, Any]], +) -> list[dict[str, Any]]: + """Return payload-safe structured comparisons without message content.""" + comparisons = [] + for opportunity_id, v1 in v1_decisions.items(): + v2 = v2_rows.get(opportunity_id) + v1_action = str(v1.get("action_code") or "") or None + v2_action = (v2 or {}).get("business_next_action") + comparisons.append({ + "opportunity_id": opportunity_id, + "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")), + "v1_action": v1_action, + "v1_reason_code": v1.get("reason_if_blocked") or v1.get("decision_version"), + "v2_business_state": (v2 or {}).get("business_state"), + "v2_business_next_action": v2_action, + "v2_reason_code": (v2 or {}).get("reason_code"), + "v2_diagnostic_status": (v2 or {}).get("diagnostic_status"), + "v2_confidence": (v2 or {}).get("confidence"), + "projection_present": v2 is not None, + "action_agrees": v2 is not None and v1_action == v2_action, + "returned_source": "v1", + }) + return comparisons + + +def _observe_flow_v2(v1_decisions: Dict[str, Dict[str, Any]]) -> None: + mode = _flow_v2_mode() + if mode == "authoritative": + raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed") + if mode != "compare" or not v1_decisions: + return + rows = _load_v2_comparison_rows(list(v1_decisions)) + for comparison in compare_v1_v2_decisions(v1_decisions, rows): + logger.info("blif_flow_v2_compare %s", json.dumps(comparison, sort_keys=True, default=str)) + + def _first_row(conn: Any, sql: str, params: Dict[str, Any]) -> Optional[Dict[str, Any]]: row = conn.execute(text(sql), params).mappings().first() return dict(row) if row else None @@ -160,6 +227,8 @@ def get_opportunity_next_action( decision is now produced by the company workflow engine. """ + if _flow_v2_mode() == "authoritative": + raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed") evidence = (_build_db_evidence(opportunity_id, preloaded=preloaded) if preloaded is not None else _build_db_evidence(opportunity_id)) if evidence is None: @@ -172,8 +241,9 @@ def get_opportunity_next_action( reason_if_blocked="opportunity_not_found", ).to_dict() - decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile)) - return decision.to_dict() + decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile)).to_dict() + _observe_flow_v2({opportunity_id: decision}) + return decision def _bulk_statement(sql: str): @@ -298,6 +368,8 @@ def _bulk_operation_snapshots( def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Dict[str, Any]]: """Return the same decisions as the single-item API with a fixed query count.""" + if _flow_v2_mode() == "authoritative": + raise RuntimeError("BLIF Flow v2 authoritative mode is disabled; cutover boundary is fail-closed") ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip())) if not ids: return {} @@ -376,4 +448,5 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di company_profile="blif", ) decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict() + _observe_flow_v2(decisions) return decisions diff --git a/scripts/audit_blif_flow_v2_cutover.py b/scripts/audit_blif_flow_v2_cutover.py new file mode 100644 index 0000000..6b3c3ce --- /dev/null +++ b/scripts/audit_blif_flow_v2_cutover.py @@ -0,0 +1,228 @@ +#!/usr/bin/env python3 +"""Read-only compatibility/cutover audit for BLIF Flow v2.""" +from __future__ import annotations + +import json +import sys +from collections import Counter +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + +from sqlalchemy import text + +from app.db import engine +from app.opportunity_next_action_service import ( + _load_v2_comparison_rows, compare_v1_v2_decisions, get_opportunity_next_actions, +) +from scripts.simulate_blif_flow_v2 import collect + + +EXPECTED_DATABASE = "clientflow_codex_test" +CUTOVER = Path("/tmp/blif_flow_v2_cutover_report.json") +INVENTORY = Path("/tmp/blif_flow_v2_legacy_field_inventory.json") +COMPARE = Path("/tmp/blif_flow_v2_compare_report.json") + + +FIELD_INVENTORY = [ + { + "field": "stage", "kind": "derived_compatibility", "factual": False, + "writers": ["app/opportunity_service.py", "app/admin_ui/pages/opportunities.py", + "app/reconciliation_service.py", "app/odoo_service.py", "app/external_reconciliation_sync.py"], + "readers": ["app/domain/opportunity_flow/evidence.py (V1 only)", "app/workflow_guard.py", + "app/admin_ui/pages/opportunities.py", "app/admin_ui/pages/orders.py", + "app/admin_dashboard.py", "app/revenue_forecast_service.py"], + "flow_v2_replacement": "opportunity_flow_state_v2.business_state", + "compatibility_required": True, "one_time_migration_required": False, + "eventually_deprecatable": True, + "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.", + }, + { + "field": "lifecycle_state", "kind": "operational_compatibility", "factual": False, + "writers": ["app/followup_service.py", "app/opportunity_service.py"], + "readers": ["app/operational_eligibility.py", "app/operations_service.py", + "app/admin_ui/pages/opportunities.py"], + "flow_v2_replacement": "business state plus active operational obligations", + "compatibility_required": True, "one_time_migration_required": False, + "eventually_deprecatable": True, + "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"], + ["business next action", "opportunity_flow_state_v2.business_next_action"], + ["time-sensitive operational queue", "OperationalEligibility"], + ["scheduled follow-up", "pending task action_code + due_at"], + ["formal document state", "commercial_documents + active opportunity_document_links"], + ["payment", "confirmed factual payment evidence/operation link"], + ["Odoo execution", "operation_links / factual Odoo evidence"], + ["messages/customer response", "messages and communications chronology"], + ["legacy stage", "compatibility only"], + ["tasks", "operator obligations/history; never business fact proof"], +] + + +def _write(path: Path, value: Any) -> None: + 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]: + with engine.connect() as conn: + conn.exec_driver_sql("BEGIN READ ONLY") + try: + identity = conn.execute(text("SELECT current_database(),current_user,current_setting('transaction_read_only')")).one() + if identity[0] != EXPECTED_DATABASE or identity[2] != "on": + raise RuntimeError(f"refusing unexpected/non-read-only database: {identity!r}") + ids = [str(value) for value in conn.execute(text("SELECT opportunity_id FROM opportunity_flow_state_v2 ORDER BY opportunity_id")).scalars()] + historical = conn.execute(text(""" + SELECT count(*) FROM opportunities o + WHERE o.next_follow_up_at IS NOT NULL + AND NOT EXISTS ( + SELECT 1 FROM tasks t WHERE t.opportunity_id=o.id AND t.status='pending' + AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN + ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE')) + ) + """)).scalar_one() + finally: + conn.rollback() + return {"database": identity[0], "user": identity[1], "transaction_read_only": identity[2]}, ids, int(historical) + + +def main() -> None: + identity, ids, historical_timestamps = _identity_and_ids() + if len(ids) != 328: + raise RuntimeError(f"expected 328 persisted Flow v2 opportunities, found {len(ids)}") + v1 = get_opportunity_next_actions(ids) + v2 = _load_v2_comparison_rows(ids) + comparisons = compare_v1_v2_decisions(v1, v2) + counts = Counter("agree" if row["action_agrees"] else "different" for row in comparisons) + counts["missing_projection"] = sum(not row["projection_present"] for row in comparisons) + compare_report = { + "generated_at": datetime.now(timezone.utc), "database": identity, + "mode_semantics": "V1 returned; persisted V2 observed; no business mutation", + "opportunity_count": len(comparisons), "counts": dict(counts), "comparisons": comparisons, + } + _write(COMPARE, compare_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 = {} + for name in ("INSTALBEIRA", "PANORAMIC SUCCESS", "ENGEXICON", "CONSTRURECUP", "X MAT", "RZSOLAR"): + matches = [row for row in projection["opportunities"] + if name.casefold() in f"{row.get('title')} {row.get('customer')}".casefold()] + named[name] = [{"opportunity_id": row["opportunity_id"], "material_process_key": row["material_process_key"], + "canonical_process_id": row["canonical_process_id"], "business_state": row["safe_v2"]["business_state"], + "effective_action": row["safe_v2"]["effective_operational_action"], + "queue": row["safe_v2"]["effective_operational_queue"], + "duplicate_suppressed": row["safe_v2"]["precedence"] == "duplicate_representation"} + for row in matches] + cutover = { + "generated_at": datetime.now(timezone.utc), "database": identity, + "source_of_truth": [{"concern": concern, "source": source} for concern, source in SOURCE_OF_TRUTH], + "data_migration": {"opportunities_requiring_mutation_before_shadow": 0, + "opportunities_requiring_no_mutation_before_shadow": len(ids), + "broad_repair_required": False, + "stage_write_migration_required": False, + "lifecycle_write_migration_required": False, + "followup_timestamp_cleanup_required": False}, + "mode_contract": { + "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": [ + "Production migration 011 presence was not and must not be checked from this environment.", + "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) + print(json.dumps({"database": identity, "opportunities": len(ids), + "compare": compare_report["counts"], "historical_timestamps": historical_timestamps, + "mutation_required": 0, "outputs": [str(CUTOVER), str(INVENTORY), str(COMPARE)]}, indent=2)) + + +if __name__ == "__main__": + main() diff --git a/tests/test_blif_flow_v2_cutover.py b/tests/test_blif_flow_v2_cutover.py new file mode 100644 index 0000000..1f5752a --- /dev/null +++ b/tests/test_blif_flow_v2_cutover.py @@ -0,0 +1,89 @@ +from datetime import datetime, timedelta, timezone +from inspect import getsource +from types import SimpleNamespace + +import pytest + +import app.opportunity_next_action_service as service +from app.admin_ui.pages.opportunities import _opportunity_lifecycle_state +from app.domain.opportunity_flow.v2 import EffectiveOperationalDecision, suppress_duplicate_representation + + +class Decision: + def to_dict(self): + return {"action_code": "V1_ACTION", "label": "V1", "description": "legacy"} + + +def test_compare_mode_returns_v1_behavior(monkeypatch, caplog): + caplog.set_level("INFO") + monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "compare") + monkeypatch.setattr(service, "_build_db_evidence", lambda *args, **kwargs: SimpleNamespace(company_profile="blif")) + monkeypatch.setattr(service, "load_company_profile", lambda *_: object()) + monkeypatch.setattr(service, "decide_opportunity_next_action", lambda *_: Decision()) + monkeypatch.setattr(service, "_load_v2_comparison_rows", lambda ids: { + "opp": {"business_state": "COMPLETED", "business_next_action": None, + "reason_code": "BUSINESS_TRANSITION"} + }) + assert service.get_opportunity_next_action("opp", preloaded={})["action_code"] == "V1_ACTION" + assert "blif_flow_v2_compare" in caplog.text + + +def test_compare_mode_has_no_business_mutation_sql(): + source = getsource(service._load_v2_comparison_rows).upper() + getsource(service._observe_flow_v2).upper() + assert "UPDATE " not in source + assert "INSERT " not in source + assert "DELETE " not in source + + +def test_stale_legacy_stage_cannot_override_v2_factual_state(): + rows = service.compare_v1_v2_decisions( + {"opp": {"action_code": "CREATE_JASMIN_QUOTE"}}, + {"opp": {"business_state": "PROFORMA_REQUIRED", "business_next_action": "CREATE_PROFORMA", + "reason_code": "BUSINESS_TRANSITION"}}, + ) + assert rows[0]["v2_business_state"] == "PROFORMA_REQUIRED" + assert rows[0]["returned_source"] == "v1" + + +def test_historical_next_followup_without_active_task_creates_no_work(): + assert _opportunity_lifecycle_state({ + "lifecycle_state": "active", "next_follow_up_at": "2020-01-01T00:00:00+00:00", + "pending_follow_up_action_code": None, + }) == "active" + + +def test_due_active_followup_still_creates_work(): + assert _opportunity_lifecycle_state({ + "lifecycle_state": "awaiting_customer", "next_follow_up_at": "2020-01-01T00:00:00+00:00", + "pending_follow_up_action_code": "FOLLOW_UP_CUSTOMER_REVIEW", + "pending_follow_up_due_at": datetime.now(timezone.utc) - timedelta(days=1), + }) == "follow_up_due" + + +def test_future_active_followup_remains_scheduled(): + assert _opportunity_lifecycle_state({ + "lifecycle_state": "awaiting_customer", + "pending_follow_up_action_code": "FOLLOW_UP_CUSTOMER_REVIEW", + "pending_follow_up_due_at": datetime.now(timezone.utc) + timedelta(days=1), + }) == "scheduled_follow_up" + + +def test_duplicate_representation_stays_suppressed(): + raw = EffectiveOperationalDecision("REVIEW_REQUIRED", "REVIEW_REQUIRED", "REVIEW_REQUIRED", + "review", "review", "medium", "business_transition") + suppressed = suppress_duplicate_representation(raw, canonical_process_id="canonical") + assert suppressed.effective_operational_action is None + assert suppressed.effective_operational_queue == "not_current" + + +def test_authoritative_mode_is_fail_closed(monkeypatch): + monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "authoritative") + with pytest.raises(RuntimeError, match="disabled"): + service.get_opportunity_next_actions([]) + + +def test_shadow_mode_has_no_decision_side_effect(monkeypatch): + monkeypatch.setattr(service.settings, "blif_flow_v2_mode", "shadow") + monkeypatch.setattr(service, "_load_v2_comparison_rows", + lambda ids: pytest.fail("shadow must not read comparison projection")) + service._observe_flow_v2({"opp": {"action_code": "V1"}}) diff --git a/tests/test_blif_flow_v2_persistence.py b/tests/test_blif_flow_v2_persistence.py index 7eac097..2ab6747 100644 --- a/tests/test_blif_flow_v2_persistence.py +++ b/tests/test_blif_flow_v2_persistence.py @@ -80,7 +80,14 @@ def test_off_mode_is_a_noop_without_deriving_or_connecting(): } -@pytest.mark.parametrize("mode", ["compare", "authoritative", "invalid"]) +@pytest.mark.parametrize("mode", ["authoritative", "invalid"]) def test_unimplemented_modes_fail_closed(mode): - with pytest.raises(RuntimeError, match="not implemented"): + with pytest.raises(RuntimeError, match="disabled"): rebuild_blif_flow_v2_projection(mode=mode, derived_rows=[]) + + +def test_compare_mode_uses_the_same_additive_projection_path(): + from inspect import getsource + source = getsource(rebuild_blif_flow_v2_projection) + assert 'selected_mode not in {"shadow", "compare"}' in source + assert '"mode": selected_mode' in source