perf: add opportunity detail read context

This commit is contained in:
plx
2026-08-14 13:24:37 +00:00
parent 07dc9db198
commit 7d138eccdc
13 changed files with 484 additions and 121 deletions

View File

@@ -1210,8 +1210,11 @@ def _outbox_items_for_opportunity(opportunity_id: str, *, target_system: str | N
if len(filtered) >= limit:
break
return filtered
def opportunity_integrations_panel_html(opportunity_id: str) -> str:
items = _outbox_items_for_opportunity(opportunity_id, target_system=None, limit=16)
def opportunity_integrations_panel_html(
opportunity_id: str, *, preloaded_items: Optional[list[dict]] = None,
) -> str:
items = (list(preloaded_items)[:16] if preloaded_items is not None
else _outbox_items_for_opportunity(opportunity_id, target_system=None, limit=16))
counts = {"pending": 0, "failed": 0, "blocked": 0, "dry_run": 0}
for item in items:
status = str(item.get("status") or "pending")
@@ -1270,8 +1273,13 @@ def opportunity_integrations_panel_html(opportunity_id: str) -> str:
</div>
</section>
'''
def opportunity_outbox_panel_html(opportunity_id: str, *, target_system: str = "jasmin") -> str:
items = _outbox_items_for_opportunity(opportunity_id, target_system=target_system, limit=12)
def opportunity_outbox_panel_html(
opportunity_id: str, *, target_system: str = "jasmin",
preloaded_items: Optional[list[dict]] = None,
) -> str:
items = ([item for item in preloaded_items if item.get("target_system") == target_system][:12]
if preloaded_items is not None
else _outbox_items_for_opportunity(opportunity_id, target_system=target_system, limit=12))
if not items:
return ""
rows = ""
@@ -1306,10 +1314,17 @@ def opportunity_outbox_panel_html(opportunity_id: str, *, target_system: str = "
</div>
</div>
'''
def jasmin_documents_html(opportunity_id: str, *, notice: str = "", error_notice: str = "") -> str:
def jasmin_documents_html(
opportunity_id: str, *, notice: str = "", error_notice: str = "",
preloaded_documents: Optional[list[dict]] = None,
preloaded_candidates: Optional[list[dict]] = None,
preloaded_customer: Optional[dict] = None,
preloaded_outbox: Optional[list[dict]] = None,
) -> str:
try:
from app.commercial_service import list_commercial_documents
docs = list_commercial_documents(opportunity_id=opportunity_id, limit=20)
docs = (list(preloaded_documents)[:20] if preloaded_documents is not None
else list_commercial_documents(opportunity_id=opportunity_id, limit=20))
except Exception as exc:
docs = []
error = str(exc)
@@ -1317,7 +1332,8 @@ def jasmin_documents_html(opportunity_id: str, *, notice: str = "", error_notice
error = ""
try:
from app.jasmin_backfill_service import find_jasmin_document_candidates_for_opportunity
jasmin_candidates = find_jasmin_document_candidates_for_opportunity(opportunity_id, limit=8)
jasmin_candidates = (list(preloaded_candidates)[:8] if preloaded_candidates is not None
else find_jasmin_document_candidates_for_opportunity(opportunity_id, limit=8))
except Exception as exc:
jasmin_candidates = []
if not error:
@@ -1325,7 +1341,8 @@ def jasmin_documents_html(opportunity_id: str, *, notice: str = "", error_notice
linked_tax_id = ""
try:
from app.commercial_service import get_customer_for_opportunity, normalize_tax_id
linked_customer = get_customer_for_opportunity(opportunity_id)
linked_customer = (preloaded_customer if preloaded_customer is not None
else get_customer_for_opportunity(opportunity_id))
linked_tax_id = normalize_tax_id((linked_customer or {}).get("tax_id"))
except Exception:
linked_tax_id = ""
@@ -1539,7 +1556,9 @@ def jasmin_documents_html(opportunity_id: str, *, notice: str = "", error_notice
notice_html = f'<div class="alert alert-info py-2 small mb-2">{esc(notice)}</div>' if notice else ''
error_notice_html = f'<div class="alert alert-danger py-2 small mb-2"><strong>Não foi possível pedir a ação.</strong><br>{esc(error_notice).replace(chr(10), "<br>")}</div>' if error_notice else ''
error_html = f'<div class="alert alert-warning py-2 small mb-2">{esc(error)}</div>' if error else ''
outbox_html = opportunity_outbox_panel_html(opportunity_id, target_system="jasmin")
outbox_html = opportunity_outbox_panel_html(
opportunity_id, target_system="jasmin", preloaded_items=preloaded_outbox,
)
refreshed_at = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC")
if invoice_source_exists:
convert_invoice_button_html = (
@@ -1762,10 +1781,16 @@ def opportunity_add_item_form_html(opportunity_id: str, products: list[dict]) ->
</div>
</form>
'''
def opportunity_products_panel_html(opportunity_id: str, *, notice: str = "", error_notice: str = "") -> str:
def opportunity_products_panel_html(
opportunity_id: str, *, notice: str = "", error_notice: str = "",
preloaded_items: Optional[list[dict]] = None,
preloaded_products: Optional[list[dict]] = None,
) -> str:
try:
items = list_opportunity_items(opportunity_id)
active_products = list_products(active="true", limit=300)
items = (list(preloaded_items) if preloaded_items is not None
else list_opportunity_items(opportunity_id))
active_products = (list(preloaded_products) if preloaded_products is not None
else list_products(active="true", limit=300))
except Exception as exc:
return f'<section id="opportunity-products-panel" class="card cf-card"><div class="card-body"><div class="alert alert-danger">Erro ao carregar produtos: {esc(exc)}</div></div></section>'
total = sum(

View File

@@ -12,13 +12,14 @@ from urllib.parse import quote
import json
import time
import uuid
from dataclasses import dataclass, field
from typing import Optional
from datetime import datetime, timezone
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, 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,
reconstructed_review_required,
@@ -48,6 +49,119 @@ _opportunity_board_column_for_stage = legacy._opportunity_board_column_for_stage
router = APIRouter()
@dataclass
class OpportunityDetailContext:
"""Request-scoped, read-only data shared by opportunity-detail renderers."""
opportunity_id: str
opportunity: dict
customer: Optional[dict] = None
tasks: list[dict] = field(default_factory=list)
events: list[dict] = field(default_factory=list)
opportunity_items: list[dict] = field(default_factory=list)
active_products: list[dict] = field(default_factory=list)
resolved_documents: list[dict] = field(default_factory=list)
operation_links: list[dict] = field(default_factory=list)
operation_snapshot: dict = field(default_factory=dict)
communications: list[dict] = field(default_factory=list)
jasmin_candidates: list[dict] = field(default_factory=list)
odoo_candidates: list[dict] = field(default_factory=list)
outbox_items: list[dict] = field(default_factory=list)
def next_action_preloaded(self) -> dict:
evidence_opportunity = dict(self.opportunity)
evidence_opportunity["fiscal_customer_id"] = (
evidence_opportunity.get("local_customer_id")
or evidence_opportunity.get("linked_customer_id")
)
evidence_opportunity["customer_id"] = evidence_opportunity.get("fiscal_customer_id")
evidence_customer = None
if self.customer:
evidence_customer = {
**self.customer,
"billing_email": self.customer.get("email"),
"address": self.customer.get("street_name"),
"postal_code": self.customer.get("postal_zone"),
"city": self.customer.get("city_name"),
}
return {
"opportunity": evidence_opportunity,
"customer": evidence_customer,
"resolved_documents": self.resolved_documents,
"operation_snapshot": self.operation_snapshot,
}
def load_opportunity_detail_context(opportunity_id: str) -> Optional[OpportunityDetailContext]:
"""Load each reusable opportunity-detail dataset at most once."""
opportunity = get_opportunity(opportunity_id)
if not opportunity:
return None
from app.commercial_service import get_customer_for_opportunity, list_commercial_documents
from app.integration_outbox_service import list_outbox
from app.operation_service import get_operation_links, get_operation_snapshot
from app.jasmin_backfill_service import find_jasmin_document_candidates_for_opportunity
customer = get_customer_for_opportunity(opportunity_id)
tasks = list_opportunity_tasks(opportunity_id, limit=100)
events = list_opportunity_events(opportunity_id, limit=100)
items = list_opportunity_items(opportunity_id)
products = list_products(active="true", limit=300)
documents = list_commercial_documents(opportunity_id=opportunity_id, limit=100)
operation_links = get_operation_links(opportunity_id)
snapshot = get_operation_snapshot(
opportunity_id, links=operation_links, resolved_documents=documents,
)
try:
communications = list_communications_for_opportunity(opportunity_id, limit=12)
except Exception:
communications = []
try:
outbox_all = list_outbox(limit=300)
outbox_items = []
for item in outbox_all:
payload = item.get("payload") or {}
if isinstance(payload, str):
try:
payload = json.loads(payload)
except Exception:
payload = {}
if str(payload.get("opportunity_id") or "") == opportunity_id:
data = dict(item)
data["payload"] = payload
outbox_items.append(data)
except Exception:
outbox_items = []
candidate_opportunity = {
**opportunity,
"fiscal_customer_name": (customer or {}).get("name"),
"fiscal_customer_tax_id": (customer or {}).get("tax_id"),
"fiscal_customer_email": (customer or {}).get("email"),
}
try:
jasmin_candidates = find_jasmin_document_candidates_for_opportunity(
opportunity_id, limit=8, opportunity=candidate_opportunity,
current_documents=documents,
)
except Exception:
jasmin_candidates = []
odoo_links = [row for row in operation_links if row.get("system") == "odoo"]
try:
_links, odoo_candidates = _opportunity_odoo_rows(
opportunity_id, preloaded_links=odoo_links,
)
except Exception:
odoo_candidates = []
return OpportunityDetailContext(
opportunity_id=opportunity_id, opportunity=opportunity, customer=customer,
tasks=tasks, events=events, opportunity_items=items, active_products=products,
resolved_documents=documents, operation_links=operation_links,
operation_snapshot=snapshot, communications=communications,
jasmin_candidates=jasmin_candidates, odoo_candidates=odoo_candidates,
outbox_items=outbox_items,
)
def _authenticated_actor(request: Request) -> str:
actor = str(getattr(request.state, "clientflow_admin_user", "") or "").strip()
if not actor:
@@ -838,14 +952,16 @@ def _odoo_status_badge(status: object) -> str:
return f'<span class="badge {cls}">{esc(status or "—")}</span>'
def _opportunity_odoo_rows(opportunity_id: str) -> tuple[list[dict], list[dict]]:
def _opportunity_odoo_rows(
opportunity_id: str, *, preloaded_links: Optional[list[dict]] = None,
) -> tuple[list[dict], list[dict]]:
"""Return linked Odoo operation links and recent reconciliation candidates.
Read-only. The panel must not call Odoo on page load; the operator uses
the explicit sync button to refresh live Odoo state.
"""
with engine.begin() as conn:
links = conn.execute(text("""
links = (preloaded_links if preloaded_links is not None else conn.execute(text("""
SELECT id::text, system, external_type, external_id, external_name,
external_url, status, payload, last_synced_at, updated_at
FROM operation_links
@@ -860,7 +976,7 @@ def _opportunity_odoo_rows(opportunity_id: str) -> tuple[list[dict], list[dict]]
ELSE 9
END,
updated_at DESC
"""), {"opportunity_id": opportunity_id}).mappings().all()
"""), {"opportunity_id": opportunity_id}).mappings().all())
candidates = conn.execute(text("""
SELECT id::text, source_system, external_type, external_id,
document_number, title, status, amount, currency,
@@ -882,9 +998,17 @@ def _opportunity_odoo_rows(opportunity_id: str) -> tuple[list[dict], list[dict]]
return [dict(r) for r in links], [dict(r) for r in candidates]
def odoo_status_panel_html(opportunity_id: str, *, notice: str = "", error_notice: str = "") -> str:
def odoo_status_panel_html(
opportunity_id: str, *, notice: str = "", error_notice: str = "",
preloaded_links: Optional[list[dict]] = None,
preloaded_candidates: Optional[list[dict]] = None,
) -> str:
try:
links, candidates = _opportunity_odoo_rows(opportunity_id)
if preloaded_candidates is None:
links, candidates = _opportunity_odoo_rows(opportunity_id, preloaded_links=preloaded_links)
else:
links = list(preloaded_links or [])
candidates = list(preloaded_candidates)
except Exception as exc:
return f'<section id="odoo-status-panel" class="card cf-card"><div class="card-body"><div class="alert alert-danger">Erro ao carregar Odoo: {esc(exc)}</div></div></section>'
@@ -1059,10 +1183,14 @@ def _task_href_with_return_to(task_id: str, return_to: str) -> str:
return href
def _render_email_identity_review(opportunity_id: str, linked_customer: dict | None) -> str:
def _render_email_identity_review(
opportunity_id: str, linked_customer: dict | None, *, opportunity: Optional[dict] = None,
) -> str:
try:
from app.fiscal_enrichment_service import email_identity_review_for_opportunity
review = email_identity_review_for_opportunity(opportunity_id, refresh=False)
review = email_identity_review_for_opportunity(
opportunity_id, refresh=False, opportunity=opportunity,
)
except Exception as exc:
return f"""
<div class="cf-soft-box mt-3">
@@ -1307,9 +1435,22 @@ def _opportunity_jasmin_state(opportunity_id: str) -> dict:
return {"jasmin_documents": 0, "item_count": 0}
def _opportunity_consistency_alert_html(opportunity: dict, tasks: list[dict], opportunity_items: list[dict], opportunity_id: str) -> str:
def _opportunity_consistency_alert_html(
opportunity: dict, tasks: list[dict], opportunity_items: list[dict], opportunity_id: str,
*, preloaded_documents: Optional[list[dict]] = None,
) -> str:
# Surface soft inconsistencies without blocking the operator.
state = _opportunity_jasmin_state(opportunity_id)
if preloaded_documents is None:
state = _opportunity_jasmin_state(opportunity_id)
else:
jasmin_docs = [row for row in preloaded_documents if row.get("system") == "jasmin"]
by_kind = {kind: sum(row.get("document_kind") == kind for row in jasmin_docs)
for kind in ("quotation", "proforma", "invoice")}
state = {
"jasmin_documents": len(jasmin_docs), "quotations": by_kind["quotation"],
"proformas": by_kind["proforma"], "invoices": by_kind["invoice"],
"item_count": len(opportunity_items),
}
stage = str(opportunity.get("stage") or "")
pending_action_codes = {str(t.get("action_code") or "") for t in tasks if str(t.get("status") or "") == "pending"}
has_payment_task = bool({"CONFIRM_PAYMENT", "CONFIRM_PAYMENT_AND_PREPARE_SHIPMENT"} & pending_action_codes)
@@ -1794,12 +1935,12 @@ async def opportunities_page(
async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = None):
if not is_uuid_text(opportunity_id):
return PlainTextResponse("Identificador de oportunidade inválido.", status_code=422)
opportunity = get_opportunity(opportunity_id)
if not opportunity:
context = load_opportunity_detail_context(opportunity_id)
if context is None:
return layout("Oportunidade não encontrada", "Pipeline comercial", '<section class="cf-empty">Oportunidade não encontrada.</section>', "opportunities")
tasks = list_opportunity_tasks(opportunity_id, limit=100)
events = list_opportunity_events(opportunity_id, limit=100)
opportunity = context.opportunity
tasks = context.tasks
events = context.events
stage = str(opportunity.get("stage") or "NEW_LEAD")
terminal_stage = str(opportunity.get("status") or "").lower() == "closed" or stage in {"WON", "LOST", "NO_INTEREST", "DELIVERED"}
all_pending_tasks = [t for t in tasks if str(t.get("status")) == "pending"]
@@ -1821,13 +1962,9 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
else:
pending_tasks = [t for t in all_pending_tasks if not _is_obsolete_after_payment_task(t, payment_confirmed_for_ui)]
next_task = pending_tasks[0] if pending_tasks else None
opportunity_items = list_opportunity_items(opportunity_id)
active_products = list_products(active="true", limit=200)
try:
from app.commercial_service import list_commercial_documents
linked_documents = list_commercial_documents(opportunity_id=opportunity_id, limit=8)
except Exception:
linked_documents = []
opportunity_items = context.opportunity_items
active_products = context.active_products[:200]
linked_documents = context.resolved_documents[:8]
primary_document = next(
(
doc for doc in linked_documents
@@ -1854,13 +1991,10 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
value_source = "documento principal" if document_value else (
"valor manual/operacional aprovado" if opportunity_items_total else "documento principal por definir"
)
operation_snapshot = get_operation_snapshot(opportunity_id)
operation_snapshot = context.operation_snapshot
opportunity_for_cockpit = dict(opportunity)
opportunity_for_cockpit["pending_task_count"] = len(pending_tasks)
try:
opportunity_communications = list_communications_for_opportunity(opportunity_id, limit=12)
except Exception:
opportunity_communications = []
opportunity_communications = context.communications
notice_html = f'<div class="alert alert-info">{esc(notice)}</div>' if notice else ''
metadata = _opportunity_metadata(opportunity)
payment_term = str(metadata.get("payment_terms") or "before_shipping")
@@ -1897,7 +2031,9 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
opportunity_return_to = f"/opportunities/{opportunity_id}"
try:
next_action = get_opportunity_next_action(opportunity_id)
next_action = get_opportunity_next_action(
opportunity_id, preloaded=context.next_action_preloaded(),
)
except Exception:
next_action = {}
@@ -1942,41 +2078,8 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
"source": "pending_task",
}
# Materialize human-only central actions into actual pending tasks.
# The top-level next action should not be an abstract label when the workbench
# expects an operator to perform it. The helper is idempotent and currently
# creates SEND_INVOICE/FOLLOW_UP_PAYMENT tasks when needed.
try:
materialized_task = ensure_pending_task_for_next_action(
opportunity_id,
next_action if isinstance(next_action, dict) else {},
source="opportunity_detail",
actor="system",
)
except Exception:
materialized_task = {"created": False}
if materialized_task.get("created"):
tasks = list_opportunity_tasks(opportunity_id, limit=100)
all_pending_tasks = [t for t in tasks if str(t.get("status")) == "pending"]
if terminal_stage:
pending_tasks = [
t for t in all_pending_tasks
if str(t.get("action_code") or "").upper() not in {
"FOLLOW_UP_QUOTE",
"FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT",
"FOLLOW_UP_CUSTOMER_REVIEW",
"FOLLOW_UP_GENERIC",
}
]
else:
pending_tasks = [t for t in all_pending_tasks if not _is_obsolete_after_payment_task(t, payment_confirmed_for_ui)]
next_task = pending_tasks[0] if pending_tasks else None
opportunity_for_cockpit["pending_task_count"] = len(pending_tasks)
try:
next_action = get_opportunity_next_action(opportunity_id)
except Exception:
pass
# GET detail is strictly read-only. SEND_INVOICE task materialization remains
# available to explicit write/repair workflows, never during page rendering.
if isinstance(next_action, dict):
# v1.5.107: keep the legacy operational cockpit aligned with the
# central decision engine. Without this, CLOSE_OPPORTUNITY could show
@@ -2090,7 +2193,7 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
customer_options = '<option value="">Selecionar cliente...</option>'
try:
from app.commercial_service import get_customer_for_opportunity, list_customers
linked_customer = get_customer_for_opportunity(opportunity_id)
linked_customer = context.customer
for c in list_customers(limit=150):
selected = "selected" if linked_customer and str(c.get("id")) == str(linked_customer.get("id")) else ""
label = f"{c.get('name') or 'Cliente'} · {c.get('tax_id') or 'sem NIF'}"
@@ -2099,7 +2202,9 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
linked_customer = None
fiscal_suggestions_html = _render_fiscal_suggestions(opportunity_id, linked_customer)
email_identity_html = _render_email_identity_review(opportunity_id, linked_customer)
email_identity_html = _render_email_identity_review(
opportunity_id, linked_customer, opportunity=opportunity,
)
try:
from app.jasmin_fiscal_sync_service import get_jasmin_fiscal_sync_preview
jasmin_fiscal_preview = get_jasmin_fiscal_sync_preview(opportunity_id)
@@ -2144,7 +2249,10 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
ok_text="Dados mínimos de envio completos.",
blocked_text="Envio deve aguardar correção destes dados.",
)
consistency_alert_html = _opportunity_consistency_alert_html(opportunity, tasks, opportunity_items, opportunity_id)
consistency_alert_html = _opportunity_consistency_alert_html(
opportunity, tasks, opportunity_items, opportunity_id,
preloaded_documents=context.resolved_documents,
)
if primary_document:
document_label = commercial_document_display_number(primary_document, fallback="número por atualizar")
document_kind = {
@@ -2589,13 +2697,26 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
{operation_cockpit_html(opportunity_id, opportunity_for_cockpit, operation_snapshot)}
<div id="documentos">{jasmin_documents_html(opportunity_id)}</div>
<div id="documentos">{jasmin_documents_html(
opportunity_id, preloaded_documents=context.resolved_documents,
preloaded_candidates=context.jasmin_candidates,
preloaded_customer=context.customer, preloaded_outbox=context.outbox_items,
)}</div>
<div id="produtos">{opportunity_products_panel_html(opportunity_id)}</div>
<div id="produtos">{opportunity_products_panel_html(
opportunity_id, preloaded_items=context.opportunity_items,
preloaded_products=context.active_products,
)}</div>
<div id="odoo">{odoo_status_panel_html(opportunity_id)}</div>
<div id="odoo">{odoo_status_panel_html(
opportunity_id,
preloaded_links=[row for row in context.operation_links if row.get("system") == "odoo"],
preloaded_candidates=context.odoo_candidates,
)}</div>
<div id="outbox">{opportunity_integrations_panel_html(opportunity_id)}</div>
<div id="outbox">{opportunity_integrations_panel_html(
opportunity_id, preloaded_items=context.outbox_items,
)}</div>
<section id="tasks-list" class="card cf-card"><div class="card-body p-0"><div class="p-3 border-bottom"><h2 class="cf-section-title">Tasks relacionadas</h2><div class="small text-secondary">Ações humanas já criadas para esta oportunidade.</div></div><div class="cf-table-wrap border-0 rounded-0"><table class="table cf-table"><thead><tr><th>Ação</th><th>Fila</th><th>Estado</th><th>Data</th></tr></thead><tbody>{task_rows}</tbody></table></div></div></section>

View File

@@ -718,7 +718,6 @@ def list_customers(q: Optional[str] = None, limit: int = 200) -> List[Dict[str,
A página Customers passa a ser a fonte principal dos dados fiscais.
"""
ensure_commercial_schema()
params: Dict[str, Any] = {"limit": int(limit)}
where = ""
if q:

View File

@@ -776,7 +776,6 @@ def _normalize_stored_identity(identity: Optional[Dict[str, Any]]) -> Optional[D
def latest_identity_for_opportunity(opportunity_id: str) -> Optional[Dict[str, Any]]:
ensure_email_identity_schema()
with engine.begin() as conn:
row = conn.execute(text("""
SELECT id::text, opportunity_id::text, task_id::text, message_id::text,

View File

@@ -36,7 +36,11 @@ from app.commercial_service import (
upsert_customer,
)
from app.opportunity_service import ensure_opportunity_schema, get_opportunity, list_opportunities
from app.email_identity_extraction_service import extract_identity_for_opportunity, is_plausible_company_mention
from app.email_identity_extraction_service import (
extract_identity_for_opportunity,
is_plausible_company_mention,
latest_identity_for_opportunity,
)
CACHE_SOURCE = "informa_pipeline_api"
@@ -823,7 +827,10 @@ def _identity_conflicts_with_linked_customer(identity: Optional[Dict[str, Any]],
return not _identity_mentions_match_name(identity, linked_name)
def email_identity_review_for_opportunity(opportunity_id: str, *, refresh: bool = False) -> Dict[str, Any]:
def email_identity_review_for_opportunity(
opportunity_id: str, *, refresh: bool = False,
opportunity: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
"""Return the operator-facing identity review for an opportunity.
This is read-only except when refresh=True, where it stores a new extraction.
@@ -831,16 +838,18 @@ def email_identity_review_for_opportunity(opportunity_id: str, *, refresh: bool
can show cases like: email mentions Dietimport S.A. but the opportunity is
linked to Fmrl - Imobiliária S.A.
"""
ensure_fiscal_enrichment_schema()
opportunity = get_opportunity(opportunity_id)
opportunity = opportunity or get_opportunity(opportunity_id)
if not opportunity:
return {"ok": False, "reason": "opportunity_not_found"}
try:
identity = extract_identity_for_opportunity(
opportunity_id,
refresh=refresh,
use_llm=bool(getattr(settings, "email_identity_extraction_use_llm", True)),
)
if refresh:
identity = extract_identity_for_opportunity(
opportunity_id,
refresh=True,
use_llm=bool(getattr(settings, "email_identity_extraction_use_llm", True)),
)
else:
identity = latest_identity_for_opportunity(opportunity_id)
except Exception as exc:
return {"ok": False, "reason": f"email_identity_extraction_failed: {exc}"}
@@ -1266,7 +1275,6 @@ def reject_fiscal_suggestion(suggestion_id: str, *, actor: str = "operator", rea
def list_fiscal_suggestions_for_opportunity(opportunity_id: str, *, limit: int = 5) -> List[Dict[str, Any]]:
ensure_fiscal_enrichment_schema()
with engine.begin() as conn:
rows = conn.execute(text("""
SELECT

View File

@@ -437,7 +437,11 @@ def _email_domain(email: Any) -> str:
return text_value.rsplit("@", 1)[-1].strip()
def find_jasmin_document_candidates_for_opportunity(opportunity_id: str, *, limit: int = 5) -> List[Dict[str, Any]]:
def find_jasmin_document_candidates_for_opportunity(
opportunity_id: str, *, limit: int = 5,
opportunity: Optional[Dict[str, Any]] = None,
current_documents: Optional[List[Dict[str, Any]]] = None,
) -> List[Dict[str, Any]]:
"""Find Jasmin documents that probably belong to an opportunity but are not imported yet.
This is deliberately conservative and prioritizes fiscal identity (NIF/customer_id)
@@ -446,7 +450,7 @@ def find_jasmin_document_candidates_for_opportunity(opportunity_id: str, *, limi
quotation instead of duplicating it.
"""
with engine.begin() as conn:
opp = conn.execute(text("""
opp = opportunity or conn.execute(text("""
SELECT
o.id::text,
o.customer_name,
@@ -470,7 +474,8 @@ def find_jasmin_document_candidates_for_opportunity(opportunity_id: str, *, limi
raw_domain = _email_domain(fiscal_email)
domain = "" if is_public_email_domain(raw_domain) else raw_domain
current_docs = conn.execute(text("""
if current_documents is None:
current_docs = conn.execute(text("""
SELECT
id::text, document_number, external_id, document_kind, document_date,
total_amount, amount, payload, created_at, updated_at
@@ -479,7 +484,13 @@ def find_jasmin_document_candidates_for_opportunity(opportunity_id: str, *, limi
AND system = 'jasmin'
AND document_kind IN ('quotation', 'proforma')
ORDER BY document_date DESC NULLS LAST, created_at DESC
"""), {"opportunity_id": opportunity_id}).mappings().all()
"""), {"opportunity_id": opportunity_id}).mappings().all()
else:
current_docs = [
row for row in current_documents
if str(row.get("system") or "") == "jasmin"
and str(row.get("document_kind") or "") in {"quotation", "proforma"}
]
candidate_columns = """
ri.id::text,

View File

@@ -182,8 +182,6 @@ def _select_candidate_doc(conn: Any, opportunity_id: str) -> Optional[Dict[str,
def get_jasmin_fiscal_sync_preview(opportunity_id: str) -> Dict[str, Any]:
"""Read-only preview for fiscal data available in linked Jasmin documents."""
ensure_opportunity_schema()
ensure_commercial_schema()
with engine.begin() as conn:
opp = conn.execute(text("""
SELECT o.id::text, o.local_customer_id::text, o.title,

View File

@@ -86,7 +86,9 @@ def get_operation_links(opportunity_id: str) -> List[Dict[str, Any]]:
"""), {"opportunity_id": opportunity_id}).mappings().all()
return [dict(r) for r in rows]
def _commercial_document_operation_fallbacks(opportunity_id: str) -> Dict[tuple, Dict[str, Any]]:
def _commercial_document_operation_fallbacks(
opportunity_id: str, resolved_documents: Optional[List[Dict[str, Any]]] = None,
) -> Dict[tuple, Dict[str, Any]]:
"""Infer operation cards from linked commercial_documents when operation_links lag.
Reconstructed opportunities often have commercial_documents imported from
@@ -95,7 +97,7 @@ def _commercial_document_operation_fallbacks(opportunity_id: str) -> Dict[tuple,
"""
try:
from app.document_reconciliation_service import resolve_document_links, select_valid_primary
resolved = resolve_document_links(opportunity_id)
resolved = resolved_documents if resolved_documents is not None else resolve_document_links(opportunity_id)
rows = [row for kind in ("invoice", "quotation", "quote", "proforma")
if (row := select_valid_primary(resolved, kind)) is not None]
except Exception:
@@ -132,10 +134,13 @@ def _commercial_document_operation_fallbacks(opportunity_id: str) -> Dict[tuple,
return fallbacks
def get_operation_snapshot(opportunity_id: str) -> Dict[str, Any]:
links = get_operation_links(opportunity_id)
def get_operation_snapshot(
opportunity_id: str, *, links: Optional[List[Dict[str, Any]]] = None,
resolved_documents: Optional[List[Dict[str, Any]]] = None,
) -> Dict[str, Any]:
links = list(links) if links is not None else get_operation_links(opportunity_id)
by_key = {(x["system"], x["external_type"]): x for x in links}
commercial_fallbacks = _commercial_document_operation_fallbacks(opportunity_id)
commercial_fallbacks = _commercial_document_operation_fallbacks(opportunity_id, resolved_documents)
cards = []
for card in OPERATION_CARDS:
link = by_key.get((card["system"], card["external_type"])) or commercial_fallbacks.get((card["system"], card["external_type"]))

View File

@@ -55,10 +55,13 @@ def _operation_snapshot_safe(opportunity_id: str) -> dict[str, Any]:
return {"cards": [], "links": []}
def _build_db_evidence(opportunity_id: str) -> OpportunityEvidence | None:
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 = _first_row(conn, """
opp = preloaded.get("opportunity") or _first_row(conn, """
SELECT
id::text,
stage,
@@ -73,7 +76,9 @@ def _build_db_evidence(opportunity_id: str) -> OpportunityEvidence | None:
if not opp:
return None
tasks = _rows(conn, """
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)
@@ -88,14 +93,16 @@ def _build_db_evidence(opportunity_id: str) -> OpportunityEvidence | None:
# 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
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 = None
linked_customer = preloaded.get("customer")
customer_id = opp.get("fiscal_customer_id") or opp.get("customer_id")
if customer_id:
if customer_id and linked_customer is None:
linked_customer = _first_row(conn, """
SELECT
id::text,
@@ -120,7 +127,7 @@ def _build_db_evidence(opportunity_id: str) -> OpportunityEvidence | None:
LIMIT 1
""", params)
snapshot = _operation_snapshot_safe(opportunity_id)
snapshot = preloaded.get("operation_snapshot") or _operation_snapshot_safe(opportunity_id)
fiscal_complete = False
if linked_customer:
fiscal_complete = bool(
@@ -144,14 +151,17 @@ def _build_db_evidence(opportunity_id: str) -> OpportunityEvidence | None:
)
def get_opportunity_next_action(opportunity_id: str) -> Dict[str, Any]:
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.
"""
evidence = _build_db_evidence(opportunity_id)
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",