Files
clientflow_backend/app/opportunity_next_action_service.py

369 lines
17 KiB
Python

"""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