"""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 from typing import Any, Dict, Iterable, Optional from sqlalchemy import bindparam, text from app.db import engine # 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, ) @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 _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) -> OpportunityEvidence | None: params = {"opportunity_id": opportunity_id} with engine.begin() as conn: opp = _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 = _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. from app.document_reconciliation_service import resolve_document_links docs = [row for row in resolve_document_links(opportunity_id, conn=conn) if row.get("system") == "jasmin" and row.get("relationship") in {"PRIMARY", "SECONDARY", "HISTORICAL"}] linked_customer = None customer_id = opp.get("fiscal_customer_id") or opp.get("customer_id") if customer_id: 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 = _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) -> 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. """ evidence = _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)) return decision.to_dict() 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 {} from app.operation_service import ensure_operation_schema ensure_operation_schema() 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() return decisions