fix: canonicalize operations work items

This commit is contained in:
plx
2026-08-15 00:32:38 +00:00
parent 9b3c3b32ef
commit fabadbffdc
4 changed files with 1041 additions and 16 deletions

474
app/canonical_operations.py Normal file
View File

@@ -0,0 +1,474 @@
"""Pure canonical read model for the operator workbench.
This module deliberately performs no I/O. It projects database source rows and
already-computed opportunity decisions into one current work item per process.
Persistent tasks, communications, reconciliation records and outbox rows remain
untouched and are retained as source evidence.
"""
from __future__ import annotations
from collections import defaultdict
import json
from typing import Any, Iterable, Mapping
from app.operation_noise import is_low_value_no_opportunity_item, is_noise_operation_item
from app.work_center_action_policy import canonical_action_code
WAITING_ACTION_CODES = {
"WAIT_PRODUCTION",
"WAIT_CUSTOMER",
"WAIT_PAYMENT",
"WAIT_SUPPLIER",
"WAIT_LOGISTICS",
"WAIT_SCHEDULED_DATE",
}
NO_WORK_ACTION_CODES = {"NO_ACTION", "NOT_FOUND"}
def _s(value: Any) -> str:
return str(value or "").strip()
def _upper(value: Any) -> str:
return _s(value).upper()
def _source_ref(row: Mapping[str, Any]) -> dict[str, Any]:
metadata = row.get("item_metadata") if isinstance(row.get("item_metadata"), dict) else {}
ref = {
"source": _s(row.get("source")),
"id": _s(row.get("id")),
"status": _s(row.get("status")),
"action_code": _upper(row.get("action_code")),
}
source_id_fields = {
"task": "task_id",
"communication": "communication_id",
"reconciliation": "reconciliation_item_id",
"outbox": "outbox_id",
}
if ref["id"] and ref["source"] in source_id_fields:
ref[source_id_fields[ref["source"]]] = ref["id"]
for key in (
"task_id", "communication_id", "message_id", "raw_event_id",
"action_run_id", "conversation_id", "reconciliation_item_id",
"outbox_id", "source_system", "source_event_id",
):
value = row.get(key) or metadata.get(key)
if value not in (None, ""):
ref[key] = str(value)
return ref
def _priority_rank(value: Any) -> int:
return {"urgente": 0, "critical": 0, "alta": 1, "high": 1, "normal": 2, "baixa": 3}.get(
_s(value).lower(), 3
)
def _best_priority(rows: Iterable[Mapping[str, Any]], decision: Mapping[str, Any] | None = None) -> str:
values = [_s((decision or {}).get("priority"))]
values.extend(_s(row.get("priority")) for row in rows)
values = [value for value in values if value]
return min(values, key=_priority_rank) if values else "normal"
def _canonical_code(value: Any) -> str:
code = canonical_action_code(value)
return _upper(code)
def _is_waiting(decision: Mapping[str, Any]) -> bool:
code = _canonical_code(decision.get("action_code"))
# Unknown WAIT_* actions remain visible until explicitly classified. This
# favors a safe false positive over silently hiding a new human action.
return code in WAITING_ACTION_CODES
def _matching_task(rows: Iterable[Mapping[str, Any]], action_code: str) -> Mapping[str, Any] | None:
for row in rows:
if row.get("source") == "task" and _canonical_code(row.get("action_code")) == action_code:
return row
return None
def _outbox_matches_action(row: Mapping[str, Any], action_code: str) -> bool:
metadata = row.get("item_metadata") if isinstance(row.get("item_metadata"), dict) else {}
if metadata.get("independent_exception") is True or metadata.get("human_intervention_scope") == "independent":
return False
candidates = {
_canonical_code(row.get("action_code")),
_canonical_code(metadata.get("action_code")),
_canonical_code(metadata.get("task_action_code")),
_canonical_code(metadata.get("current_action_code")),
}
candidates.discard("")
return action_code in candidates
def _explicit_task_id(row: Mapping[str, Any]) -> str:
metadata = row.get("item_metadata") if isinstance(row.get("item_metadata"), dict) else {}
return _s(row.get("task_id") or metadata.get("task_id"))
def _standalone_group_key(
row: Mapping[str, Any],
*,
explicitly_referenced_task_ids: set[str],
) -> str:
"""Return a conservative grouping key for rows without an opportunity.
Conversation identity is a safe process boundary, but not enough to prove
that different actions are equivalent. Cross-action merging is allowed
only through an explicit communication→task reference.
"""
process_key = _standalone_process_key(row)
if not process_key.startswith("conversation:"):
return process_key
related_task_id = _explicit_task_id(row)
if related_task_id:
return f"{process_key}:related-task:{related_task_id}"
row_id = _s(row.get("id"))
if row.get("source") == "task" and row_id in explicitly_referenced_task_ids:
return f"{process_key}:related-task:{row_id}"
return f"{process_key}:action:{_canonical_code(row.get('action_code')) or 'REVIEW_MANUALLY'}"
def _has_explicit_unresolved_review(rows: Iterable[Mapping[str, Any]]) -> bool:
review_codes = {
"REVIEW",
"REVIEW_MANUALLY",
"ASSOCIATE_OPPORTUNITY",
"REVIEW_ASSOCIATION",
"LINK_DOCUMENT",
}
review_statuses = {"needs_review", "review_required", "ambiguous", "conflict"}
review_flags = {
"needs_review",
"review_required",
"association_review_required",
"customer_association_review_required",
"opportunity_association_review_required",
}
for row in rows:
metadata = row.get("item_metadata") if isinstance(row.get("item_metadata"), dict) else {}
if _canonical_code(row.get("action_code")) in review_codes:
return True
if _s(row.get("status")).lower() in review_statuses:
return True
if _s(row.get("opportunity_linking_status")).lower() in review_statuses:
return True
if any(metadata.get(flag) is True for flag in review_flags):
return True
return False
def opportunity_ids_from_work_seeds(rows: Iterable[Mapping[str, Any]]) -> list[str]:
"""Return decision candidates from eligible seeds only."""
return sorted({_s(row.get("opportunity_id")) for row in rows if _s(row.get("opportunity_id"))})
def partition_canonical_items(
items: Iterable[Mapping[str, Any]], *, display_limit: int
) -> dict[str, Any]:
all_items = [dict(item) for item in items]
actionable = [item for item in all_items if item.get("operational_queue") != "waiting"]
waiting = [item for item in all_items if item.get("operational_queue") == "waiting"]
return {
"actionable_items": actionable,
"waiting_items": waiting,
"visible_actionable_items": actionable[:display_limit],
"visible_waiting_items": waiting[:display_limit],
"work_queue_total": len(actionable),
"waiting_total": len(waiting),
}
def _unique_refs(items: Iterable[Mapping[str, Any]], field: str) -> list[dict[str, Any]]:
refs: list[dict[str, Any]] = []
seen: set[str] = set()
for item in items:
for ref in item.get(field) or []:
value = dict(ref) if isinstance(ref, Mapping) else {"value": str(ref)}
source = _s(value.get("source"))
source_id = _s(value.get("id"))
fingerprint = (
f"{source}:{source_id}"
if source and source_id
else json.dumps(value, ensure_ascii=False, sort_keys=True, default=str)
)
if fingerprint in seen:
continue
seen.add(fingerprint)
refs.append(value)
return refs
def _has_task_execution_handle(item: Mapping[str, Any]) -> bool:
return item.get("source") == "task" or _s(item.get("href")).startswith("/tasks/")
def _useful_identity(value: Any) -> bool:
normalized = _s(value).casefold()
return normalized not in {"", "contacto sem identificação", "contacto", "cliente", "geral"}
def _merge_exact_work_items(items: Iterable[Mapping[str, Any]]) -> list[dict[str, Any]]:
"""Enforce one item for an already-resolved process/action identity."""
groups: dict[str, list[Mapping[str, Any]]] = defaultdict(list)
for item in items:
groups[_s(item.get("work_item_key"))].append(item)
merged_items: list[dict[str, Any]] = []
for work_item_key, group in groups.items():
# A task-backed item is the best execution handle. Within the same
# handle class, retain the highest existing priority.
primary = min(
group,
key=lambda item: (
0 if _has_task_execution_handle(item) else 1,
_priority_rank(item.get("priority")),
),
)
merged = dict(primary)
source_refs = _unique_refs(group, "source_refs")
stale_task_refs = _unique_refs(group, "stale_task_refs")
merged["source_refs"] = source_refs
merged["stale_task_refs"] = stale_task_refs
merged["raw_source_count"] = len(source_refs)
merged["priority"] = _best_priority(group)
merged["created_at"] = max((_s(item.get("created_at")) for item in group), default="")
for field in ("customer_name", "contact_display_name", "fiscal_customer_name", "customer_email"):
if _useful_identity(merged.get(field)):
continue
replacement = next((item.get(field) for item in group if _useful_identity(item.get(field))), None)
if replacement:
merged[field] = replacement
# Defensive fallback for unusual projected fixtures: never discard a
# task URL when a task-backed duplicate contains one.
if not _s(merged.get("href")).startswith("/tasks/"):
task_href = next(
(_s(item.get("href")) for item in group if _s(item.get("href")).startswith("/tasks/")),
"",
)
if task_href:
merged["href"] = task_href
merged_items.append(merged)
keys = [_s(item.get("work_item_key")) for item in merged_items]
if len(keys) != len(set(keys)):
raise AssertionError("canonical Operations invariant violated: duplicate work_item_key")
return merged_items
def _standalone_process_key(row: Mapping[str, Any]) -> str:
source = _s(row.get("source")) or "unknown"
source_system = _s(row.get("source_system")) or source
conversation_id = _s(row.get("conversation_id"))
if conversation_id and source in {"task", "communication"}:
return f"conversation:{source_system}:{conversation_id}"
if source == "reconciliation":
group_key = _s(row.get("process_group_key") or row.get("source_event_id") or row.get("id"))
return f"reconciliation:{group_key}"
if source == "outbox":
return f"integration-exception:{source_system}:{_s(row.get('id'))}"
return f"{source}:{_s(row.get('id'))}"
def _standalone_item(row: Mapping[str, Any]) -> dict[str, Any]:
process_key = _standalone_process_key(row)
action_code = _canonical_code(row.get("action_code")) or "REVIEW_MANUALLY"
queue = "exception" if row.get("source") == "outbox" else _s(row.get("queue")) or "rever"
item = dict(row)
item.update({
"process_key": process_key,
"work_item_key": f"{process_key}:action:{action_code}",
"current_action_code": action_code,
"action_code": action_code,
"operational_queue": queue,
"why_human_required": _s(row.get("detail")) or "A intervenção ainda não está associada com segurança a um processo comercial.",
"source_refs": [_source_ref(row)],
"raw_source_count": 1,
})
return item
def _merge_standalone_conversation(rows: list[Mapping[str, Any]]) -> dict[str, Any]:
# A conversation is a safe identity boundary. Prefer a task as the execution
# handle, but retain every communication/message row as evidence.
ordered = sorted(rows, key=lambda row: (row.get("source") != "task", _priority_rank(row.get("priority"))))
primary = ordered[0]
item = _standalone_item(primary)
item["source_refs"] = [_source_ref(row) for row in rows]
item["raw_source_count"] = len(rows)
item["priority"] = _best_priority(rows)
return item
def _opportunity_item(
opportunity_id: str,
rows: list[Mapping[str, Any]],
decision: Mapping[str, Any],
) -> tuple[dict[str, Any] | None, list[dict[str, Any]]]:
action_code = _canonical_code(decision.get("action_code")) or "REVIEW"
independent_exceptions: list[dict[str, Any]] = []
evidence_rows: list[Mapping[str, Any]] = []
for row in rows:
if row.get("source") == "outbox" and not row.get("_evidence_only") and not _outbox_matches_action(row, action_code):
independent_exceptions.append(_standalone_item(row))
else:
evidence_rows.append(row)
if action_code in NO_WORK_ACTION_CODES:
if not _has_explicit_unresolved_review(evidence_rows):
return None, independent_exceptions
# The central decision normally wins, but an explicit unresolved review
# marker must not be hidden by NO_ACTION/NOT_FOUND.
action_code = "REVIEW"
decision = {
**dict(decision),
"action_code": action_code,
"label": "Rever evidência pendente",
"description": "A decisão central indica que não há ação, mas existe evidência explicitamente marcada para revisão.",
"reason": "Existe revisão humana não resolvida apesar da decisão central sem ação.",
"priority": "normal",
"target_url": f"/opportunities/{opportunity_id}",
}
matching_task = _matching_task(evidence_rows, action_code)
primary = matching_task or next((row for row in evidence_rows if row.get("source") == "task"), None)
primary = primary or (evidence_rows[0] if evidence_rows else {})
process_key = f"opportunity:{opportunity_id}"
waiting = _is_waiting(decision)
opportunity_row = next((row for row in evidence_rows if _s(row.get("fiscal_customer_name"))), None)
opportunity_row = opportunity_row or next((row for row in evidence_rows if _s(row.get("opportunity_title"))), None)
opportunity_row = opportunity_row or primary
item = dict(primary)
item.update({
"id": _s((matching_task or primary).get("id")) or opportunity_id,
"source": "opportunity",
"process_key": process_key,
"work_item_key": f"{process_key}:action:{action_code}",
"current_action_code": action_code,
"action_code": action_code,
"title": _s(decision.get("label")) or _s(primary.get("title")) or action_code,
"detail": _s(decision.get("description") or decision.get("reason")),
"why_human_required": _s(decision.get("reason") or decision.get("description")),
"priority": _s(decision.get("priority")) or _best_priority(evidence_rows),
"operational_queue": "waiting" if waiting else "do_now",
"queue": _s((matching_task or primary).get("queue")) or _s(primary.get("queue")) or "rever",
"status": "waiting" if waiting else "pending",
"href": _s((matching_task or {}).get("href")) or _s(decision.get("target_url")) or f"/opportunities/{opportunity_id}",
"action_label": _s(decision.get("label")) or _s(primary.get("action_label")) or "Abrir",
"opportunity_id": opportunity_id,
"opportunity_title": _s(opportunity_row.get("opportunity_title")),
"customer_name": _s(opportunity_row.get("fiscal_customer_name") or opportunity_row.get("customer_name")),
"contact_display_name": _s(opportunity_row.get("fiscal_customer_name") or opportunity_row.get("contact_display_name")),
"source_refs": [_source_ref(row) for row in evidence_rows],
"raw_source_count": len(evidence_rows),
"stale_task_refs": [
_source_ref(row) for row in evidence_rows
if row.get("source") == "task" and _canonical_code(row.get("action_code")) != action_code
],
"decision_version": decision.get("decision_version"),
"decision": dict(decision),
})
return item, independent_exceptions
def canonicalize_operations(
source_rows: Iterable[Mapping[str, Any]],
opportunity_decisions: Mapping[str, Mapping[str, Any]],
*,
evidence_rows: Iterable[Mapping[str, Any]] | None = None,
display_limit: int | None = None,
) -> dict[str, Any]:
"""Project eligible work seeds into stable, canonical work items.
``evidence_rows`` can enrich only opportunities already present in the
cleaned seed set. Evidence can never select an opportunity or create a
standalone/exception item.
"""
raw_rows = [dict(row) for row in source_rows]
clean_rows = [row for row in raw_rows if not is_noise_operation_item(row) and not is_low_value_no_opportunity_item(row)]
raw_evidence_rows = [dict(row) for row in (evidence_rows or [])]
clean_evidence_rows = [
row for row in raw_evidence_rows
if not is_noise_operation_item(row) and not is_low_value_no_opportunity_item(row)
]
by_opportunity: dict[str, list[Mapping[str, Any]]] = defaultdict(list)
standalone_groups: dict[str, list[Mapping[str, Any]]] = defaultdict(list)
explicitly_referenced_task_ids = {
_explicit_task_id(row) for row in clean_rows
if row.get("source") == "communication" and _explicit_task_id(row)
}
for row in clean_rows:
opportunity_id = _s(row.get("opportunity_id"))
if opportunity_id:
by_opportunity[opportunity_id].append(row)
else:
standalone_groups[_standalone_group_key(
row, explicitly_referenced_task_ids=explicitly_referenced_task_ids
)].append(row)
# Crucial eligibility boundary: evidence-only opportunity IDs are ignored.
# Mark retained rows so downstream exception handling cannot turn evidence
# into a second work item.
selected_opportunity_ids = set(by_opportunity)
retained_evidence_rows = []
for row in clean_evidence_rows:
opportunity_id = _s(row.get("opportunity_id"))
if opportunity_id not in selected_opportunity_ids:
continue
evidence = {**row, "_evidence_only": True}
by_opportunity[opportunity_id].append(evidence)
retained_evidence_rows.append(evidence)
items: list[dict[str, Any]] = []
for opportunity_id, rows in by_opportunity.items():
decision = opportunity_decisions.get(opportunity_id)
if not decision:
# A calculation failure must not silently hide customer work.
fallback = _merge_standalone_conversation(rows)
fallback["process_key"] = f"opportunity:{opportunity_id}"
fallback["work_item_key"] = f"opportunity:{opportunity_id}:action:REVIEW"
fallback["current_action_code"] = fallback["action_code"] = "REVIEW"
fallback["why_human_required"] = "Não foi possível calcular com segurança a ação atual da oportunidade."
fallback["opportunity_id"] = opportunity_id
items.append(fallback)
continue
item, exceptions = _opportunity_item(opportunity_id, rows, decision)
if item is not None:
items.append(item)
items.extend(exceptions)
for rows in standalone_groups.values():
# Only a reliable conversation identity is merged. Other process keys are
# record-specific and therefore remain separate conservatively.
items.append(_merge_standalone_conversation(rows))
# Relationship-aware standalone grouping may intentionally produce separate
# intermediate groups. Exact process/action identity is nevertheless the
# final safe uniqueness boundary.
items = _merge_exact_work_items(items)
# Preserve the existing Operations ordering: newest first inside each
# priority band, without allowing supporting/stale rows to set the band.
items.sort(key=lambda item: _s(item.get("created_at")), reverse=True)
items.sort(key=lambda item: _priority_rank(item.get("priority")))
all_items = items
visible_items = all_items[: max(0, int(display_limit))] if display_limit is not None else all_items
return {
"items": visible_items,
"all_items": all_items,
"raw_source_count": len(raw_rows) + len(raw_evidence_rows),
"clean_source_count": len(clean_rows) + len(retained_evidence_rows),
"seed_source_count": len(clean_rows),
"evidence_source_count": len(retained_evidence_rows),
"canonical_count": len(all_items),
"visible_count": len(visible_items),
}

