2 Commits

Author SHA1 Message Date
plx
91ba85fc92 perf: optimize revenue forecast effective document links 2026-08-14 12:06:26 +00:00
plx
69116f620c perf: eliminate N+1 in opportunities next actions 2026-08-14 11:56:03 +00:00
7 changed files with 519 additions and 38 deletions

View File

@@ -17,7 +17,7 @@ import app.admin_dashboard as legacy
from app.admin_dashboard import * # noqa: F401,F403 from app.admin_dashboard import * # noqa: F401,F403
from app.admin_ui.labels import primary_action_label from app.admin_ui.labels import primary_action_label
from app.operation_noise import is_noise_operation_item 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.opportunity_action_task_materializer import ensure_pending_task_for_next_action
from app.work_center_action_policy import ( from app.work_center_action_policy import (
canonical_action_code, 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 "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. central engine fails for one card, the card falls back to the legacy text.
""" """
for opp in opportunities: pending = {
if isinstance(opp.get("clientflow_next_action"), dict): str(opp.get("id") or "").strip(): opp
continue for opp in opportunities
oid = str(opp.get("id") or "").strip() if not isinstance(opp.get("clientflow_next_action"), dict)
if not oid: and str(opp.get("id") or "").strip()
continue }
try: if not pending:
decision = get_opportunity_next_action(oid) return
except Exception as exc: try:
decision = { decisions = get_opportunity_next_actions(pending)
except Exception as exc:
decisions = {
oid: {
"action_code": "DECISION_ERROR", "action_code": "DECISION_ERROR",
"label": opportunity_next_action_text(opp), "label": opportunity_next_action_text(opp),
"description": f"Falha ao calcular próxima ação central: {exc}", "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): if isinstance(decision, dict):
opp["clientflow_next_action"] = decision opp["clientflow_next_action"] = decision

View File

@@ -181,28 +181,78 @@ def prepare_effective_document_links(conn: Any) -> str:
result sets. The temporary table is session-local and is populated solely result sets. The temporary table is session-local and is populated solely
through the group-aware resolver. through the group-aware resolver.
""" """
opportunity_ids = {str(row[0]) for row in conn.execute(text("""
SELECT DISTINCT opportunity_id::text FROM commercial_documents
WHERE opportunity_id IS NOT NULL
""")).all()}
if document_reconciliation_v2_available(conn):
opportunity_ids.update(str(row[0]) for row in conn.execute(text("""
SELECT DISTINCT opportunity_id::text FROM opportunity_document_links
WHERE ended_at IS NULL
""")).all())
conn.execute(text("""CREATE TEMP TABLE IF NOT EXISTS _effective_document_links ( conn.execute(text("""CREATE TEMP TABLE IF NOT EXISTS _effective_document_links (
opportunity_id UUID NOT NULL, document_id UUID NOT NULL, document_kind TEXT NOT NULL, opportunity_id UUID NOT NULL, document_id UUID NOT NULL, document_kind TEXT NOT NULL,
relationship TEXT NOT NULL, ended_at TIMESTAMPTZ) ON COMMIT DROP""")) relationship TEXT NOT NULL, ended_at TIMESTAMPTZ) ON COMMIT DROP"""))
conn.execute(text("TRUNCATE _effective_document_links")) conn.execute(text("TRUNCATE _effective_document_links"))
for opportunity_id in sorted(opportunity_ids): legacy_relationship = """CASE
for row in resolve_document_links(opportunity_id, conn=conn): WHEN d.is_active IS FALSE OR lower(COALESCE(NULLIF(d.role, ''), 'current')) = 'detached' THEN 'REMOVED'
conn.execute(text("""INSERT INTO _effective_document_links WHEN d.is_primary IS TRUE AND lower(COALESCE(NULLIF(d.role, ''), 'current')) IN ('current','accepted') THEN 'PRIMARY'
(opportunity_id,document_id,document_kind,relationship,ended_at) WHEN lower(COALESCE(NULLIF(d.role, ''), 'current')) IN ('historical','history','superseded') THEN 'HISTORICAL'
VALUES(CAST(:oid AS UUID),CAST(:did AS UUID),:kind,:relationship,:ended_at)"""), ELSE 'SECONDARY'
{"oid": opportunity_id, "did": row["document_id"], END"""
"kind": row.get("document_kind") or "unknown", if document_reconciliation_v2_available(conn):
"relationship": row.get("relationship") or "SECONDARY", # This is the relational form of resolve_document_links(): a legacy
"ended_at": row.get("ended_at")}) # group switches to v2 only when every legacy document in that exact
# opportunity/kind group has a current v2 link. V2-only groups remain
# absent, matching the rollout resolver's current behaviour.
conn.execute(text(f"""
INSERT INTO _effective_document_links
(opportunity_id, document_id, document_kind, relationship, ended_at)
WITH legacy_groups AS (
SELECT DISTINCT d.opportunity_id, COALESCE(d.document_kind, '') AS document_kind
FROM commercial_documents d
WHERE d.opportunity_id IS NOT NULL
), complete_groups AS (
SELECT g.opportunity_id, g.document_kind
FROM legacy_groups g
WHERE NOT EXISTS (
SELECT 1
FROM commercial_documents d
WHERE d.opportunity_id = g.opportunity_id
AND COALESCE(d.document_kind, '') = g.document_kind
AND NOT EXISTS (
SELECT 1
FROM opportunity_document_links l
WHERE l.opportunity_id = g.opportunity_id
AND COALESCE(l.document_kind, '') = g.document_kind
AND l.document_id = d.id
AND l.ended_at IS NULL
)
)
), effective AS (
SELECT l.opportunity_id, l.document_id, l.document_kind,
l.relationship, l.ended_at
FROM opportunity_document_links l
JOIN complete_groups g
ON g.opportunity_id = l.opportunity_id
AND g.document_kind = COALESCE(l.document_kind, '')
WHERE l.ended_at IS NULL
UNION ALL
SELECT d.opportunity_id, d.id, d.document_kind,
{legacy_relationship} AS relationship,
NULL::timestamptz AS ended_at
FROM commercial_documents d
LEFT JOIN complete_groups g
ON g.opportunity_id = d.opportunity_id
AND g.document_kind = COALESCE(d.document_kind, '')
WHERE d.opportunity_id IS NOT NULL
AND g.opportunity_id IS NULL
)
SELECT opportunity_id, document_id, COALESCE(NULLIF(document_kind, ''), 'unknown'),
relationship, ended_at
FROM effective
"""))
else:
conn.execute(text(f"""
INSERT INTO _effective_document_links
(opportunity_id, document_id, document_kind, relationship, ended_at)
SELECT d.opportunity_id, d.id, COALESCE(NULLIF(d.document_kind, ''), 'unknown'),
{legacy_relationship} AS relationship,
NULL::timestamptz
FROM commercial_documents d
WHERE d.opportunity_id IS NOT NULL
"""))
return "_effective_document_links" return "_effective_document_links"

View File

@@ -6,9 +6,10 @@ opportunity pages, tasks and future audits consume the same decision vocabulary.
from __future__ import annotations from __future__ import annotations
from dataclasses import asdict, dataclass 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 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. # 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)) decision = decide_opportunity_next_action(evidence, load_company_profile(evidence.company_profile))
return decision.to_dict() 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

View File

@@ -320,9 +320,13 @@ def _historical_stage_rates(conn: Any) -> dict[str, dict[str, Any]]:
} }
def _realised_for_period(conn: Any, *, metric: str, period_start: date, period_end: date) -> dict[str, Any]: def _realised_for_period(
conn: Any, *, metric: str, period_start: date, period_end: date,
effective_links_prepared: bool = False,
) -> dict[str, Any]:
from app.document_reconciliation_service import prepare_effective_document_links from app.document_reconciliation_service import prepare_effective_document_links
prepare_effective_document_links(conn) if not effective_links_prepared:
prepare_effective_document_links(conn)
if metric == "cash_received": if metric == "cash_received":
rows = conn.execute(text(""" rows = conn.execute(text("""
WITH latest_doc AS ( WITH latest_doc AS (
@@ -385,9 +389,12 @@ def _realised_for_period(conn: Any, *, metric: str, period_start: date, period_e
} }
def _already_realised_ids(conn: Any, *, metric: str) -> set[str]: def _already_realised_ids(
conn: Any, *, metric: str, effective_links_prepared: bool = False,
) -> set[str]:
from app.document_reconciliation_service import prepare_effective_document_links from app.document_reconciliation_service import prepare_effective_document_links
prepare_effective_document_links(conn) if not effective_links_prepared:
prepare_effective_document_links(conn)
if metric == "cash_received": if metric == "cash_received":
sql = """ sql = """
SELECT DISTINCT opportunity_id::text SELECT DISTINCT opportunity_id::text
@@ -441,8 +448,13 @@ def get_revenue_forecast(*, limit: int = 1000, month: str | None = None, metric:
from app.document_reconciliation_service import prepare_effective_document_links from app.document_reconciliation_service import prepare_effective_document_links
prepare_effective_document_links(conn) prepare_effective_document_links(conn)
historical = _historical_stage_rates(conn) historical = _historical_stage_rates(conn)
realised = _realised_for_period(conn, metric=metric, period_start=period_start, period_end=period_end) realised = _realised_for_period(
already_realised_ids = _already_realised_ids(conn, metric=metric) conn, metric=metric, period_start=period_start, period_end=period_end,
effective_links_prepared=True,
)
already_realised_ids = _already_realised_ids(
conn, metric=metric, effective_links_prepared=True,
)
rows = conn.execute(text(""" rows = conn.execute(text("""
WITH latest_doc AS ( WITH latest_doc AS (
SELECT DISTINCT ON (l.opportunity_id) SELECT DISTINCT ON (l.opportunity_id)

View File

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

View File

@@ -0,0 +1,102 @@
from contextlib import nullcontext
from datetime import date
import app.document_reconciliation_service as document_service
import app.revenue_forecast_service as forecast_service
class _Result:
def __init__(self, rows=(), scalar_value=False):
self.rows = list(rows)
self.scalar_value = scalar_value
def mappings(self):
return self
def all(self):
return self.rows
def scalar(self):
return self.scalar_value
class _Connection:
def __init__(self, *, v2=False):
self.v2 = v2
self.statements = []
def execute(self, statement, params=None):
sql = str(statement)
self.statements.append(sql)
if "to_regclass" in sql:
return _Result(scalar_value=self.v2)
return _Result()
class _Engine:
def __init__(self, connection):
self.connection = connection
def begin(self):
return nullcontext(self.connection)
def test_forecast_prepares_effective_links_once(monkeypatch):
connection = _Connection()
calls = []
flags = []
monkeypatch.setattr(forecast_service, "engine", _Engine(connection))
monkeypatch.setattr(forecast_service, "ensure_revenue_forecast_schema", lambda: None)
monkeypatch.setattr(forecast_service, "get_sales_target", lambda **kwargs: {"target_amount": 0})
monkeypatch.setattr(document_service, "prepare_effective_document_links", lambda conn: calls.append(conn) or "_effective_document_links")
monkeypatch.setattr(forecast_service, "_historical_stage_rates", lambda conn: {})
def realised(conn, **kwargs):
flags.append(kwargs["effective_links_prepared"])
return {"amount": 0.0, "count": 0, "opportunity_ids": set(), "items": []}
def realised_ids(conn, **kwargs):
flags.append(kwargs["effective_links_prepared"])
return set()
monkeypatch.setattr(forecast_service, "_realised_for_period", realised)
monkeypatch.setattr(forecast_service, "_already_realised_ids", realised_ids)
forecast_service.get_revenue_forecast()
assert calls == [connection]
assert flags == [True, True]
def test_helpers_still_prepare_when_called_in_isolation(monkeypatch):
connection = _Connection()
calls = []
monkeypatch.setattr(document_service, "prepare_effective_document_links", lambda conn: calls.append(conn))
forecast_service._realised_for_period(
connection, metric="invoiced", period_start=date(2026, 8, 1), period_end=date(2026, 8, 31),
)
forecast_service._already_realised_ids(connection, metric="invoiced")
assert calls == [connection, connection]
def test_set_based_preparation_has_constant_statement_count_and_no_resolver(monkeypatch):
monkeypatch.setattr(
document_service, "resolve_document_links",
lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("per-opportunity resolver called")),
)
for v2 in (False, True):
connection = _Connection(v2=v2)
document_service.prepare_effective_document_links(connection)
assert len(connection.statements) == 4
assert sum("INSERT INTO _effective_document_links" in sql for sql in connection.statements) == 1
def test_set_based_legacy_relationship_case_matches_resolver_vocabulary():
source = document_service.prepare_effective_document_links.__code__.co_consts
sql_fragments = " ".join(value for value in source if isinstance(value, str))
for relationship in ("REMOVED", "PRIMARY", "HISTORICAL", "SECONDARY"):
assert relationship in sql_fragments
assert "NOT EXISTS" in sql_fragments
assert "l.ended_at IS NULL" in sql_fragments

View File

@@ -4,7 +4,7 @@ from pathlib import Path
def test_opportunity_board_uses_central_next_action(): def test_opportunity_board_uses_central_next_action():
text = Path('app/admin_ui/pages/opportunities.py').read_text() text = Path('app/admin_ui/pages/opportunities.py').read_text()
assert 'def _attach_central_next_actions' in 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 '_attach_central_next_actions(opportunities)' in text
assert 'central.get("label")' in text assert 'central.get("label")' in text
assert 'CLOSE_OPPORTUNITY' in text assert 'CLOSE_OPPORTUNITY' in text