From fabadbffdc9fffc95e52ff0230174ce5f571009b Mon Sep 17 00:00:00 2001 From: plx Date: Sat, 15 Aug 2026 00:32:38 +0000 Subject: [PATCH] fix: canonicalize operations work items --- app/canonical_operations.py | 474 ++++++++++++++++++ app/operations_service.py | 163 +++++- app/opportunity_next_action_service.py | 5 +- tests/test_canonical_operations_read_model.py | 415 +++++++++++++++ 4 files changed, 1041 insertions(+), 16 deletions(-) create mode 100644 app/canonical_operations.py create mode 100644 tests/test_canonical_operations_read_model.py diff --git a/app/canonical_operations.py b/app/canonical_operations.py new file mode 100644 index 0000000..28a835b --- /dev/null +++ b/app/canonical_operations.py @@ -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), + } diff --git a/app/operations_service.py b/app/operations_service.py index 9972f70..0a0b3fc 100644 --- a/app/operations_service.py +++ b/app/operations_service.py @@ -21,6 +21,13 @@ from app.work_center_action_policy import ( canonical_action_code, 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: @@ -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 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: counts = conn.execute(text(""" SELECT @@ -383,7 +395,7 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: LIMIT :limit """), {"limit": int(limit)}).mappings().all() - work_items = conn.execute(text(""" + work_seed_rows = conn.execute(text(""" SELECT * FROM ( SELECT 'task' AS source, t.id::text AS id, t.created_at, t.due_at, CASE @@ -433,14 +445,19 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: COALESCE(o.stage, '') AS opportunity_stage, COALESCE(o.value_amount, 0) AS opportunity_value_amount, 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 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 messages m ON m.id = t.message_id LEFT JOIN raw_events re ON re.id = t.raw_event_id - WHERE t.status = 'pending' - AND NOT (t.action_code LIKE 'FOLLOW_UP_%' AND t.due_at IS NOT NULL AND t.due_at > now()) + WHERE TRUE UNION ALL @@ -477,11 +494,16 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: COALESCE(o.stage, '') AS opportunity_stage, COALESCE(o.value_amount, 0) AS opportunity_value_amount, 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 customers cu ON cu.id = o.local_customer_id - WHERE io.status IN ('pending','failed','blocked') - AND NOT (io.status = 'failed' AND (COALESCE(io.last_error,'') ILIKE '%limpo manualmente%' OR COALESCE(io.last_error,'') ILIKE '%resolvido manualmente%')) + WHERE TRUE UNION ALL @@ -519,29 +541,139 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]: '/communications/' || c.id::text AS href, CASE WHEN c.customer_id IS NULL THEN 'Associar cliente' ELSE 'Abrir' END AS action_label, ''::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.stage, '') AS opportunity_stage, COALESCE(o.value_amount, 0) AS opportunity_value_amount, 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 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 ORDER BY 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 - LIMIT :limit - """), {"limit": int(limit)}).mappings().all() + """), {"candidate_limit": candidate_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()} # v4.9.0: the visible Operations total should match the queue the # operator can actually act on, not raw pending tasks that include mailbox # 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 { "counts": cleaned_counts, "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], "recent_communications": [dict(r) for r in recent_communications], "work_items": cleaned_work_items, + "waiting_items": waiting_items, + **diagnostics, + "projection_metrics": diagnostics, } diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index 653720b..e512250 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -302,8 +302,9 @@ def get_opportunity_next_actions(opportunity_ids: Iterable[str]) -> Dict[str, Di if not ids: return {} - from app.operation_service import ensure_operation_schema - ensure_operation_schema() + # Read path only: schema creation belongs to startup/migrations. In + # particular, GET /operations must never perform DDL while calculating its + # canonical projection. with engine.begin() as conn: opportunities = _bulk_rows(conn, """ SELECT id::text, stage, status, title, diff --git a/tests/test_canonical_operations_read_model.py b/tests/test_canonical_operations_read_model.py new file mode 100644 index 0000000..a44d6d6 --- /dev/null +++ b/tests/test_canonical_operations_read_model.py @@ -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