View File

@@ -21,6 +21,13 @@ from app.work_center_action_policy import (
canonical_action_code, canonical_action_code,
reconstructed_review_required, reconstructed_review_required,
) )
from app.canonical_operations import (
canonicalize_operations,
opportunity_ids_from_work_seeds,
partition_canonical_items,
)
from app.opportunity_next_action_service import get_opportunity_next_actions
from app.document_reconciliation_service import active_document_link_exclusion_sql
def _int(value: Any) -> int: def _int(value: Any) -> int:
@@ -303,6 +310,11 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
Technical lists remain available in their own pages and should only appear Technical lists remain available in their own pages and should only appear
here when they block an operator action. here when they block an operator action.
""" """
display_limit = max(1, min(int(limit), 200))
# Bound each entity source independently, then apply the user-facing limit
# after canonicalization. This prevents duplicate rows from consuming the
# display limit while still avoiding an unbounded history scan.
candidate_limit = max(200, min(display_limit * 20, 1000))
with engine.begin() as conn: with engine.begin() as conn:
counts = conn.execute(text(""" counts = conn.execute(text("""
SELECT SELECT
@@ -383,7 +395,7 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
LIMIT :limit LIMIT :limit
"""), {"limit": int(limit)}).mappings().all() """), {"limit": int(limit)}).mappings().all()
work_items = conn.execute(text(""" work_seed_rows = conn.execute(text("""
SELECT * FROM ( SELECT * FROM (
SELECT 'task' AS source, t.id::text AS id, t.created_at, t.due_at, SELECT 'task' AS source, t.id::text AS id, t.created_at, t.due_at,
CASE CASE
@@ -433,14 +445,19 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(o.stage, '') AS opportunity_stage, COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.value_amount, 0) AS opportunity_value_amount, COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM tasks t FROM (
SELECT * FROM tasks
WHERE status = 'pending'
AND NOT (action_code LIKE 'FOLLOW_UP_%' AND due_at IS NOT NULL AND due_at > now())
ORDER BY created_at DESC
LIMIT :candidate_limit
) t
LEFT JOIN opportunities o ON o.id = t.opportunity_id LEFT JOIN opportunities o ON o.id = t.opportunity_id
LEFT JOIN customers cu_opp ON cu_opp.id = o.local_customer_id LEFT JOIN customers cu_opp ON cu_opp.id = o.local_customer_id
LEFT JOIN customers cu_task ON cu_task.id::text = t.customer_id LEFT JOIN customers cu_task ON cu_task.id::text = t.customer_id
LEFT JOIN messages m ON m.id = t.message_id LEFT JOIN messages m ON m.id = t.message_id
LEFT JOIN raw_events re ON re.id = t.raw_event_id LEFT JOIN raw_events re ON re.id = t.raw_event_id
WHERE t.status = 'pending' WHERE TRUE
AND NOT (t.action_code LIKE 'FOLLOW_UP_%' AND t.due_at IS NOT NULL AND t.due_at > now())
UNION ALL UNION ALL
@@ -477,11 +494,16 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(o.stage, '') AS opportunity_stage, COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.value_amount, 0) AS opportunity_value_amount, COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM integration_outbox io FROM (
SELECT * FROM integration_outbox
WHERE status IN ('failed','blocked')
AND NOT (status = 'failed' AND (COALESCE(last_error,'') ILIKE '%limpo manualmente%' OR COALESCE(last_error,'') ILIKE '%resolvido manualmente%'))
ORDER BY created_at DESC
LIMIT :candidate_limit
) io
LEFT JOIN opportunities o ON o.id::text = NULLIF(io.payload->>'opportunity_id','') LEFT JOIN opportunities o ON o.id::text = NULLIF(io.payload->>'opportunity_id','')
LEFT JOIN customers cu ON cu.id = o.local_customer_id LEFT JOIN customers cu ON cu.id = o.local_customer_id
WHERE io.status IN ('pending','failed','blocked') WHERE TRUE
AND NOT (io.status = 'failed' AND (COALESCE(io.last_error,'') ILIKE '%limpo manualmente%' OR COALESCE(io.last_error,'') ILIKE '%resolvido manualmente%'))
UNION ALL UNION ALL
@@ -519,29 +541,139 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
'/communications/' || c.id::text AS href, '/communications/' || c.id::text AS href,
CASE WHEN c.customer_id IS NULL THEN 'Associar cliente' ELSE 'Abrir' END AS action_label, CASE WHEN c.customer_id IS NULL THEN 'Associar cliente' ELSE 'Abrir' END AS action_label,
''::text AS opportunity_linking_status, ''::text AS opportunity_linking_status,
COALESCE(c.metadata, '{}'::jsonb) AS item_metadata, COALESCE(c.metadata, '{}'::jsonb)
|| CASE WHEN c.task_id IS NOT NULL
THEN jsonb_build_object('task_id', c.task_id::text)
ELSE '{}'::jsonb END AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata, COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage, COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.value_amount, 0) AS opportunity_value_amount, COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM communications c FROM (
SELECT * FROM communications
WHERE status IN ('new','classified','needs_review')
ORDER BY created_at DESC
LIMIT :candidate_limit
) c
LEFT JOIN customers cu ON cu.id = c.customer_id LEFT JOIN customers cu ON cu.id = c.customer_id
LEFT JOIN opportunities o ON o.id = c.opportunity_id LEFT JOIN opportunities o ON o.id = c.opportunity_id
WHERE c.status IN ('new','classified','needs_review') WHERE TRUE
UNION ALL
SELECT 'reconciliation' AS source, ri.id::text AS id, ri.created_at, NULL::timestamptz AS due_at,
COALESCE(ri.priority, 'normal') AS priority,
'rever' AS queue,
upper(COALESCE(ri.suggested_action, 'RECONCILE_DOCUMENTS')) AS action_code,
COALESCE(ri.title, 'Rever reconciliação') AS title,
COALESCE(ri.description, ri.resolution_note, '') AS detail,
ri.status,
ri.source_system,
NULL::text AS conversation_id,
NULL::text AS contact_id,
COALESCE(ri.description, '') AS request_text,
ri.opportunity_id::text,
COALESCE(o.title, '') AS opportunity_title,
COALESCE(cu.name, ri.customer_name, ri.customer_email, '') AS customer_name,
COALESCE(ri.customer_name, '') AS sender_name,
COALESCE(ri.customer_email, '') AS sender_email,
COALESCE(cu.name, '') AS fiscal_customer_name,
COALESCE(cu.email, '') AS fiscal_customer_email,
COALESCE(cu.tax_id, '') AS fiscal_customer_tax_id,
COALESCE(cu.street_name, '') AS fiscal_customer_street_name,
COALESCE(cu.postal_zone, '') AS fiscal_customer_postal_zone,
COALESCE(cu.city_name, '') AS fiscal_customer_city_name,
COALESCE(cu.name, ri.customer_name, ri.customer_email, '') AS contact_display_name,
COALESCE(ri.customer_email, '') AS customer_email,
''::text AS no_opportunity_reason,
'/reconciliation' AS href,
'Rever' AS action_label,
''::text AS opportunity_linking_status,
COALESCE(ri.payload, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.value_amount, ri.amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, ri.currency, 'EUR') AS opportunity_currency
FROM (
SELECT * FROM reconciliation_items
WHERE status IN ('open','needs_review','conflict')
AND """ + active_document_link_exclusion_sql("reconciliation_items") + """
ORDER BY created_at DESC
LIMIT :candidate_limit
) ri
LEFT JOIN opportunities o ON o.id = ri.opportunity_id
LEFT JOIN customers cu ON cu.id = COALESCE(o.local_customer_id, ri.customer_id)
WHERE TRUE
) items ) items
ORDER BY ORDER BY
CASE lower(priority) WHEN 'alta' THEN 1 WHEN 'high' THEN 1 WHEN 'urgente' THEN 0 WHEN 'normal' THEN 2 ELSE 3 END, CASE lower(priority) WHEN 'alta' THEN 1 WHEN 'high' THEN 1 WHEN 'urgente' THEN 0 WHEN 'normal' THEN 2 ELSE 3 END,
created_at DESC created_at DESC
LIMIT :limit """), {"candidate_limit": candidate_limit}).mappings().all()
"""), {"limit": int(limit)}).mappings().all()
cleaned_work_items = _attach_operation_urls(_normalise_work_item_intent([dict(r) for r in work_items])) seed_opportunity_ids = opportunity_ids_from_work_seeds(work_seed_rows)
evidence_rows = []
if seed_opportunity_ids:
evidence_rows = conn.execute(text("""
SELECT * FROM (
SELECT 'communication' AS source, c.id::text AS id, c.created_at,
c.status, upper(COALESCE(c.classification, 'REVIEW_MANUALLY')) AS action_code,
c.opportunity_id::text, c.source_system, c.conversation_id,
COALESCE(c.metadata, '{}'::jsonb)
|| CASE WHEN c.task_id IS NOT NULL
THEN jsonb_build_object('task_id', c.task_id::text)
ELSE '{}'::jsonb END AS item_metadata
FROM communications c
WHERE c.opportunity_id = ANY(CAST(:opportunity_ids AS UUID[]))
AND c.status NOT IN ('new','classified','needs_review')
ORDER BY c.created_at DESC
LIMIT :candidate_limit
) communication_evidence
UNION ALL
SELECT * FROM (
SELECT 'reconciliation' AS source, ri.id::text AS id, ri.created_at,
ri.status, upper(COALESCE(ri.suggested_action, 'RECONCILE_DOCUMENTS')) AS action_code,
ri.opportunity_id::text, ri.source_system, NULL::text AS conversation_id,
COALESCE(ri.payload, '{}'::jsonb) AS item_metadata
FROM reconciliation_items ri
WHERE ri.opportunity_id = ANY(CAST(:opportunity_ids AS UUID[]))
AND ri.status IN ('linked','resolved')
ORDER BY ri.created_at DESC
LIMIT :candidate_limit
) reconciliation_evidence
"""), {
"opportunity_ids": seed_opportunity_ids,
"candidate_limit": candidate_limit,
}).mappings().all()
normalized_rows = _normalise_work_item_intent([dict(r) for r in work_seed_rows])
opportunity_ids = opportunity_ids_from_work_seeds(normalized_rows)
decisions = get_opportunity_next_actions(opportunity_ids) if opportunity_ids else {}
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
canonical_items = _attach_operation_urls(list(projection["items"]))
partition = partition_canonical_items(canonical_items, display_limit=display_limit)
actionable_items = partition["actionable_items"]
all_waiting_items = partition["waiting_items"]
cleaned_work_items = partition["visible_actionable_items"]
waiting_items = partition["visible_waiting_items"]
cleaned_counts = {k: _int(v) for k, v in dict(counts).items()} cleaned_counts = {k: _int(v) for k, v in dict(counts).items()}
# v4.9.0: the visible Operations total should match the queue the # v4.9.0: the visible Operations total should match the queue the
# operator can actually act on, not raw pending tasks that include mailbox # operator can actually act on, not raw pending tasks that include mailbox
# noise awaiting cleanup. The cleanup script still fixes the data source. # noise awaiting cleanup. The cleanup script still fixes the data source.
cleaned_counts["work_queue_total"] = len(cleaned_work_items) cleaned_counts["work_queue_total"] = partition["work_queue_total"]
cleaned_counts["waiting_total"] = partition["waiting_total"]
cleaned_counts["raw_source_count"] = projection["raw_source_count"]
cleaned_counts["canonical_total"] = len(canonical_items)
diagnostics = {
"candidate_limit": candidate_limit,
"raw_source_count": projection["raw_source_count"],
"clean_source_count": projection["clean_source_count"],
"seed_source_count": projection["seed_source_count"],
"evidence_source_count": projection["evidence_source_count"],
"canonical_count": len(canonical_items),
"visible_count": len(cleaned_work_items),
"waiting_total": len(all_waiting_items),
}
return { return {
"counts": cleaned_counts, "counts": cleaned_counts,
"recent_outbox": [dict(r) for r in recent_outbox], "recent_outbox": [dict(r) for r in recent_outbox],
@@ -550,6 +682,9 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"incomplete_customers": [dict(r) for r in incomplete_customers], "incomplete_customers": [dict(r) for r in incomplete_customers],
"recent_communications": [dict(r) for r in recent_communications], "recent_communications": [dict(r) for r in recent_communications],
"work_items": cleaned_work_items, "work_items": cleaned_work_items,
"waiting_items": waiting_items,
**diagnostics,
"projection_metrics": diagnostics,
} }

