feat: add guarded production Flow v2 authoritative simulation
This commit is contained in:
@@ -8,6 +8,7 @@ from __future__ import annotations
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
from contextlib import nullcontext
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from sqlalchemy import text
|
||||
@@ -26,7 +27,7 @@ from app.canonical_operations import (
|
||||
opportunity_ids_from_work_seeds,
|
||||
partition_canonical_items,
|
||||
)
|
||||
from app.opportunity_next_action_service import get_opportunity_next_actions
|
||||
from app.opportunity_next_action_service import get_opportunity_next_actions, get_opportunity_next_actions_for_mode
|
||||
from app.document_reconciliation_service import active_document_link_exclusion_sql
|
||||
|
||||
|
||||
@@ -104,7 +105,7 @@ def _norm_identity(value: Any) -> str:
|
||||
|
||||
|
||||
def _identity_tokens(value: Any) -> set[str]:
|
||||
return {
|
||||
result = {
|
||||
token
|
||||
for token in _norm_identity(value).split()
|
||||
if len(token) >= 3 and token not in _IDENTITY_STOPWORDS
|
||||
@@ -302,7 +303,9 @@ def _attach_operation_urls(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
||||
return _sanitize_operation_identities(cleaned)
|
||||
|
||||
|
||||
def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
def get_operations_summary(
|
||||
limit: int = 24, *, flow_mode: str | None = None, connection: Any = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Build the /operations work queue summary.
|
||||
|
||||
/operations is intentionally not a mini-dashboard. It returns a compact
|
||||
@@ -315,7 +318,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
# after canonicalization. This prevents duplicate rows from consuming the
|
||||
# display limit while still avoiding an unbounded history scan.
|
||||
candidate_limit = max(200, min(display_limit * 20, 1000))
|
||||
with engine.begin() as conn:
|
||||
active_flow_mode = str(flow_mode if flow_mode is not None else settings.blif_flow_v2_mode or "off").strip().lower()
|
||||
with (nullcontext(connection) if connection is not None else engine.begin()) as conn:
|
||||
counts = conn.execute(text("""
|
||||
SELECT
|
||||
(SELECT COUNT(*) FROM opportunities WHERE status = 'open')::int AS open_opportunities,
|
||||
@@ -675,9 +679,9 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
# In authoritative mode the factual projection, not a legacy task/stage,
|
||||
# selects material processes that have current business work. Synthetic
|
||||
# rows are read-model seeds only and are never persisted.
|
||||
if str(settings.blif_flow_v2_mode or "off").strip().lower() == "authoritative":
|
||||
if active_flow_mode == "authoritative":
|
||||
seeded = {str(row.get("opportunity_id") or "") for row in seed_rows}
|
||||
with engine.connect() as conn:
|
||||
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
|
||||
v2_seeds = conn.execute(text("""
|
||||
SELECT p.opportunity_id::text, p.business_next_action, p.derived_at,
|
||||
o.title, o.value_amount, o.currency, o.lifecycle_state,
|
||||
@@ -713,7 +717,9 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")})
|
||||
normalized_rows = _normalise_work_item_intent(seed_rows)
|
||||
opportunity_ids = opportunity_ids_from_work_seeds(normalized_rows)
|
||||
decisions = get_opportunity_next_actions(opportunity_ids) if opportunity_ids else {}
|
||||
decisions = (get_opportunity_next_actions_for_mode(
|
||||
opportunity_ids, flow_mode=active_flow_mode, connection=connection,
|
||||
) if opportunity_ids else {})
|
||||
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
|
||||
canonical_items = _attach_operation_urls(list(projection["items"]))
|
||||
partition = partition_canonical_items(canonical_items, display_limit=display_limit)
|
||||
@@ -758,6 +764,11 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
|
||||
**diagnostics,
|
||||
"projection_metrics": diagnostics,
|
||||
}
|
||||
if flow_mode is not None:
|
||||
# Complete, untruncated read-model population for offline simulation;
|
||||
# normal API/GET callers do not receive this audit-only field.
|
||||
result["simulation_all_items"] = canonical_items
|
||||
return result
|
||||
|
||||
|
||||
def get_system_health_summary() -> Dict[str, Any]:
|
||||
|
||||
@@ -376,13 +376,24 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
|
||||
if not ids:
|
||||
return {}
|
||||
|
||||
if _flow_v2_mode() == "authoritative":
|
||||
return _get_authoritative_opportunity_decisions(ids)
|
||||
return get_opportunity_next_actions_for_mode(ids, flow_mode=_flow_v2_mode())
|
||||
|
||||
|
||||
def get_opportunity_next_actions_for_mode(
|
||||
opportunity_ids: Iterable[str], *, flow_mode: str, connection: Any = None,
|
||||
) -> Dict[str, Dict[str, Any]]:
|
||||
"""Explicit-mode read boundary used by the read-only cutover simulation."""
|
||||
ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip()))
|
||||
if not ids:
|
||||
return {}
|
||||
if flow_mode == "authoritative":
|
||||
return _get_authoritative_opportunity_decisions(ids, connection=connection)
|
||||
|
||||
# Read path only: schema creation belongs to startup/migrations. In
|
||||
# particular, GET /operations must never perform DDL while calculating its
|
||||
# canonical projection.
|
||||
with engine.begin() as conn:
|
||||
from contextlib import nullcontext
|
||||
with (nullcontext(connection) if connection is not None else engine.begin()) as conn:
|
||||
opportunities = _bulk_rows(conn, """
|
||||
SELECT id::text, stage, status, title,
|
||||
local_customer_id::text AS fiscal_customer_id,
|
||||
@@ -453,13 +464,15 @@ 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)
|
||||
if flow_mode == "compare":
|
||||
_observe_flow_v2(decisions)
|
||||
return decisions
|
||||
|
||||
|
||||
def _get_authoritative_opportunity_decisions(ids: list[str]) -> Dict[str, Dict[str, Any]]:
|
||||
def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any = None) -> Dict[str, Dict[str, Any]]:
|
||||
"""Read-only set-oriented adapter input loader; never derives or persists V2."""
|
||||
with engine.connect() as conn:
|
||||
from contextlib import nullcontext
|
||||
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
|
||||
projections = _bulk_rows(conn, """
|
||||
SELECT opportunity_id::text, material_process_key,
|
||||
canonical_opportunity_id::text, is_duplicate_representation,
|
||||
|
||||
Reference in New Issue
Block a user