diff --git a/app/admin_ui/pages/opportunities.py b/app/admin_ui/pages/opportunities.py index 60ff0de..640b074 100644 --- a/app/admin_ui/pages/opportunities.py +++ b/app/admin_ui/pages/opportunities.py @@ -17,7 +17,7 @@ import app.admin_dashboard as legacy from app.admin_dashboard import * # noqa: F401,F403 from app.admin_ui.labels import primary_action_label from app.operation_noise import is_noise_operation_item -from app.opportunity_next_action_service import get_opportunity_next_action +from app.opportunity_next_action_service import get_opportunity_next_action, get_opportunity_next_actions from app.opportunity_action_task_materializer import ensure_pending_task_for_next_action from app.work_center_action_policy import ( canonical_action_code, @@ -1451,20 +1451,27 @@ def _attach_central_next_actions(opportunities: list[dict]) -> None: "Acompanhar produção". Keep this read-only and best-effort: if the central engine fails for one card, the card falls back to the legacy text. """ - for opp in opportunities: - if isinstance(opp.get("clientflow_next_action"), dict): - continue - oid = str(opp.get("id") or "").strip() - if not oid: - continue - try: - decision = get_opportunity_next_action(oid) - except Exception as exc: - decision = { + pending = { + str(opp.get("id") or "").strip(): opp + for opp in opportunities + if not isinstance(opp.get("clientflow_next_action"), dict) + and str(opp.get("id") or "").strip() + } + if not pending: + return + try: + decisions = get_opportunity_next_actions(pending) + except Exception as exc: + decisions = { + oid: { "action_code": "DECISION_ERROR", "label": opportunity_next_action_text(opp), "description": f"Falha ao calcular próxima ação central: {exc}", } + for oid, opp in pending.items() + } + for oid, opp in pending.items(): + decision = decisions.get(oid) if isinstance(decision, dict): opp["clientflow_next_action"] = decision diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index 6225e05..4010cb4 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -6,9 +6,10 @@ opportunity pages, tasks and future audits consume the same decision vocabulary. from __future__ import annotations from dataclasses import asdict, dataclass -from typing import Any, Dict, Optional +from collections import defaultdict +from typing import Any, Dict, Iterable, Optional -from sqlalchemy import text +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. @@ -163,3 +164,205 @@ def get_opportunity_next_action(opportunity_id: str) -> Dict[str, Any]: 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 diff --git a/tests/test_opportunity_next_action_bulk.py b/tests/test_opportunity_next_action_bulk.py new file mode 100644 index 0000000..babf0bc --- /dev/null +++ b/tests/test_opportunity_next_action_bulk.py @@ -0,0 +1,107 @@ +from contextlib import nullcontext + +import pytest + +import app.document_reconciliation_service as document_service +import app.operation_service as operation_service +import app.opportunity_next_action_service as service +from app.domain.opportunity_flow import build_opportunity_evidence + + +class _Result: + def __init__(self, rows): + self._rows = rows + + def mappings(self): + return self + + def all(self): + return self._rows + + +class _Connection: + def __init__(self, opportunities, *, tasks=None, operation_links=None): + self.opportunities = opportunities + self.tasks = tasks or [] + self.operation_links = operation_links or [] + self.query_count = 0 + + def execute(self, statement, params=None): + self.query_count += 1 + sql = str(statement) + if "FROM opportunities WHERE" in sql: + wanted = set((params or {}).get("opportunity_ids", [])) + return _Result([row for oid, row in self.opportunities.items() if oid in wanted]) + if "FROM tasks WHERE" in sql: + return _Result(self.tasks) + if "FROM operation_links WHERE" in sql: + return _Result(self.operation_links) + return _Result([]) + + +class _Engine: + def __init__(self, connection): + self.connection = connection + + def begin(self): + return nullcontext(self.connection) + + +def _opportunity(oid): + return { + "id": oid, + "stage": "NEW_LEAD", + "status": "open", + "title": "Teste", + "fiscal_customer_id": None, + "customer_id": None, + "metadata": {}, + } + + +@pytest.mark.parametrize("tasks,operation_links", [ + ([], []), + ([{"id": "task-1", "opportunity_id": "11111111-1111-1111-1111-111111111111", + "action_code": "CONTACT_CUSTOMER", "action": "Contactar cliente", "note": "", + "priority": "alta", "route": "/tasks/task-1", "status": "pending", + "due_at": None, "created_at": None, "metadata": {}, "rn": 1}], []), + ([], [{"id": "link-1", "opportunity_id": "11111111-1111-1111-1111-111111111111", + "system": "clientflow", "external_type": "payment", "external_id": None, + "external_name": "Pagamento", "external_url": None, "status": "confirmed", + "payload": {}, "last_synced_at": None, "created_at": None, "updated_at": None}]), +]) +def test_bulk_decision_is_equivalent_to_individual(monkeypatch, tasks, operation_links): + oid = "11111111-1111-1111-1111-111111111111" + row = _opportunity(oid) + connection = _Connection({oid: row}, tasks=tasks, operation_links=operation_links) + monkeypatch.setattr(service, "engine", _Engine(connection)) + monkeypatch.setattr(operation_service, "ensure_operation_schema", lambda: None) + monkeypatch.setattr(document_service, "document_reconciliation_v2_available", lambda conn=None: False) + + snapshots = service._bulk_operation_snapshots([oid], {oid: operation_links}, {oid: []}) + evidence = build_opportunity_evidence( + row, linked_documents=[], tasks=tasks, operation_snapshot=snapshots[oid], + linked_customer=None, fiscal_data_complete=False, + has_reconciliation_candidate=False, company_profile="blif", + ) + monkeypatch.setattr(service, "_build_db_evidence", lambda opportunity_id: evidence) + + assert service.get_opportunity_next_actions([oid])[oid] == service.get_opportunity_next_action(oid) + + +def test_bulk_query_count_is_constant_for_batch_size(monkeypatch): + monkeypatch.setattr(operation_service, "ensure_operation_schema", lambda: None) + monkeypatch.setattr(document_service, "document_reconciliation_v2_available", lambda conn=None: False) + + def measured(size): + opportunities = { + f"00000000-0000-0000-0000-{number:012d}": + _opportunity(f"00000000-0000-0000-0000-{number:012d}") + for number in range(1, size + 1) + } + connection = _Connection(opportunities) + monkeypatch.setattr(service, "engine", _Engine(connection)) + service.get_opportunity_next_actions(opportunities) + return connection.query_count + + assert measured(1) == measured(50) == 6 diff --git a/tests/test_v4928_1_5_108_pipeline_and_whout_indicators_static.py b/tests/test_v4928_1_5_108_pipeline_and_whout_indicators_static.py index 11da6b8..e4b12f7 100644 --- a/tests/test_v4928_1_5_108_pipeline_and_whout_indicators_static.py +++ b/tests/test_v4928_1_5_108_pipeline_and_whout_indicators_static.py @@ -4,7 +4,7 @@ from pathlib import Path def test_opportunity_board_uses_central_next_action(): text = Path('app/admin_ui/pages/opportunities.py').read_text() assert 'def _attach_central_next_actions' in text - assert 'get_opportunity_next_action(oid)' in text + assert 'get_opportunity_next_actions(pending)' in text assert '_attach_central_next_actions(opportunities)' in text assert 'central.get("label")' in text assert 'CLOSE_OPPORTUNITY' in text