Files
clientflow_backend/app/operational_eligibility.py

168 lines
8.4 KiB
Python

"""Pure operational eligibility policy layered after process decisions."""
from __future__ import annotations
from dataclasses import asdict, dataclass
from datetime import datetime, timezone
from typing import Any, Iterable, Mapping
@dataclass(frozen=True)
class OperationalEligibility:
eligible: bool
queue: str
reason_code: str
reason_text: str
blocking_action_code: str | None = None
obligation_source_refs: tuple[dict[str, Any], ...] = ()
confidence: str = "high"
def to_dict(self) -> dict[str, Any]:
value = asdict(self)
value["obligation_source_refs"] = list(self.obligation_source_refs)
return value
RESPONSE_ACTIONS = {"SEND_INFO", "SUPPORT"}
DOCUMENT_ACTIONS = {"SEND_QUOTE", "SEND_PROFORMA", "SEND_INVOICE", "CREATE_JASMIN_QUOTE"}
CURRENT_QUEUES = {"do_now", "review", "blocked", "exception"}
def _s(value: Any) -> str:
return str(value or "").strip()
def _code(value: Any) -> str:
return _s(value).upper()
def _dt(value: Any) -> datetime | None:
if not value:
return None
if isinstance(value, datetime):
parsed = value
else:
try:
parsed = datetime.fromisoformat(_s(value).replace("Z", "+00:00"))
except ValueError:
return None
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
def _refs(item: Mapping[str, Any]) -> list[dict[str, Any]]:
return [dict(ref) for ref in item.get("source_refs") or [] if isinstance(ref, Mapping)]
def _task_refs(item: Mapping[str, Any]) -> list[dict[str, Any]]:
return [ref for ref in _refs(item) if ref.get("source") == "task" and _s(ref.get("status")).lower() == "pending"]
def _result(eligible: bool, queue: str, code: str, text: str, *, refs=(), blocker=None, confidence="high") -> OperationalEligibility:
return OperationalEligibility(eligible, queue, code, text, blocker, tuple(dict(ref) for ref in refs), confidence)
def evaluate_operational_eligibility(
item: Mapping[str, Any], *, now: datetime | None = None,
) -> OperationalEligibility:
"""Classify current work using structured state only; never mutates input."""
now = now or datetime.now(timezone.utc)
action = _code(item.get("current_action_code") or item.get("action_code"))
refs = _refs(item)
tasks = _task_refs(item)
task_codes = {_code(ref.get("action_code")) for ref in tasks}
decision = item.get("decision") if isinstance(item.get("decision"), Mapping) else {}
metadata = item.get("item_metadata") if isinstance(item.get("item_metadata"), Mapping) else {}
if _s(item.get("source")).lower() == "outbox" or _s(item.get("operational_queue")).lower() == "exception":
return _result(True, "exception", "INTEGRATION_FAILURE", "A integração falhou e requer intervenção.", refs=refs)
reconciliation_refs = [ref for ref in refs if ref.get("source") == "reconciliation"]
reconciliation_only = bool(reconciliation_refs) and len(reconciliation_refs) == len(refs)
structured_obligation = any(
value is True
for value in (
item.get("current_downstream_obligation"),
item.get("reconciliation_blocks_current_action"),
metadata.get("current_downstream_obligation"),
metadata.get("reconciliation_blocks_current_action"),
)
)
if reconciliation_only and not item.get("opportunity_id") and not tasks and not structured_obligation:
statuses = {_s(ref.get("status")).lower() for ref in reconciliation_refs}
if statuses & {"needs_review", "conflict"}:
return _result(
True, "review", "RECONCILIATION_ASSOCIATION_REQUIRED",
"A associação do documento exige confirmação humana.", refs=reconciliation_refs,
)
return _result(
False, "backlog", "DOCUMENT_RECONCILIATION_BACKLOG",
"Registo documental histórico sem obrigação comercial atual estruturada.",
refs=reconciliation_refs,
)
if action == "CALL_CUSTOMER":
active_call_refs = [ref for ref in tasks if _code(ref.get("action_code")) == "CALL_CUSTOMER"]
if not active_call_refs:
return _result(
False, "not_current", "NO_ACTIVE_SCHEDULED_CALL",
"Não existe uma tarefa CALL_CUSTOMER pendente ativa.", refs=refs,
)
due = _dt(item.get("due_at"))
if due and due > now:
return _result(True, "waiting", "FOLLOW_UP_NOT_DUE", "Contacto agendado; ainda não chegou a data.", refs=active_call_refs)
return _result(True, "do_now", "FOLLOW_UP_DUE", "Chegou a data agendada para contactar o cliente.", refs=active_call_refs)
if action in {"WAIT_PRODUCTION", "WAIT_CUSTOMER", "WAIT_PAYMENT", "WAIT_SUPPLIER", "WAIT_LOGISTICS", "WAIT_SCHEDULED_DATE"}:
return _result(True, "waiting", "WAITING_EXTERNAL_EVENT", "O processo aguarda um evento externo.", refs=refs)
if action in RESPONSE_ACTIONS and not tasks:
inbound = _dt(item.get("latest_public_inbound"))
outbound = _dt(item.get("latest_public_outbound"))
if inbound and outbound and outbound > inbound:
return _result(False, "not_current", "ALREADY_ANSWERED", "Existe resposta pública posterior ao último pedido do cliente.", refs=refs)
if action in {"REVIEW", "REVIEW_MANUALLY", "REVIEW_RECONSTRUCTED_PROCESS", "ASSOCIATE_OPPORTUNITY"}:
return _result(True, "review", "MANUAL_REVIEW_REQUIRED", _s(item.get("why_human_required")) or "A evidência exige decisão humana.", refs=refs)
if action == "MARK_NO_INTEREST":
return _result(True, "review", "COMMERCIAL_STATE_CONFIRMATION_REQUIRED", "Confirmar explicitamente a alteração do estado comercial.", refs=refs)
if action == "VALIDATE_FISCAL_CUSTOMER":
financial_state = _s(decision.get("financial_state")).lower()
blocking = bool(task_codes & (DOCUMENT_ACTIONS | {"VALIDATE_FISCAL_CUSTOMER"})) or financial_state == "payment_confirmed"
if blocking:
return _result(True, "do_now", "CURRENT_ACTION_BLOCKED_BY_FISCAL_IDENTITY", "A associação fiscal bloqueia uma operação documental atual.", refs=tasks, blocker="VALIDATE_FISCAL_CUSTOMER")
return _result(False, "backlog", "FISCAL_DATA_NOT_CURRENTLY_BLOCKING", "Os dados fiscais estão incompletos, mas não bloqueiam uma transição atual.", refs=refs)
if action == "RECONCILE_DOCUMENTS":
review_refs = [ref for ref in refs if ref.get("source") == "reconciliation" and _s(ref.get("status")).lower() in {"needs_review", "conflict"}]
blocking = bool(task_codes & (DOCUMENT_ACTIONS | {"CONFIRM_PAYMENT", "RECONCILE_DOCUMENTS"}))
if review_refs or blocking:
return _result(True, "review", "CURRENT_ACTION_BLOCKED_BY_DOCUMENT_LINK", "É necessário confirmar a ligação documental antes da transição atual.", refs=review_refs or tasks, blocker="RECONCILE_DOCUMENTS")
return _result(False, "backlog", "DOCUMENT_LINK_DATA_HYGIENE_ONLY", "A reconciliação melhora o histórico, mas não bloqueia trabalho atual.", refs=refs)
lifecycle = _s(item.get("opportunity_lifecycle_state")).lower()
if lifecycle in {"awaiting_customer", "nurture"} and not tasks:
return _result(True, "waiting", "WAITING_CUSTOMER", "A próxima iniciativa é esperada do cliente ou da data agendada.", refs=refs)
if action == "CONFIRM_PAYMENT" and not tasks:
proof_codes = {"COMPROVATIVO_PAGAMENTO", "PAYMENT_PROOF", "PAYMENT_RECEIVED"}
if not any(_code(ref.get("action_code")) in proof_codes for ref in refs):
return _result(True, "waiting", "WAITING_CUSTOMER_PAYMENT", "Aguardar pagamento ou comprovativo do cliente.", refs=refs)
if tasks:
return _result(True, "do_now", "EXPLICIT_PENDING_TASK", "Existe uma tarefa pendente explícita.", refs=tasks)
return _result(True, "do_now", "CURRENT_ACTION_DUE", "A evidência atual requer intervenção do operador.", refs=refs)
def apply_operational_eligibility(items: Iterable[Mapping[str, Any]], *, now: datetime | None = None) -> list[dict[str, Any]]:
projected = []
for source in items:
item = dict(source)
eligibility = evaluate_operational_eligibility(item, now=now)
item["eligibility"] = eligibility.to_dict()
item["operational_queue"] = eligibility.queue
item["why_human_required"] = eligibility.reason_text
item["eligibility_reason_code"] = eligibility.reason_code
item["eligible"] = eligibility.eligible
projected.append(item)
return projected