View File

@@ -302,8 +302,9 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di
if not ids: if not ids:
return {} return {}
from app.operation_service import ensure_operation_schema # Read path only: schema creation belongs to startup/migrations. In
ensure_operation_schema() # particular, GET /operations must never perform DDL while calculating its
# canonical projection.
with engine.begin() as conn: with engine.begin() as conn:
opportunities = _bulk_rows(conn, """ opportunities = _bulk_rows(conn, """
SELECT id::text, stage, status, title, SELECT id::text, stage, status, title,

View File

@@ -0,0 +1,415 @@
from __future__ import annotations
from pathlib import Path
from app.canonical_operations import (
_merge_exact_work_items,
canonicalize_operations,
opportunity_ids_from_work_seeds,
partition_canonical_items,
)
def row(source="task", id="1", opportunity_id="opp-1", action_code="SEND_QUOTE", **extra):
value = {
"source": source,
"id": id,
"opportunity_id": opportunity_id,
"action_code": action_code,
"status": "pending",
"priority": "normal",
"queue": "vendas",
"title": action_code,
"detail": "Requer intervenção",
"source_system": "chatwoot" if source in {"task", "communication"} else source,
"created_at": "2026-08-15T10:00:00+00:00",
"item_metadata": {},
}
value.update(extra)
return value
def decision(code="SEND_QUOTE", **extra):
value = {
"action_code": code,
"label": code.replace("_", " ").title(),
"description": "Ação canónica calculada.",
"reason": "O processo requer esta ação agora.",
"priority": "normal",
"target_url": "/opportunities/opp-1",
"decision_version": "test-v1",
}
value.update(extra)
return value
def project(rows, decisions=None, limit=None, evidence_rows=None):
return canonicalize_operations(rows, decisions or {}, evidence_rows=evidence_rows, display_limit=limit)
def test_task_and_communication_for_same_opportunity_are_one_card_and_evidence():
result = project([
row(id="task-1"),
row(source="communication", id="comm-1"),
], {"opp-1": decision()})
assert result["canonical_count"] == 1
assert {ref["source"] for ref in result["items"][0]["source_refs"]} == {"task", "communication"}
def test_only_task_matching_canonical_action_is_current_and_stale_task_is_evidence():
result = project([
row(id="stale", action_code="SEND_PROFORMA"),
row(id="current", action_code="CONFIRM_PAYMENT"),
], {"opp-1": decision("CONFIRM_PAYMENT")})
item = result["items"][0]
assert item["id"] == "current"
assert item["current_action_code"] == "CONFIRM_PAYMENT"
assert [ref["task_id"] for ref in item["stale_task_refs"]] == ["stale"]
def test_linked_communication_is_evidence_not_duplicate():
result = project([row(source="communication", id="comm-1")], {"opp-1": decision()})
assert len(result["items"]) == 1
assert result["items"][0]["source_refs"][0]["communication_id"] == "comm-1"
def test_standalone_actionable_and_uncertain_communications_remain_visible():
result = project([
row(source="communication", id="a", opportunity_id="", conversation_id="10", action_code="SUPPORT"),
row(source="communication", id="b", opportunity_id="", conversation_id="11", action_code="REVIEW_MANUALLY"),
])
assert result["canonical_count"] == 2
assert {item["current_action_code"] for item in result["items"]} == {"SUPPORT", "REVIEW_MANUALLY"}
def test_same_reliable_standalone_conversation_is_one_process_but_subject_is_not_used():
result = project([
row(source="task", id="a", opportunity_id="", conversation_id="10", title="same"),
row(source="communication", id="b", opportunity_id="", conversation_id="10", title="same"),
row(source="communication", id="c", opportunity_id="", conversation_id="11", title="same"),
])
assert result["canonical_count"] == 2
assert any(item["raw_source_count"] == 2 for item in result["items"])
def test_same_conversation_with_different_actions_remains_two_work_items():
result = project([
row(source="task", id="quote", opportunity_id="", conversation_id="10", action_code="SEND_QUOTE"),
row(source="communication", id="support", opportunity_id="", conversation_id="10", action_code="SUPPORT"),
])
assert result["canonical_count"] == 2
assert {item["current_action_code"] for item in result["items"]} == {"SEND_QUOTE", "SUPPORT"}
def test_explicit_communication_task_relation_can_merge_different_actions():
result = project([
row(source="task", id="task-1", opportunity_id="", conversation_id="10", action_code="SEND_QUOTE"),
row(
source="communication", id="comm-1", opportunity_id="", conversation_id="10",
action_code="REVIEW_MANUALLY", item_metadata={"task_id": "task-1"},
),
])
assert result["canonical_count"] == 1
assert result["items"][0]["current_action_code"] == "SEND_QUOTE"
assert {ref["source"] for ref in result["items"][0]["source_refs"]} == {"task", "communication"}
def test_final_exact_key_merge_combines_independent_same_action_groups_and_keeps_task_handle():
result = project([
row(
source="task", id="task-1", opportunity_id="", conversation_id="2433",
action_code="SEND_INFO", href="/tasks/task-1", priority="normal",
),
row(
source="communication", id="related", opportunity_id="", conversation_id="2433",
action_code="SEND_INFO", item_metadata={"task_id": "task-1"},
),
row(
source="communication", id="independent", opportunity_id="", conversation_id="2433",
action_code="SEND_INFO", priority="alta", customer_name="Cliente conhecido",
),
])
assert result["canonical_count"] == 1
item = result["items"][0]
assert item["work_item_key"] == "conversation:chatwoot:2433:action:SEND_INFO"
assert item["source"] == "task"
assert item["href"] == "/tasks/task-1"
assert item["priority"] == "alta"
assert item["customer_name"] == "Cliente conhecido"
assert item["created_at"] == "2026-08-15T10:00:00+00:00"
assert item["raw_source_count"] == 3
assert {(ref["source"], ref["id"]) for ref in item["source_refs"]} == {
("task", "task-1"), ("communication", "related"), ("communication", "independent")
}
def test_exact_key_merge_unions_unique_refs_and_stale_task_refs():
duplicate_ref = {"source": "communication", "id": "same"}
item = _merge_exact_work_items([
{
"work_item_key": "conversation:chatwoot:2433:action:SEND_INFO",
"process_key": "conversation:chatwoot:2433", "current_action_code": "SEND_INFO",
"source_refs": [duplicate_ref, {"source": "task", "id": "task-1"}],
"stale_task_refs": [{"source": "task", "id": "stale-1"}],
"priority": "normal",
},
{
"work_item_key": "conversation:chatwoot:2433:action:SEND_INFO",
"process_key": "conversation:chatwoot:2433", "current_action_code": "SEND_INFO",
"source_refs": [duplicate_ref, {"source": "communication", "id": "comm-2"}],
"stale_task_refs": [
{"source": "task", "id": "stale-1"},
{"source": "task", "id": "stale-2"},
],
"priority": "normal",
},
])[0]
assert len(item["source_refs"]) == 3
assert {(ref["source"], ref["id"]) for ref in item["source_refs"]} == {
("communication", "same"), ("task", "task-1"), ("communication", "comm-2")
}
assert {ref["id"] for ref in item["stale_task_refs"]} == {"stale-1", "stale-2"}
def test_deterministic_noise_is_excluded():
result = project([
row(opportunity_id="", action_code="IGNORE_BOUNCE", title="Mailer-Daemon"),
row(opportunity_id="", action_code="IGNORE_SPAM"),
row(opportunity_id="", action_code="NO_ACTION", title="Automated notification"),
])
assert result["canonical_count"] == 0
def test_wait_decision_is_classified_waiting():
item = project([row()], {"opp-1": decision("WAIT_PRODUCTION")})["items"][0]
assert item["operational_queue"] == "waiting"
assert item["status"] == "waiting"
def test_waiting_allowlist_does_not_hide_unknown_wait_action():
for code in ("WAIT_CUSTOMER", "WAIT_PRODUCTION"):
assert project([row()], {"opp-1": decision(code)})["items"][0]["operational_queue"] == "waiting"
unknown = project([row()], {"opp-1": decision("WAIT_MANUAL_REVIEW")})["items"][0]
assert unknown["operational_queue"] == "do_now"
assert unknown["current_action_code"] == "WAIT_MANUAL_REVIEW"
def test_blocker_is_only_current_action_and_downstream_task_is_evidence():
result = project([
row(id="downstream", action_code="SEND_INVOICE"),
row(id="blocker", action_code="VALIDATE_FISCAL_CUSTOMER"),
], {"opp-1": decision("VALIDATE_FISCAL_CUSTOMER", reason="Falta cliente fiscal.")})
assert [item["current_action_code"] for item in result["items"]] == ["VALIDATE_FISCAL_CUSTOMER"]
assert result["items"][0]["stale_task_refs"][0]["task_id"] == "downstream"
def test_linked_reconciliation_is_evidence_and_standalone_decision_remains_visible():
result = project([
row(source="reconciliation", id="loose", opportunity_id="", action_code="RECONCILE_DOCUMENTS"),
row(id="task-seed", action_code="RECONCILE_DOCUMENTS"),
], {"opp-1": decision("RECONCILE_DOCUMENTS")}, evidence_rows=[
row(source="reconciliation", id="linked", status="linked", action_code="RECONCILE_DOCUMENTS"),
])
assert result["canonical_count"] == 2
linked = next(item for item in result["items"] if item.get("opportunity_id"))
assert any(ref.get("reconciliation_item_id") == "linked" for ref in linked["source_refs"])
def test_historical_communication_does_not_seed_but_enriches_a_seeded_opportunity():
historical = row(source="communication", id="history", status="done", action_code="SEND_INFO")
evidence_only = project([], {"opp-1": decision()}, evidence_rows=[historical])
assert evidence_only["canonical_count"] == 0
assert evidence_only["evidence_source_count"] == 0
seeded = project([row(id="seed")], {"opp-1": decision()}, evidence_rows=[historical])
assert seeded["canonical_count"] == 1
assert seeded["evidence_source_count"] == 1
assert any(ref.get("communication_id") == "history" for ref in seeded["items"][0]["source_refs"])
def test_linked_and_resolved_reconciliation_do_not_seed_but_can_be_evidence():
linked = row(source="reconciliation", id="linked", status="linked", action_code="RECONCILE_DOCUMENTS")
resolved = row(source="reconciliation", id="resolved", status="resolved", action_code="RECONCILE_DOCUMENTS")
assert project([], {"opp-1": decision()}, evidence_rows=[linked])["canonical_count"] == 0
assert project([], {"opp-1": decision()}, evidence_rows=[resolved])["canonical_count"] == 0
result = project([row(id="seed")], {"opp-1": decision()}, evidence_rows=[linked, resolved])
assert result["canonical_count"] == 1
refs = result["items"][0]["source_refs"]
assert {ref.get("reconciliation_item_id") for ref in refs} >= {"linked", "resolved"}
def test_current_needs_review_communication_and_open_reconciliation_seed_work():
result = project([
row(source="communication", id="review", opportunity_id="", status="needs_review", action_code="REVIEW_MANUALLY"),
row(source="reconciliation", id="open", opportunity_id="", status="open", action_code="RECONCILE_DOCUMENTS"),
])
assert result["canonical_count"] == 2
assert result["seed_source_count"] == 2
def test_evidence_only_opportunity_ids_are_not_decision_candidates():
seeds = [row(id="seed", opportunity_id="seeded-opportunity")]
evidence = [row(source="communication", id="history", opportunity_id="historical-opportunity", status="done")]
assert opportunity_ids_from_work_seeds(seeds) == ["seeded-opportunity"]
assert "historical-opportunity" not in opportunity_ids_from_work_seeds(seeds)
result = project(seeds, {"seeded-opportunity": decision()}, evidence_rows=evidence)
assert result["canonical_count"] == 1
assert result["items"][0]["opportunity_id"] == "seeded-opportunity"
assert result["evidence_source_count"] == 0
def test_evidence_does_not_change_queue_counts_and_limit_is_post_canonical():
seeds = [row(id=str(index), opportunity_id=f"opp-{index}") for index in range(40)]
decisions = {f"opp-{index}": decision() for index in range(40)}
evidence = [
row(source="communication", id=f"history-{index}", opportunity_id=f"opp-{index}", status="done")
for index in range(40)
]
result = project(seeds, decisions, evidence_rows=evidence)
partition = partition_canonical_items(result["items"], display_limit=30)
assert result["seed_source_count"] == 40
assert result["evidence_source_count"] == 40
assert partition["work_queue_total"] == 40
assert len(partition["visible_actionable_items"]) == 30
def test_canonical_count_and_display_limit_follow_final_exact_key_merge():
rows = []
for index in range(35):
conversation_id = str(3000 + index)
rows.append(row(source="task", id=f"task-{index}", opportunity_id="", conversation_id=conversation_id, action_code="SEND_INFO"))
rows.append(row(source="communication", id=f"comm-{index}", opportunity_id="", conversation_id=conversation_id, action_code="SEND_INFO"))
result = project(rows, limit=30)
assert result["canonical_count"] == 35
assert result["visible_count"] == 30
assert len({item["work_item_key"] for item in result["all_items"]}) == 35
partition = partition_canonical_items(result["all_items"], display_limit=30)
assert partition["work_queue_total"] == 35
def test_every_canonical_output_has_unique_work_item_key():
result = project([
row(source="task", id="task", opportunity_id="", conversation_id="10", action_code="SEND_INFO"),
row(source="communication", id="comm", opportunity_id="", conversation_id="10", action_code="SEND_INFO"),
row(source="communication", id="review", opportunity_id="", conversation_id="10", action_code="REVIEW_MANUALLY"),
])
keys = [item["work_item_key"] for item in result["all_items"]]
assert len(keys) == len(set(keys))
assert len(keys) == 2
def test_pending_outbox_is_not_a_candidate_and_failed_exception_is_visible():
# SQL excludes pending rows; pure projection also preserves only what it is
# given, so assert the service query enforces that source boundary.
source = Path("app/operations_service.py").read_text(encoding="utf-8")
work_query = source.split("work_seed_rows = conn.execute", 1)[1].split("normalized_rows =", 1)[0]
assert "WHERE status IN ('failed','blocked')" in work_query
failed = row(source="outbox", id="failure", opportunity_id="", action_code="JASMIN_CREATE", status="failed")
assert project([failed])["items"][0]["operational_queue"] == "exception"
def test_matching_integration_failure_is_evidence_but_independent_failure_is_second_card():
matching = row(source="outbox", id="match", action_code="SEND_INVOICE", status="failed")
independent = row(source="outbox", id="other", action_code="JASMIN_TAX_FAILURE", status="failed")
result = project([row(action_code="SEND_INVOICE"), matching, independent], {"opp-1": decision("SEND_INVOICE")})
assert result["canonical_count"] == 2
opportunity_item = next(item for item in result["items"] if item["source"] == "opportunity")
assert any(ref.get("outbox_id") == "match" for ref in opportunity_item["source_refs"])
assert any(item["process_key"].startswith("integration-exception:") for item in result["items"])
def test_outbox_error_text_does_not_create_an_unstructured_action_match():
failure = row(
source="outbox", id="failure", action_code="JASMIN_FAILURE", status="failed",
title="SEND_INVOICE failed", detail="Error while processing SEND_INVOICE",
)
result = project([row(action_code="SEND_INVOICE"), failure], {"opp-1": decision("SEND_INVOICE")})
assert result["canonical_count"] == 2
opportunity_item = next(item for item in result["items"] if item["source"] == "opportunity")
assert not any(ref.get("outbox_id") == "failure" for ref in opportunity_item["source_refs"])
assert any(item["process_key"].startswith("integration-exception:") for item in result["items"])
def test_no_action_hides_resolved_evidence_but_preserves_explicit_review():
resolved = row(source="communication", id="resolved", status="done", action_code="SEND_INFO")
assert project([resolved], {"opp-1": decision("NO_ACTION")})["canonical_count"] == 0
needs_review = row(source="communication", id="review", status="needs_review", action_code="SEND_INFO")
review_result = project([needs_review], {"opp-1": decision("NO_ACTION")})
assert review_result["canonical_count"] == 1
assert review_result["items"][0]["current_action_code"] == "REVIEW"
assert "decisão central sem ação" in review_result["items"][0]["why_human_required"].lower()
def test_no_action_preserves_explicit_ambiguous_association_review():
ambiguous = row(
id="ambiguous", action_code="SEND_INVOICE",
opportunity_linking_status="ambiguous",
)
result = project([ambiguous], {"opp-1": decision("NO_ACTION")})
assert result["canonical_count"] == 1
assert result["items"][0]["current_action_code"] == "REVIEW"
def test_opportunity_customer_identity_wins_over_blank_communication_identity():
result = project([
row(source="communication", id="comm", customer_name="", contact_display_name=""),
row(id="task", fiscal_customer_name="CONSTRURECUP", customer_name=""),
], {"opp-1": decision()})
assert result["items"][0]["customer_name"] == "CONSTRURECUP"
assert result["items"][0]["contact_display_name"] == "CONSTRURECUP"
def test_display_limit_is_applied_after_canonicalization():
rows = []
decisions = {}
for index in range(40):
opportunity_id = f"opp-{index}"
decisions[opportunity_id] = decision()
rows.extend(row(id=f"{index}-{duplicate}", opportunity_id=opportunity_id) for duplicate in range(10))
result = project(rows, decisions, limit=30)
assert result["raw_source_count"] == 400
assert result["canonical_count"] == 40
assert result["visible_count"] == 30
def test_canonical_counts_and_keys_are_stable_and_action_sensitive():
rows = [row()]
first = project(rows, {"opp-1": decision("SEND_QUOTE")})
again = project(rows, {"opp-1": decision("SEND_QUOTE")})
changed = project(rows, {"opp-1": decision("CONFIRM_PAYMENT")})
assert first["canonical_count"] == len(first["all_items"]) == 1
assert first["items"][0]["process_key"] == again["items"][0]["process_key"] == "opportunity:opp-1"
assert first["items"][0]["work_item_key"] == again["items"][0]["work_item_key"]
assert first["items"][0]["work_item_key"] != changed["items"][0]["work_item_key"]
def test_projection_and_operations_get_path_contain_no_writes():
projection = Path("app/canonical_operations.py").read_text(encoding="utf-8").upper()
service = Path("app/operations_service.py").read_text(encoding="utf-8")
get_body = service.split("def get_operations_summary", 1)[1].split("def get_system_health_summary", 1)[0].upper()
next_actions = Path("app/opportunity_next_action_service.py").read_text(encoding="utf-8")
bulk_body = next_actions.split("def get_opportunity_next_actions", 1)[1]
for token in ("INSERT INTO", "UPDATE TASKS", "DELETE FROM", "ENSURE_PENDING_TASK"):
assert token not in projection
assert token not in get_body
assert "ensure_operation_schema()" not in bulk_body
def test_operations_query_separates_seed_eligibility_and_exposes_top_level_diagnostics():
source = Path("app/operations_service.py").read_text(encoding="utf-8")
seed_query = source.split("work_seed_rows = conn.execute", 1)[1].split("seed_opportunity_ids =", 1)[0]
evidence_query = source.split("evidence_rows = conn.execute", 1)[1].split("normalized_rows =", 1)[0]
assert "WHERE status IN ('new','classified','needs_review')" in seed_query
assert "OR opportunity_id IS NOT NULL" not in seed_query
assert "WHERE status IN ('open','needs_review','conflict')" in seed_query
assert "active_document_link_exclusion_sql" in source
assert "c.status NOT IN ('new','classified','needs_review')" in evidence_query
assert "ri.status IN ('linked','resolved')" in evidence_query
for key in (
"raw_source_count", "clean_source_count", "seed_source_count",
"evidence_source_count", "canonical_count", "visible_count",
"candidate_limit", "waiting_total",
):
assert f'"{key}"' in source