500 lines
22 KiB
Python
500 lines
22 KiB
Python
"""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
|
|
from app.operational_eligibility import apply_operational_eligibility
|
|
|
|
|
|
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_queues = {"do_now", "review", "blocked", "exception"}
|
|
actionable = [item for item in all_items if item.get("operational_queue") in actionable_queues]
|
|
waiting = [item for item in all_items if item.get("operational_queue") == "waiting"]
|
|
backlog = [item for item in all_items if item.get("operational_queue") == "backlog"]
|
|
not_current = [item for item in all_items if item.get("operational_queue") == "not_current"]
|
|
return {
|
|
"actionable_items": actionable,
|
|
"waiting_items": waiting,
|
|
"visible_actionable_items": actionable[:display_limit],
|
|
"visible_waiting_items": waiting[:display_limit],
|
|
"backlog_items": backlog,
|
|
"not_current_items": not_current,
|
|
"work_queue_total": len(actionable),
|
|
"waiting_total": len(waiting),
|
|
"backlog_total": len(backlog),
|
|
"not_current_total": len(not_current),
|
|
}
|
|
|
|
|
|
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]]]:
|
|
scheduled_call = next((
|
|
row for row in rows
|
|
if row.get("source") == "task"
|
|
and _s(row.get("status")).lower() == "pending"
|
|
and _canonical_code(row.get("action_code")) == "CALL_CUSTOMER"
|
|
), None)
|
|
if scheduled_call:
|
|
decision = {
|
|
**dict(decision),
|
|
"action_code": "CALL_CUSTOMER",
|
|
"label": "Ligar ao cliente",
|
|
"description": _s(scheduled_call.get("detail")) or "Contacto telefónico agendado pelo operador.",
|
|
"reason": "Existe um contacto telefónico explicitamente agendado.",
|
|
"priority": _s(scheduled_call.get("priority")) or "normal",
|
|
"target_url": _s(scheduled_call.get("href")),
|
|
}
|
|
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 = scheduled_call or _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)
|
|
items = apply_operational_eligibility(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),
|
|
}
|