"""Central decision service for opportunity next actions. v4928.1.5.60 delegates operational flow to app.domain.opportunity_flow so opportunity pages, tasks and future audits consume the same decision vocabulary. """ 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 from app.authoritative_operational_adapter import decide_authoritative_operation # 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, build_opportunity_evidence, decide_opportunity_next_action, load_company_profile, ) logger = logging.getLogger(__name__) @dataclass class OpportunityNextAction: action_code: str label: str description: str priority: str = "normal" target_url: Optional[str] = None can_execute: bool = True reason_if_blocked: Optional[str] = None document_id: Optional[str] = None document_number: Optional[str] = None def to_dict(self) -> Dict[str, Any]: 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 {} # Lightweight repository fakes used by pure unit tests intentionally expose # only ``begin``. Comparison observation is optional and must not alter the # V1 contract or its query count when that read capability is absent. if not hasattr(engine, "connect"): 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 != "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 def _rows(conn: Any, sql: str, params: Dict[str, Any]) -> list[Dict[str, Any]]: return [dict(r) for r in conn.execute(text(sql), params).mappings().all()] def _operation_snapshot_safe(opportunity_id: str) -> dict[str, Any]: try: from app.operation_service import get_operation_snapshot return get_operation_snapshot(opportunity_id) except Exception: return {"cards": [], "links": []} def _build_db_evidence( opportunity_id: str, preloaded: Optional[Dict[str, Any]] = None, ) -> OpportunityEvidence | None: preloaded = preloaded or {} params = {"opportunity_id": opportunity_id} with engine.begin() as conn: opp = preloaded.get("opportunity") or _first_row(conn, """ SELECT id::text, stage, status, title, local_customer_id::text AS fiscal_customer_id, local_customer_id::text AS customer_id, metadata FROM opportunities WHERE id = CAST(:opportunity_id AS UUID) """, params) if not opp: return None tasks = preloaded.get("tasks") if tasks is None: tasks = _rows(conn, """ SELECT id::text, action_code, action, note, priority, route, status, due_at, created_at, metadata FROM tasks WHERE opportunity_id = CAST(:opportunity_id AS UUID) ORDER BY CASE COALESCE(priority, 'normal') WHEN 'alta' THEN 1 WHEN 'normal' THEN 2 WHEN 'baixa' THEN 3 ELSE 4 END, due_at NULLS LAST, created_at DESC LIMIT 20 """, params) # DTO contract retained from the legacy loader: # SELECT id::text, external_id, document_kind # The SQL now qualifies these fields because links and documents both # have ids; relationship comes exclusively from the canonical link. docs = preloaded.get("resolved_documents") if docs is None: from app.document_reconciliation_service import resolve_document_links docs = resolve_document_links(opportunity_id, conn=conn) docs = [row for row in docs if row.get("system") == "jasmin" and row.get("relationship") in {"PRIMARY", "SECONDARY", "HISTORICAL"}] linked_customer = preloaded.get("customer") customer_id = opp.get("fiscal_customer_id") or opp.get("customer_id") if customer_id and linked_customer is None: linked_customer = _first_row(conn, """ SELECT id::text, name, tax_id, email, email AS billing_email, street_name AS address, postal_zone AS postal_code, city_name AS city, phone FROM customers WHERE id = CAST(:customer_id AS UUID) """, {"customer_id": customer_id}) candidate = _first_row(conn, """ SELECT id::text, source_system, external_type, document_number, title, confidence FROM reconciliation_items WHERE opportunity_id = CAST(:opportunity_id AS UUID) AND status = 'open' ORDER BY confidence DESC NULLS LAST, created_at DESC LIMIT 1 """, params) snapshot = preloaded.get("operation_snapshot") or _operation_snapshot_safe(opportunity_id) fiscal_complete = False if linked_customer: fiscal_complete = bool( linked_customer.get("tax_id") and (linked_customer.get("billing_email") or linked_customer.get("email")) and linked_customer.get("address") and linked_customer.get("postal_code") and linked_customer.get("city") ) return build_opportunity_evidence( opp, linked_documents=docs, tasks=tasks, operation_snapshot=snapshot, linked_customer=linked_customer, fiscal_data_complete=fiscal_complete, has_reconciliation_candidate=bool(candidate), reconciliation_label=(candidate or {}).get("document_number") or (candidate or {}).get("title"), company_profile="blif", ) def get_opportunity_next_action( opportunity_id: str, *, preloaded: Optional[Dict[str, Any]] = None, ) -> Dict[str, Any]: """Return the recommended operator action for one opportunity. This remains a read-only service and returns the legacy dict shape, but the decision is now produced by the company workflow engine. """ if _flow_v2_mode() == "authoritative": return get_opportunity_next_actions([opportunity_id]).get(opportunity_id) or decide_authoritative_operation(None).to_dict() evidence = (_build_db_evidence(opportunity_id, preloaded=preloaded) if preloaded is not None else _build_db_evidence(opportunity_id)) if evidence is None: return OpportunityNextAction( action_code="NOT_FOUND", label="Oportunidade não encontrada", description="Não foi possível encontrar esta oportunidade.", priority="baixa", can_execute=False, reason_if_blocked="opportunity_not_found", ).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): return text(sql).bindparams(bindparam("opportunity_ids", expanding=True)) def _bulk_rows(conn: Any, sql: str, opportunity_ids: list[str]) -> list[Dict[str, Any]]: return [dict(row) for row in conn.execute( _bulk_statement(sql), {"opportunity_ids": opportunity_ids} ).mappings().all()] def _group_by_opportunity(rows: Iterable[Dict[str, Any]]) -> dict[str, list[Dict[str, Any]]]: grouped: dict[str, list[Dict[str, Any]]] = defaultdict(list) for row in rows: grouped[str(row.get("opportunity_id") or "")].append(row) return grouped def _resolve_bulk_documents( opportunity_ids: list[str], conn: Any ) -> dict[str, list[Dict[str, Any]]]: """Set-oriented equivalent of resolve_document_links(..., include_ended=False).""" from app.document_reconciliation_service import ( _LINK_SELECT, _legacy_relationship, document_reconciliation_v2_available, ) legacy_rows = _bulk_rows(conn, """ SELECT d.id::text, d.customer_id::text, d.opportunity_id::text, d.document_kind, d.external_id, d.document_number, d.status AS jasmin_status, d.status, d.total_amount, d.amount, d.tax_amount, d.currency, d.document_date, d.due_date, d.external_url, d.company, d.document_type, d.serie, d.series_number, d.customer_party_key, d.parent_document_id::text, d.version_number, d.payload, d.system, d.role, d.is_primary, d.is_active, d.created_at, d.updated_at FROM commercial_documents d WHERE d.opportunity_id::text IN :opportunity_ids ORDER BY d.opportunity_id, d.document_kind, d.created_at, d.id """, opportunity_ids) legacy = _group_by_opportunity(legacy_rows) if not document_reconciliation_v2_available(conn): return { oid: [row | {"document_id": row["id"], "link_id": None, "relationship": _legacy_relationship(row), "resolution_source": "legacy_rollout"} for row in legacy.get(oid, [])] for oid in opportunity_ids } v2_rows = [dict(row) for row in conn.execute( _bulk_statement(_LINK_SELECT + """ WHERE l.opportunity_id::text IN :opportunity_ids AND l.ended_at IS NULL ORDER BY l.opportunity_id, l.document_kind, l.updated_at DESC """), {"opportunity_ids": opportunity_ids}).mappings().all()] v2 = _group_by_opportunity(v2_rows) resolved: dict[str, list[Dict[str, Any]]] = {} for oid in opportunity_ids: legacy_groups = _group_by_opportunity( [row | {"opportunity_id": row.get("document_kind")} for row in legacy.get(oid, [])] ) v2_groups = _group_by_opportunity( [row | {"opportunity_id": row.get("document_kind")} for row in v2.get(oid, [])] ) documents: list[Dict[str, Any]] = [] for kind in sorted(set(legacy_groups) | set(v2_groups)): legacy_group, v2_group = legacy_groups.get(kind, []), v2_groups.get(kind, []) current_ids = {str(row["document_id"]) for row in v2_group if not row.get("ended_at")} legacy_ids = {str(row["id"]) for row in legacy_group} if legacy_ids and legacy_ids.issubset(current_ids): documents.extend(row | {"opportunity_id": oid, "resolution_source": "v2"} for row in v2_group) else: documents.extend(row | {"opportunity_id": oid, "document_id": row["id"], "link_id": None, "relationship": _legacy_relationship(row), "resolution_source": "legacy_rollout"} for row in legacy_group) resolved[oid] = documents return resolved def _bulk_operation_snapshots( opportunity_ids: list[str], links_by_opp: dict[str, list[Dict[str, Any]]], documents_by_opp: dict[str, list[Dict[str, Any]]], ) -> dict[str, Dict[str, Any]]: from app.document_reconciliation_service import select_valid_primary from app.operation_service import OPERATION_ACTIONS, OPERATION_CARDS, STATUS_LABELS, integration_settings_summary settings = integration_settings_summary() snapshots: dict[str, Dict[str, Any]] = {} for oid in opportunity_ids: links = links_by_opp.get(oid, []) by_key = {(row["system"], row["external_type"]): row for row in links} fallbacks: dict[tuple[str, str], Dict[str, Any]] = {} docs = documents_by_opp.get(oid, []) for kind in ("invoice", "quotation", "quote", "proforma"): row = select_valid_primary(docs, kind) if row is None: continue doc_kind = str(row.get("document_kind") or "").lower() key = ("jasmin", "invoice" if doc_kind == "invoice" else "quotation") if key in fallbacks: continue name = row.get("document_number") or row.get("external_id") or "Documento Jasmin" fallbacks[key] = { "id": "", "opportunity_id": oid, "system": key[0], "external_type": key[1], "external_id": row.get("external_id") or row.get("document_number") or "", "external_name": name, "external_url": row.get("external_url") or "", "status": "issued" if key[1] == "invoice" else "created", "payload": row.get("payload") or {"source": "commercial_documents_fallback"}, "last_synced_at": row.get("updated_at"), "created_at": row.get("updated_at"), "updated_at": row.get("updated_at"), } cards = [] for card in OPERATION_CARDS: link = by_key.get((card["system"], card["external_type"])) or fallbacks.get((card["system"], card["external_type"])) if link: status = link.get("status") or "pending" cards.append({**card, **link, "status_label": STATUS_LABELS.get(status, status)}) else: cards.append({**card, "status": "not_created", "status_label": card["empty"], "external_name": "", "external_url": ""}) snapshots[oid] = {"cards": cards, "links": links, "settings": settings, "actions": OPERATION_ACTIONS} return 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.""" ids = list(dict.fromkeys(str(value).strip() for value in opportunity_ids if str(value).strip())) if not ids: return {} if _flow_v2_mode() == "authoritative": return _get_authoritative_opportunity_decisions(ids) # 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: opportunities = _bulk_rows(conn, """ SELECT id::text, stage, status, title, local_customer_id::text AS fiscal_customer_id, local_customer_id::text AS customer_id, metadata FROM opportunities WHERE id::text IN :opportunity_ids """, ids) tasks = _bulk_rows(conn, """ SELECT * FROM ( SELECT id::text, opportunity_id::text, action_code, action, note, priority, route, status, due_at, created_at, metadata, row_number() OVER (PARTITION BY opportunity_id ORDER BY CASE COALESCE(priority, 'normal') WHEN 'alta' THEN 1 WHEN 'normal' THEN 2 WHEN 'baixa' THEN 3 ELSE 4 END, due_at NULLS LAST, created_at DESC) AS rn FROM tasks WHERE opportunity_id::text IN :opportunity_ids ) ranked WHERE rn <= 20 ORDER BY opportunity_id, rn """, ids) documents = _resolve_bulk_documents(ids, conn) customers = _bulk_rows(conn, """ SELECT o.id::text AS opportunity_id, c.id::text, c.name, c.tax_id, c.email, c.email AS billing_email, c.street_name AS address, c.postal_zone AS postal_code, c.city_name AS city, c.phone FROM opportunities o JOIN customers c ON c.id = o.local_customer_id WHERE o.id::text IN :opportunity_ids """, ids) candidates = _bulk_rows(conn, """ SELECT * FROM ( SELECT opportunity_id::text, id::text, source_system, external_type, document_number, title, confidence, row_number() OVER (PARTITION BY opportunity_id ORDER BY confidence DESC NULLS LAST, created_at DESC) AS rn FROM reconciliation_items WHERE opportunity_id::text IN :opportunity_ids AND status = 'open' ) ranked WHERE rn = 1 """, ids) operation_links = _bulk_rows(conn, """ SELECT id::text, opportunity_id::text, system, external_type, external_id, external_name, external_url, status, payload, last_synced_at, created_at, updated_at FROM operation_links WHERE opportunity_id::text IN :opportunity_ids ORDER BY opportunity_id, system, external_type """, ids) opp_by_id = {str(row["id"]): row for row in opportunities} tasks_by_id = _group_by_opportunity(tasks) customers_by_id = {str(row["opportunity_id"]): row for row in customers} candidates_by_id = {str(row["opportunity_id"]): row for row in candidates} links_by_id = _group_by_opportunity(operation_links) snapshots = _bulk_operation_snapshots(ids, links_by_id, documents) profile = load_company_profile("blif") decisions: Dict[str, Dict[str, Any]] = {} for oid in ids: opp = opp_by_id.get(oid) if opp is None: decisions[oid] = OpportunityNextAction( action_code="NOT_FOUND", label="Oportunidade não encontrada", description="Não foi possível encontrar esta oportunidade.", priority="baixa", can_execute=False, reason_if_blocked="opportunity_not_found", ).to_dict() continue customer = customers_by_id.get(oid) fiscal_complete = bool(customer and customer.get("tax_id") and (customer.get("billing_email") or customer.get("email")) and customer.get("address") and customer.get("postal_code") and customer.get("city")) candidate = candidates_by_id.get(oid) evidence = build_opportunity_evidence( opp, linked_documents=documents.get(oid, []), tasks=tasks_by_id.get(oid, []), operation_snapshot=snapshots.get(oid), linked_customer=customer, fiscal_data_complete=fiscal_complete, has_reconciliation_candidate=bool(candidate), reconciliation_label=(candidate or {}).get("document_number") or (candidate or {}).get("title"), company_profile="blif", ) decisions[oid] = decide_opportunity_next_action(evidence, profile).to_dict() _observe_flow_v2(decisions) return decisions def _get_authoritative_opportunity_decisions(ids: list[str]) -> Dict[str, Dict[str, Any]]: """Read-only set-oriented adapter input loader; never derives or persists V2.""" with engine.connect() as conn: projections = _bulk_rows(conn, """ 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 FROM opportunity_flow_state_v2 WHERE opportunity_id::text IN :opportunity_ids """, ids) tasks = _bulk_rows(conn, """ SELECT id::text, opportunity_id::text, action_code, status, due_at, resolved_at, superseded_by_task_id::text, created_at, metadata FROM tasks WHERE opportunity_id::text IN :opportunity_ids AND status = 'pending' """, ids) customers = _bulk_rows(conn, """ SELECT o.id::text AS opportunity_id, (c.tax_id IS NOT NULL AND c.tax_id <> '' AND c.email IS NOT NULL AND c.email <> '' AND c.street_name IS NOT NULL AND c.street_name <> '' AND c.postal_zone IS NOT NULL AND c.postal_zone <> '' AND c.city_name IS NOT NULL AND c.city_name <> '') AS fiscal_complete FROM opportunities o LEFT JOIN customers c ON c.id = o.local_customer_id WHERE o.id::text IN :opportunity_ids """, ids) reconciliation = _bulk_rows(conn, """ SELECT opportunity_id::text FROM reconciliation_items WHERE opportunity_id::text IN :opportunity_ids AND status IN ('needs_review','conflict') AND COALESCE((payload->>'reconciliation_blocks_current_action')::boolean, false) """, ids) projection_by_id = {str(row["opportunity_id"]): row for row in projections} tasks_by_id = _group_by_opportunity(tasks) fiscal_by_id = {str(row["opportunity_id"]): bool(row.get("fiscal_complete")) for row in customers} reconciliation_ids = {str(row["opportunity_id"]) for row in reconciliation} decisions = {} for oid in ids: projection = projection_by_id.get(oid) decision = decide_authoritative_operation( projection, obligations=tasks_by_id.get(oid, []), fiscal_complete=fiscal_by_id.get(oid, False), reconciliation_blocking=oid in reconciliation_ids, ) value = decision.to_dict() if projection is None: value["opportunity_id"] = oid value["target_url"] = f"/opportunities/{oid}" else: value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}" decisions[oid] = value return decisions