diff --git a/app/action_catalog.py b/app/action_catalog.py
index 937d81d..49cfccf 100644
--- a/app/action_catalog.py
+++ b/app/action_catalog.py
@@ -23,6 +23,7 @@ TRIAGE_ACTION_CODES = {
}
INTERNAL_ACTION_CODES = {
+ "CALL_CUSTOMER",
"FOLLOW_UP_QUOTE",
"FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT",
@@ -40,6 +41,13 @@ INTERNAL_ACTION_CODES = {
ACTION_CODES = TRIAGE_ACTION_CODES | INTERNAL_ACTION_CODES
ACTION_MAP = {
+ "CALL_CUSTOMER": {
+ "route": "vendas",
+ "action_required": True,
+ "action": "Ligar ao cliente",
+ "safe_to_post": False,
+ "business_event_on_done": "customer_called",
+ },
"SEND_QUOTE": {
"route": "vendas",
"action_required": True,
diff --git a/app/admin_dashboard.py b/app/admin_dashboard.py
index 8929130..37525c8 100644
--- a/app/admin_dashboard.py
+++ b/app/admin_dashboard.py
@@ -286,6 +286,7 @@ Obrigado pela sua mensagem.
Vamos analisar o pedido e responder assim que possível.
{closing}""")
ACTION_UI_LABELS = {
+ "CALL_CUSTOMER": "Ligar ao cliente",
"SEND_INFO": "Enviar informação",
"SEND_QUOTE": "Enviar orçamento",
"SEND_PROFORMA": "Enviar orçamento para pagamento",
diff --git a/app/admin_ui/labels.py b/app/admin_ui/labels.py
index 0f10630..1b208dd 100644
--- a/app/admin_ui/labels.py
+++ b/app/admin_ui/labels.py
@@ -8,6 +8,7 @@ from __future__ import annotations
from typing import Any
ACTION_LABELS = {
+ "CALL_CUSTOMER": "Ligar ao cliente",
"SEND_INFO": "Enviar informação",
"SEND_QUOTE": "Preparar orçamento",
"SEND_PROFORMA": "Enviar orçamento para pagamento",
@@ -37,6 +38,7 @@ ACTION_LABELS = {
}
PRIMARY_ACTION_LABELS = {
+ "CALL_CUSTOMER": "Ligar ao cliente",
"SEND_INFO": "Preparar resposta",
"SEND_QUOTE": "Preparar orçamento",
"SEND_PROFORMA": "Preparar orçamento para pagamento",
diff --git a/app/admin_ui/pages/opportunities.py b/app/admin_ui/pages/opportunities.py
index ad51e4d..6f93f52 100644
--- a/app/admin_ui/pages/opportunities.py
+++ b/app/admin_ui/pages/opportunities.py
@@ -18,6 +18,7 @@ from datetime import datetime, timezone
import app.admin_dashboard as legacy
from app.admin_dashboard import * # noqa: F401,F403
from app.admin_ui.labels import primary_action_label
+from app.timezone_utils import format_operator_datetime, parse_operator_local_datetime
from app.operation_noise import is_noise_operation_item
from app.opportunity_next_action_service import get_opportunity_next_action, get_opportunity_next_actions
from app.work_center_action_policy import (
@@ -1543,11 +1544,19 @@ def _parse_opportunity_dt(value: object):
def _opportunity_lifecycle_state(opp: dict) -> str:
state = str(opp.get("lifecycle_state") or "active").strip().lower() or "active"
now = datetime.now(timezone.utc)
+ pending_call = str(opp.get("pending_follow_up_action_code") or "").strip().upper() == "CALL_CUSTOMER"
+ pending_call_due = _parse_opportunity_dt(opp.get("pending_follow_up_due_at"))
+ if pending_call:
+ return "follow_up_due" if pending_call_due and pending_call_due <= now else "scheduled_follow_up"
+ # Compatibility only: a stale denormalized lifecycle marker cannot invent
+ # a scheduled call when no active CALL_CUSTOMER task exists.
+ if state == "scheduled_follow_up":
+ return "active"
nurture_until = _parse_opportunity_dt(opp.get("nurture_until"))
next_follow_up = _parse_opportunity_dt(opp.get("next_follow_up_at"))
if state == "nurture" and nurture_until and nurture_until <= now:
return "follow_up_due"
- if state in {"awaiting_customer", "active"} and next_follow_up and next_follow_up <= now:
+ if state in {"awaiting_customer", "active", "scheduled_follow_up"} and next_follow_up and next_follow_up <= now:
return "follow_up_due"
return state
@@ -1705,6 +1714,9 @@ def _opportunity_card_next_action(opp: dict) -> str:
return pending_follow_up_action or primary_action_label(pending_follow_up_code, fallback="Executar follow-up vencido")
if lifecycle_state == "nurture":
return pending_follow_up_action or "Aguardar data para retomar contacto"
+ if lifecycle_state == "scheduled_follow_up":
+ due = opp.get("pending_follow_up_due_at")
+ return f"Ligar ao cliente em {format_operator_datetime(due)}" if due else "Ligar ao cliente"
if lifecycle_state == "awaiting_customer":
due = _parse_opportunity_dt(opp.get("next_follow_up_at"))
return f"Aguardar resposta até {due.strftime('%d/%m')}" if due else "Aguardar resposta do cliente"
@@ -2041,13 +2053,13 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
lifecycle_task = next((
task for task in pending_tasks
if str(task.get("action_code") or "").upper() in {
- "CONFIRM_DELIVERY", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA",
+ "CALL_CUSTOMER", "CONFIRM_DELIVERY", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT", "FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_GENERIC",
"RECOVER_OPPORTUNITY", "REVIEW_NURTURE",
}
), None)
lifecycle_override = False
- if lifecycle_task and lifecycle_state_for_detail in {"follow_up_due", "recovery", "nurture"}:
+ if lifecycle_task and lifecycle_state_for_detail in {"follow_up_due", "scheduled_follow_up", "recovery", "nurture"}:
lifecycle_override = True
next_action = {
"action_code": str(lifecycle_task.get("action_code") or "FOLLOW_UP_GENERIC"),
@@ -2305,9 +2317,28 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
'''
+ scheduled_call_task = next((
+ task for task in pending_tasks
+ if str(task.get("action_code") or "").upper() == "CALL_CUSTOMER"
+ ), None)
+ scheduled_call_cancel_html = (
+ f'
'
+ if scheduled_call_task else ""
+ )
manual_follow_up_html = f'''
-
Criar follow-up
+
Agendar contacto
+
Agenda uma chamada sem alterar a fase comercial da oportunidade.
+
+ {scheduled_call_cancel_html}
+
Outros follow-ups
Agenda uma tarefa de follow-up. O sistema não cria confirmações de entrega automaticamente; usa “Verificar entrega” apenas quando o histórico sugere um problema real.
@@ -2335,11 +2366,12 @@ async def opportunity_detail_page(opportunity_id: str, notice: Optional[str] = N
lifecycle_state_label = lifecycle_label(lifecycle_state)
last_customer_dt = _parse_opportunity_dt(opportunity.get("last_customer_activity_at") or opportunity.get("last_message_at"))
last_operator_dt = _parse_opportunity_dt(opportunity.get("last_operator_activity_at"))
- next_follow_dt = _parse_opportunity_dt(opportunity.get("next_follow_up_at") or opportunity.get("nurture_until"))
+ next_follow_value = scheduled_call_task.get("due_at") if scheduled_call_task else (opportunity.get("next_follow_up_at") or opportunity.get("nurture_until"))
+ next_follow_dt = _parse_opportunity_dt(next_follow_value)
lifecycle_summary = [
f"Última atividade do cliente: {last_customer_dt.strftime('%d/%m/%Y %H:%M') if last_customer_dt else 'sem registo'}",
f"Último contacto do operador: {last_operator_dt.strftime('%d/%m/%Y %H:%M') if last_operator_dt else 'sem registo'}",
- f"Próximo contacto: {next_follow_dt.strftime('%d/%m/%Y') if next_follow_dt else 'não agendado'}",
+ f"Próximo contacto: {format_operator_datetime(next_follow_dt) if next_follow_dt else 'não agendado'}",
f"Tentativas: {int(opportunity.get('follow_up_attempts') or 0)}",
f"Entrega da última comunicação: {str(opportunity.get('last_delivery_status') or 'não verificada')}",
]
@@ -3016,26 +3048,45 @@ async def create_opportunity_follow_up_action(opportunity_id: str, request: Requ
form = await request.form()
follow_up_type = str(form.get("follow_up_type") or "generic").strip()
note = str(form.get("note") or "").strip()
+ due_at_text = str(form.get("due_at") or "").strip()
try:
delay_days = int(str(form.get("delay_days") or "3"))
except Exception:
delay_days = 3
try:
- from app.followup_service import create_manual_follow_up_for_opportunity
+ from app.followup_service import create_manual_follow_up_for_opportunity, schedule_customer_call
- result = create_manual_follow_up_for_opportunity(
- opportunity_id=opportunity_id,
- follow_up_type=follow_up_type,
- delay_days=delay_days,
- note=note,
- created_by="operator",
- )
- notice = "Follow-up agendado." if result.get("ok") else "Não foi possível agendar follow-up."
+ if follow_up_type == "customer_call":
+ if not due_at_text:
+ return PlainTextResponse("Data/hora obrigatória.", status_code=422)
+ due_at = parse_operator_local_datetime(due_at_text)
+ result = schedule_customer_call(
+ opportunity_id=opportunity_id, due_at=due_at, note=note, created_by="operator",
+ )
+ notice = "Contacto telefónico agendado." if result.get("ok") else "Não foi possível agendar o contacto."
+ else:
+ result = create_manual_follow_up_for_opportunity(
+ opportunity_id=opportunity_id, follow_up_type=follow_up_type,
+ delay_days=delay_days, note=note, created_by="operator",
+ )
+ notice = "Follow-up agendado." if result.get("ok") else "Não foi possível agendar follow-up."
except Exception as exc:
notice = f"Erro ao agendar follow-up: {exc}"
return RedirectResponse(f"/opportunities/{opportunity_id}?notice={quote(notice)}", status_code=303)
+@router.post("/opportunities/{opportunity_id}/scheduled-call/cancel")
+async def cancel_opportunity_scheduled_call_action(opportunity_id: str, request: Request):
+ if not is_uuid_text(opportunity_id):
+ return PlainTextResponse("Identificador de oportunidade inválido.", status_code=422)
+ from app.followup_service import cancel_scheduled_customer_call
+ result = cancel_scheduled_customer_call(
+ opportunity_id=opportunity_id, reason="Cancelado pelo operador", created_by="operator",
+ )
+ notice = "Contacto agendado cancelado." if result.get("cancelled") else "Não existe contacto agendado ativo."
+ return RedirectResponse(f"/opportunities/{opportunity_id}?notice={quote(notice)}", status_code=303)
+
+
@router.post("/opportunities/{opportunity_id}/commercial-terms")
async def update_opportunity_commercial_terms_action(opportunity_id: str, request: Request):
if not is_uuid_text(opportunity_id):
diff --git a/app/admin_ui/view_models/operations.py b/app/admin_ui/view_models/operations.py
index daefb60..55dcdd0 100644
--- a/app/admin_ui/view_models/operations.py
+++ b/app/admin_ui/view_models/operations.py
@@ -16,6 +16,8 @@ from app.admin_ui.guidance import item_requires_fiscal_customer, work_item_block
# Legacy label kept for regression context: Ambíguas.
OPERATION_FILTERS = [
("all", "A fazer", "/operations?scope=all", "/operations/partials/work-items?scope=all"),
+ ("waiting", "Aguardar", "/operations?scope=waiting", "/operations/partials/work-items?scope=waiting"),
+ ("backlog", "Backlog", "/operations?scope=backlog", "/operations/partials/work-items?scope=backlog"),
("bloqueadas", "Bloqueadas", "/operations?scope=bloqueadas", "/operations/partials/work-items?scope=bloqueadas"),
("ambiguas", "Associações por confirmar", "/operations?scope=ambiguas", "/operations/partials/work-items?scope=ambiguas"),
("atrasadas", "Atrasadas", "/operations?scope=atrasadas", "/operations/partials/work-items?scope=atrasadas"),
@@ -168,6 +170,8 @@ def is_scheduled_for_future(item: dict[str, Any]) -> bool:
def is_overdue(item: dict[str, Any]) -> bool:
+ if str(item.get("operational_queue") or "").lower() in {"waiting", "backlog", "not_current"}:
+ return False
due_dt = _parse_dt(item.get("due_at"))
if due_dt is not None:
return due_dt < datetime.now(timezone.utc)
@@ -241,6 +245,9 @@ def is_high_priority(item: dict[str, Any]) -> bool:
action_code = str(item.get("action_code") or "").upper()
priority = str(item.get("priority") or "").lower()
+ if str(item.get("operational_queue") or "").lower() in {"waiting", "backlog", "not_current"}:
+ return False
+
if is_ambiguous(item):
return True
if source == "outbox" and status in {"failed", "blocked", "stale"}:
@@ -312,7 +319,12 @@ def build_operations_view_model(scope: str = "all", *, limit: int = 30) -> dict[
counts = data.get("counts") or {}
scope = str(scope or "all").strip().lower()
work_items = list(data.get("work_items") or [])
- visible_items = [item for item in work_items if operation_item_matches_scope(item, scope)]
+ if scope == "waiting":
+ visible_items = list(data.get("waiting_items") or [])
+ elif scope == "backlog":
+ visible_items = list(data.get("backlog_items") or [])
+ else:
+ visible_items = [item for item in work_items if operation_item_matches_scope(item, scope)]
high_items = [item for item in visible_items if is_high_priority(item)]
normal_items = [item for item in visible_items if not is_high_priority(item)]
review_total = int(counts.get("review_tasks", 0) or 0) + int(counts.get("communications_needs_review", 0) or 0)
diff --git a/app/canonical_operations.py b/app/canonical_operations.py
index 28a835b..7bfa540 100644
--- a/app/canonical_operations.py
+++ b/app/canonical_operations.py
@@ -13,6 +13,7 @@ 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 = {
@@ -173,15 +174,22 @@ 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"]
+ 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),
}
@@ -313,6 +321,22 @@ def _opportunity_item(
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]] = []
@@ -338,7 +362,7 @@ def _opportunity_item(
"target_url": f"/opportunities/{opportunity_id}",
}
- matching_task = _matching_task(evidence_rows, action_code)
+ 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}"
@@ -455,6 +479,7 @@ def canonicalize_operations(
# 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.
diff --git a/app/followup_service.py b/app/followup_service.py
index 014bc03..2691d65 100644
--- a/app/followup_service.py
+++ b/app/followup_service.py
@@ -21,6 +21,7 @@ from app.db import engine
FOLLOW_UP_ACTION_CODES = {
+ "CALL_CUSTOMER",
"FOLLOW_UP_QUOTE",
"FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT",
@@ -1035,6 +1036,97 @@ def create_manual_follow_up_for_opportunity(
contact_purpose=str(step.get("purpose") or "manual_followup"),
)
+
+def schedule_customer_call(
+ *, opportunity_id: str, due_at: datetime, note: str = "", created_by: str = "operator",
+) -> Dict[str, Any]:
+ """Create or reschedule the authoritative active customer-call task.
+
+ ``lifecycle_state``/``next_follow_up_at`` are synchronized below only as
+ derived compatibility fields. Read paths must prefer the pending task.
+ """
+ if due_at.tzinfo is None:
+ due_at = due_at.replace(tzinfo=timezone.utc)
+ due_at = due_at.astimezone(timezone.utc)
+ opportunity = _load_opportunity_context(opportunity_id)
+ if not opportunity:
+ return {"ok": False, "status": "opportunity_not_found"}
+ with engine.begin() as conn:
+ existing = conn.execute(text("""
+ SELECT id::text FROM tasks
+ WHERE opportunity_id=CAST(:opportunity_id AS UUID)
+ AND status='pending' AND action_code='CALL_CUSTOMER'
+ ORDER BY created_at DESC LIMIT 1
+ FOR UPDATE
+ """), {"opportunity_id": opportunity_id}).mappings().first()
+ if existing:
+ task_id = str(existing["id"])
+ conn.execute(text("""
+ UPDATE tasks SET due_at=CAST(:due_at AS TIMESTAMPTZ),
+ note=COALESCE(NULLIF(:note, ''), note), updated_at=now(),
+ metadata=COALESCE(metadata, '{}'::jsonb) || jsonb_build_object(
+ 'scheduled_contact', true, 'contact_channel', 'phone',
+ 'rescheduled_at', now(), 'rescheduled_by', CAST(:created_by AS TEXT))
+ WHERE id=CAST(:task_id AS UUID)
+ """), {"task_id": task_id, "due_at": due_at.isoformat(), "note": note, "created_by": created_by})
+ status = "rescheduled"
+ else:
+ task_id = ""
+ status = "create"
+ if status == "create":
+ created = create_follow_up_task(
+ opportunity_id=opportunity_id, action_code="CALL_CUSTOMER", route="vendas",
+ action="Ligar ao cliente", note=note or "Contactar o cliente por telefone na data agendada.",
+ reason="CUSTOMER_REQUESTED_SCHEDULED_CALL", delay_days=0, created_by=created_by,
+ idempotency_suffix="scheduled-customer-call", follow_up_family="customer_call",
+ follow_up_stage=1, follow_up_max_stage=1, cascade=False,
+ contact_purpose="scheduled_customer_call", due_at_override=due_at,
+ )
+ if not created.get("ok"):
+ return created
+ task_id = str(created.get("task_id") or "")
+ status = str(created.get("status") or "created")
+ with engine.begin() as conn:
+ conn.execute(text("""
+ UPDATE opportunities SET lifecycle_state='scheduled_follow_up',
+ next_follow_up_at=CAST(:due_at AS TIMESTAMPTZ), updated_at=now()
+ WHERE id=CAST(:opportunity_id AS UUID)
+ """), {"opportunity_id": opportunity_id, "due_at": due_at.isoformat()})
+ if task_id:
+ conn.execute(text("""
+ INSERT INTO task_events(task_id,event_type,payload,created_by)
+ VALUES(CAST(:task_id AS UUID),'customer_call_scheduled',
+ jsonb_build_object('due_at',CAST(:due_at AS TEXT),'note',CAST(:note AS TEXT)),:created_by)
+ """), {"task_id": task_id, "due_at": due_at.isoformat(), "note": note, "created_by": created_by})
+ return {"ok": True, "status": status, "task_id": task_id, "opportunity_id": opportunity_id, "due_at": due_at.isoformat()}
+
+
+def cancel_scheduled_customer_call(
+ *, opportunity_id: str, reason: str = "Cancelado pelo operador", created_by: str = "operator",
+) -> Dict[str, Any]:
+ """Cancel only the explicit scheduled phone contact, preserving other tasks."""
+ with engine.begin() as conn:
+ rows = conn.execute(text("""
+ UPDATE tasks SET status='skipped', updated_at=now(),
+ metadata=COALESCE(metadata,'{}'::jsonb) || jsonb_build_object(
+ 'scheduled_contact_cancelled_at',now(),
+ 'scheduled_contact_cancelled_reason',CAST(:reason AS TEXT))
+ WHERE opportunity_id=CAST(:opportunity_id AS UUID)
+ AND status='pending' AND action_code='CALL_CUSTOMER'
+ RETURNING id::text
+ """), {"opportunity_id": opportunity_id, "reason": reason}).mappings().all()
+ for row in rows:
+ conn.execute(text("""
+ INSERT INTO task_events(task_id,event_type,payload,created_by)
+ VALUES(CAST(:task_id AS UUID),'customer_call_cancelled',
+ jsonb_build_object('reason',CAST(:reason AS TEXT)),:created_by)
+ """), {"task_id": row["id"], "reason": reason, "created_by": created_by})
+ conn.execute(text("""
+ UPDATE opportunities SET lifecycle_state='active', next_follow_up_at=NULL, updated_at=now()
+ WHERE id=CAST(:opportunity_id AS UUID) AND lifecycle_state='scheduled_follow_up'
+ """), {"opportunity_id": opportunity_id})
+ return {"ok": True, "status": "cancelled" if rows else "not_found", "cancelled": len(rows)}
+
def align_pending_followups_for_opportunity(
*,
opportunity_id: str,
diff --git a/app/operational_eligibility.py b/app/operational_eligibility.py
new file mode 100644
index 0000000..5798fb7
--- /dev/null
+++ b/app/operational_eligibility.py
@@ -0,0 +1,167 @@
+"""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
diff --git a/app/operations_service.py b/app/operations_service.py
index 0a0b3fc..8434826 100644
--- a/app/operations_service.py
+++ b/app/operations_service.py
@@ -443,6 +443,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(t.metadata, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
+ COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
+ o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM (
@@ -492,6 +494,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(io.payload, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
+ COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
+ o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM (
@@ -547,6 +551,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
ELSE '{}'::jsonb END AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
+ COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
+ o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM (
@@ -592,6 +598,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(ri.payload, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
+ COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
+ o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, ri.amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, ri.currency, 'EUR') AS opportunity_currency
FROM (
@@ -645,7 +653,29 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"candidate_limit": candidate_limit,
}).mappings().all()
- normalized_rows = _normalise_work_item_intent([dict(r) for r in work_seed_rows])
+ conversation_ids = sorted({
+ str(row.get("conversation_id") or "").strip()
+ for row in work_seed_rows
+ if str(row.get("conversation_id") or "").strip()
+ })
+ conversation_timeline = {}
+ if conversation_ids:
+ timeline_rows = conn.execute(text("""
+ SELECT conversation_id,
+ max(created_at) FILTER (WHERE direction = 'inbound') AS latest_public_inbound,
+ max(created_at) FILTER (WHERE direction = 'outbound') AS latest_public_outbound
+ FROM messages
+ WHERE source_system = 'chatwoot'
+ AND conversation_id = ANY(CAST(:conversation_ids AS TEXT[]))
+ GROUP BY conversation_id
+ """), {"conversation_ids": conversation_ids}).mappings().all()
+ conversation_timeline = {str(row["conversation_id"]): dict(row) for row in timeline_rows}
+
+ seed_rows = [dict(r) for r in work_seed_rows]
+ for row in seed_rows:
+ timeline = conversation_timeline.get(str(row.get("conversation_id") or ""), {})
+ row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")})
+ normalized_rows = _normalise_work_item_intent(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])
@@ -661,6 +691,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
# noise awaiting cleanup. The cleanup script still fixes the data source.
cleaned_counts["work_queue_total"] = partition["work_queue_total"]
cleaned_counts["waiting_total"] = partition["waiting_total"]
+ cleaned_counts["backlog_total"] = partition["backlog_total"]
+ cleaned_counts["not_current_total"] = partition["not_current_total"]
cleaned_counts["raw_source_count"] = projection["raw_source_count"]
cleaned_counts["canonical_total"] = len(canonical_items)
@@ -673,6 +705,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"canonical_count": len(canonical_items),
"visible_count": len(cleaned_work_items),
"waiting_total": len(all_waiting_items),
+ "backlog_total": partition["backlog_total"],
+ "not_current_total": partition["not_current_total"],
}
return {
"counts": cleaned_counts,
@@ -683,6 +717,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"recent_communications": [dict(r) for r in recent_communications],
"work_items": cleaned_work_items,
"waiting_items": waiting_items,
+ "backlog_items": partition["backlog_items"][:display_limit],
+ "not_current_items": partition["not_current_items"][:display_limit],
**diagnostics,
"projection_metrics": diagnostics,
}
diff --git a/app/opportunity_service.py b/app/opportunity_service.py
index cc415b2..960783f 100644
--- a/app/opportunity_service.py
+++ b/app/opportunity_service.py
@@ -183,6 +183,7 @@ STAGE_ON_TASK_DONE = {
_SCHEMA_READY = False
FOLLOW_UP_ACTION_CODES = {
+ "CALL_CUSTOMER",
"FOLLOW_UP_QUOTE",
"FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT",
@@ -198,6 +199,7 @@ OPPORTUNITY_LIFECYCLE_STATES = {
"active": "Ativa",
"awaiting_customer": "A aguardar cliente",
"follow_up_due": "Follow-up vencido",
+ "scheduled_follow_up": "Contacto agendado",
"recovery": "Recuperação",
"nurture": "Acompanhamento futuro",
}
@@ -1577,15 +1579,21 @@ def list_opportunities(
(
SELECT t.action_code FROM tasks t
WHERE t.opportunity_id = o.id AND t.status = 'pending'
- AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
+ AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
ORDER BY t.due_at NULLS LAST, t.created_at DESC LIMIT 1
) AS pending_follow_up_action_code,
(
SELECT t.action FROM tasks t
WHERE t.opportunity_id = o.id AND t.status = 'pending'
- AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
+ AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
ORDER BY t.due_at NULLS LAST, t.created_at DESC LIMIT 1
- ) AS pending_follow_up_action
+ ) AS pending_follow_up_action,
+ (
+ SELECT t.due_at FROM tasks t
+ WHERE t.opportunity_id = o.id AND t.status = 'pending'
+ AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
+ ORDER BY t.due_at NULLS LAST, t.created_at DESC LIMIT 1
+ ) AS pending_follow_up_due_at
FROM opportunities o
LEFT JOIN customers c ON c.id = o.local_customer_id
{where_sql}
@@ -1643,15 +1651,21 @@ def get_opportunity(opportunity_id: str) -> Optional[Dict[str, Any]]:
(
SELECT t.action_code FROM tasks t
WHERE t.opportunity_id = o.id AND t.status = 'pending'
- AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
+ AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
ORDER BY t.due_at NULLS LAST, t.created_at DESC LIMIT 1
) AS pending_follow_up_action_code,
(
SELECT t.action FROM tasks t
WHERE t.opportunity_id = o.id AND t.status = 'pending'
- AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
+ AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
ORDER BY t.due_at NULLS LAST, t.created_at DESC LIMIT 1
- ) AS pending_follow_up_action
+ ) AS pending_follow_up_action,
+ (
+ SELECT t.due_at FROM tasks t
+ WHERE t.opportunity_id = o.id AND t.status = 'pending'
+ AND (t.action_code LIKE 'FOLLOW_UP_%' OR t.action_code IN ('CALL_CUSTOMER','CONFIRM_DELIVERY','RECOVER_OPPORTUNITY','REVIEW_NURTURE'))
+ ORDER BY t.due_at NULLS LAST, t.created_at DESC LIMIT 1
+ ) AS pending_follow_up_due_at
FROM opportunities o
LEFT JOIN customers c ON c.id = o.local_customer_id
WHERE o.id = CAST(:opportunity_id AS UUID)
diff --git a/app/task_service.py b/app/task_service.py
index fc73abd..c0763a8 100644
--- a/app/task_service.py
+++ b/app/task_service.py
@@ -56,7 +56,7 @@ def _priority_for_action(action_code: str, route: str) -> str:
route = str(route or "").strip().lower()
if code in {"SEND_PROFORMA", "SEND_INVOICE", "CONFIRM_PAYMENT", "PREPARE_ORDER", "VALIDATE_PHYSICAL_ORDER", "CREATE_SHIPMENT"}:
return "alta"
- if code.startswith("FOLLOW_UP_") or code in {"CONFIRM_DELIVERY", "RECOVER_OPPORTUNITY", "REVIEW_NURTURE"}:
+ if code.startswith("FOLLOW_UP_") or code in {"CALL_CUSTOMER", "CONFIRM_DELIVERY", "RECOVER_OPPORTUNITY", "REVIEW_NURTURE"}:
return "normal"
if code in {"REVIEW_MANUALLY", "REMOVE_FROM_LIST", "MARK_NO_INTEREST", "IGNORE_SPAM", "NO_ACTION", "IGNORE_BOUNCE"}:
return "baixa"
diff --git a/app/timezone_utils.py b/app/timezone_utils.py
new file mode 100644
index 0000000..8d1811c
--- /dev/null
+++ b/app/timezone_utils.py
@@ -0,0 +1,31 @@
+"""Application civil-time helpers for Portuguese operators."""
+from __future__ import annotations
+
+from datetime import datetime, timezone
+from typing import Any
+from zoneinfo import ZoneInfo
+
+
+OPERATOR_TIMEZONE = ZoneInfo("Europe/Lisbon")
+
+
+def parse_operator_local_datetime(value: Any) -> datetime:
+ """Interpret a form datetime as Europe/Lisbon and return the UTC instant."""
+ parsed = value if isinstance(value, datetime) else datetime.fromisoformat(str(value).strip())
+ if parsed.tzinfo is None:
+ parsed = parsed.replace(tzinfo=OPERATOR_TIMEZONE)
+ return parsed.astimezone(timezone.utc)
+
+
+def to_operator_datetime(value: Any) -> datetime | None:
+ if not value:
+ return None
+ parsed = value if isinstance(value, datetime) else datetime.fromisoformat(str(value).replace("Z", "+00:00"))
+ if parsed.tzinfo is None:
+ parsed = parsed.replace(tzinfo=timezone.utc)
+ return parsed.astimezone(OPERATOR_TIMEZONE)
+
+
+def format_operator_datetime(value: Any, pattern: str = "%d/%m/%Y %H:%M") -> str:
+ parsed = to_operator_datetime(value)
+ return parsed.strftime(pattern) if parsed else ""
diff --git a/tests/test_operational_eligibility_p1b.py b/tests/test_operational_eligibility_p1b.py
new file mode 100644
index 0000000..adc173b
--- /dev/null
+++ b/tests/test_operational_eligibility_p1b.py
@@ -0,0 +1,347 @@
+from datetime import datetime, timedelta, timezone
+
+from app.operational_eligibility import evaluate_operational_eligibility
+from app.canonical_operations import canonicalize_operations, partition_canonical_items
+from app.timezone_utils import format_operator_datetime, parse_operator_local_datetime
+
+
+NOW = datetime(2026, 8, 15, 12, tzinfo=timezone.utc)
+
+
+def test_lisbon_summer_and_winter_schedule_conversion_and_display():
+ summer = parse_operator_local_datetime("2026-08-21T10:00")
+ winter = parse_operator_local_datetime("2026-12-21T10:00")
+ assert summer == datetime(2026, 8, 21, 9, tzinfo=timezone.utc)
+ assert winter == datetime(2026, 12, 21, 10, tzinfo=timezone.utc)
+ assert format_operator_datetime(summer) == "21/08/2026 10:00"
+ assert format_operator_datetime(winter) == "21/12/2026 10:00"
+
+
+def test_lisbon_scheduled_due_comparison_uses_same_utc_instant():
+ due = parse_operator_local_datetime("2026-08-21T10:00")
+ before = evaluate_operational_eligibility(item("CALL_CUSTOMER", refs=[task_ref("CALL_CUSTOMER")], due_at=due),
+ now=datetime(2026, 8, 21, 8, 59, tzinfo=timezone.utc))
+ at_due = evaluate_operational_eligibility(item("CALL_CUSTOMER", refs=[task_ref("CALL_CUSTOMER")], due_at=due),
+ now=datetime(2026, 8, 21, 9, 0, tzinfo=timezone.utc))
+ assert before.reason_code == "FOLLOW_UP_NOT_DUE"
+ assert at_due.reason_code == "FOLLOW_UP_DUE"
+
+
+def item(action, *, refs=(), due_at=None, **extra):
+ value = {
+ "current_action_code": action,
+ "action_code": action,
+ "source_refs": list(refs),
+ "due_at": due_at,
+ "operational_queue": "do_now",
+ }
+ value.update(extra)
+ return value
+
+
+def task_ref(action):
+ return {"source": "task", "id": f"task-{action}", "status": "pending", "action_code": action}
+
+
+def test_future_and_due_customer_call_are_waiting_then_do_now():
+ future = evaluate_operational_eligibility(
+ item("CALL_CUSTOMER", refs=[task_ref("CALL_CUSTOMER")], due_at=NOW + timedelta(days=1)), now=NOW,
+ )
+ assert (future.queue, future.reason_code, future.eligible) == ("waiting", "FOLLOW_UP_NOT_DUE", True)
+ due = evaluate_operational_eligibility(
+ item("CALL_CUSTOMER", refs=[task_ref("CALL_CUSTOMER")], due_at=NOW), now=NOW,
+ )
+ assert (due.queue, due.reason_code) == ("do_now", "FOLLOW_UP_DUE")
+
+
+def test_answered_standalone_info_and_support_are_not_current():
+ for action in ("SEND_INFO", "SUPPORT"):
+ result = evaluate_operational_eligibility(item(
+ action, latest_public_inbound=NOW - timedelta(hours=2),
+ latest_public_outbound=NOW - timedelta(hours=1),
+ ), now=NOW)
+ assert (result.queue, result.reason_code, result.eligible) == ("not_current", "ALREADY_ANSWERED", False)
+
+
+def test_newer_inbound_keeps_response_work_current():
+ result = evaluate_operational_eligibility(item(
+ "SEND_INFO", latest_public_inbound=NOW, latest_public_outbound=NOW - timedelta(hours=1),
+ ), now=NOW)
+ assert result.queue == "do_now"
+
+
+def test_review_and_no_interest_are_not_resolved_by_outbound():
+ for action in ("REVIEW_MANUALLY", "MARK_NO_INTEREST"):
+ result = evaluate_operational_eligibility(item(
+ action, latest_public_inbound=NOW - timedelta(hours=2),
+ latest_public_outbound=NOW - timedelta(hours=1),
+ ), now=NOW)
+ assert result.queue == "review"
+
+
+def test_explicit_response_task_is_not_hidden_by_timeline():
+ result = evaluate_operational_eligibility(item(
+ "SEND_INFO", refs=[task_ref("SEND_INFO")],
+ latest_public_inbound=NOW - timedelta(hours=2), latest_public_outbound=NOW,
+ ), now=NOW)
+ assert (result.queue, result.reason_code) == ("do_now", "EXPLICIT_PENDING_TASK")
+
+
+def test_fiscal_prerequisite_is_backlog_unless_structurally_blocking():
+ hygiene = evaluate_operational_eligibility(item(
+ "VALIDATE_FISCAL_CUSTOMER", decision={"financial_state": "no_document"},
+ ), now=NOW)
+ assert (hygiene.queue, hygiene.reason_code) == ("backlog", "FISCAL_DATA_NOT_CURRENTLY_BLOCKING")
+ blocked = evaluate_operational_eligibility(item(
+ "VALIDATE_FISCAL_CUSTOMER", refs=[task_ref("SEND_INVOICE")],
+ decision={"financial_state": "payment_confirmed"},
+ ), now=NOW)
+ assert (blocked.queue, blocked.reason_code) == ("do_now", "CURRENT_ACTION_BLOCKED_BY_FISCAL_IDENTITY")
+
+
+def test_reconciliation_hygiene_vs_current_blocker():
+ hygiene = evaluate_operational_eligibility(item(
+ "RECONCILE_DOCUMENTS",
+ refs=[{"source": "reconciliation", "id": "r1", "status": "open", "action_code": "RECONCILE_DOCUMENTS"}],
+ ), now=NOW)
+ assert hygiene.queue == "backlog"
+ blocker = evaluate_operational_eligibility(item(
+ "RECONCILE_DOCUMENTS", refs=[task_ref("SEND_INVOICE")],
+ ), now=NOW)
+ assert blocker.queue == "review"
+
+
+def reconciliation_ref(status="open", action="SEND_INVOICE"):
+ return {"source": "reconciliation", "id": f"rec-{action}", "status": status, "action_code": action}
+
+
+def test_reconciliation_only_open_invoice_and_proforma_are_backlog():
+ for action in ("SEND_INVOICE", "SEND_PROFORMA"):
+ result = evaluate_operational_eligibility(item(
+ action, refs=[reconciliation_ref("open", action)], source="reconciliation",
+ ), now=NOW)
+ assert (result.queue, result.reason_code, result.eligible) == (
+ "backlog", "DOCUMENT_RECONCILIATION_BACKLOG", False,
+ )
+
+
+def test_reconciliation_only_review_invoice_and_payment_are_review():
+ for action in ("SEND_INVOICE", "CONFIRM_PAYMENT"):
+ result = evaluate_operational_eligibility(item(
+ action, refs=[reconciliation_ref("needs_review", action)], source="reconciliation",
+ ), now=NOW)
+ assert (result.queue, result.reason_code) == (
+ "review", "RECONCILIATION_ASSOCIATION_REQUIRED",
+ )
+
+
+def test_explicit_task_or_opportunity_document_action_wins_over_reconciliation_backlog():
+ for action in ("SEND_INVOICE", "SEND_PROFORMA"):
+ task_backed = evaluate_operational_eligibility(item(
+ action, refs=[reconciliation_ref("open", action), task_ref(action)],
+ source="task", opportunity_id="opp-1",
+ ), now=NOW)
+ assert (task_backed.queue, task_backed.reason_code) == ("do_now", "EXPLICIT_PENDING_TASK")
+
+ opportunity_backed = evaluate_operational_eligibility(item(
+ action, refs=[reconciliation_ref("open", action)], source="opportunity",
+ opportunity_id="opp-1", current_downstream_obligation=True,
+ ), now=NOW)
+ assert opportunity_backed.queue == "do_now"
+
+
+def test_reconciliation_backlog_keeps_canonical_identity_unique():
+ source = {
+ "source": "reconciliation", "source_system": "jasmin", "id": "rec-1",
+ "status": "open", "action_code": "SEND_PROFORMA", "created_at": NOW,
+ "priority": "normal", "queue": "rever",
+ }
+ projected = canonicalize_operations([source, dict(source)], {})
+ assert len(projected["all_items"]) == 1
+ work = projected["all_items"][0]
+ assert work["work_item_key"] == "reconciliation:rec-1:action:SEND_PROFORMA"
+ assert (work["operational_queue"], work["eligibility_reason_code"]) == (
+ "backlog", "DOCUMENT_RECONCILIATION_BACKLOG",
+ )
+
+
+def test_waiting_and_backlog_never_become_immediate_by_age():
+ from app.admin_ui.view_models.operations import is_high_priority, is_overdue
+
+ old = item("CALL_CUSTOMER", due_at=NOW + timedelta(days=1), operational_queue="waiting",
+ created_at=NOW - timedelta(days=100), priority="alta")
+ assert is_overdue(old) is False
+ assert is_high_priority(old) is False
+
+
+def test_outbox_exception_remains_visible():
+ result = evaluate_operational_eligibility(item("JASMIN_SYNC", source="outbox", status="failed"), now=NOW)
+ assert (result.queue, result.reason_code, result.eligible) == ("exception", "INTEGRATION_FAILURE", True)
+
+
+def test_waiting_customer_without_explicit_task_is_not_immediate_work():
+ result = evaluate_operational_eligibility(item(
+ "CREATE_JASMIN_QUOTE", opportunity_lifecycle_state="awaiting_customer",
+ ), now=NOW)
+ assert (result.queue, result.reason_code) == ("waiting", "WAITING_CUSTOMER")
+
+
+def test_scheduled_call_overrides_theoretical_presentation_without_changing_process_identity():
+ task = {
+ "source": "task", "id": "task-call", "status": "pending",
+ "action_code": "CALL_CUSTOMER", "opportunity_id": "opp-1",
+ "due_at": datetime(2099, 1, 1, tzinfo=timezone.utc), "created_at": NOW - timedelta(days=2),
+ "priority": "normal", "queue": "vendas", "href": "/tasks/task-call",
+ "title": "Ligar ao cliente", "detail": "Conforme combinado",
+ }
+ projected = canonicalize_operations([task], {"opp-1": {
+ "action_code": "CREATE_JASMIN_QUOTE", "label": "Criar orçamento",
+ "reason": "Passo teórico", "priority": "normal",
+ }})
+ assert projected["canonical_count"] == 1
+ work = projected["items"][0]
+ assert work["process_key"] == "opportunity:opp-1"
+ assert work["work_item_key"] == "opportunity:opp-1:action:CALL_CUSTOMER"
+ assert work["operational_queue"] == "waiting"
+ partition = partition_canonical_items(projected["all_items"], display_limit=20)
+ assert partition["work_queue_total"] == 0
+ assert partition["waiting_total"] == 1
+
+
+def test_due_scheduled_call_appears_once_and_preserves_quote_stage_evidence():
+ task = {
+ "source": "task", "id": "task-call", "status": "pending",
+ "action_code": "CALL_CUSTOMER", "opportunity_id": "opp-1", "opportunity_stage": "QUOTE_SENT",
+ "due_at": datetime(2020, 1, 1, tzinfo=timezone.utc), "created_at": NOW - timedelta(days=2),
+ "priority": "normal", "queue": "vendas", "href": "/tasks/task-call",
+ }
+ projected = canonicalize_operations([task, dict(task)], {"opp-1": {"action_code": "CONFIRM_PAYMENT"}})
+ assert len(projected["all_items"]) == 1
+ assert projected["items"][0]["operational_queue"] == "do_now"
+ assert projected["items"][0]["opportunity_stage"] == "QUOTE_SENT"
+
+
+class _Result:
+ def __init__(self, rows=()):
+ self.rows = list(rows)
+
+ def mappings(self):
+ return self
+
+ def first(self):
+ return self.rows[0] if self.rows else None
+
+ def all(self):
+ return self.rows
+
+
+class _FollowupConnection:
+ def __init__(self, existing=True):
+ self.existing = existing
+ self.sql = []
+
+ def execute(self, statement, params=None):
+ sql = str(statement)
+ self.sql.append((sql, params or {}))
+ if "SELECT id::text FROM tasks" in sql:
+ return _Result([{"id": "00000000-0000-0000-0000-000000000010"}] if self.existing else [])
+ if "UPDATE tasks SET status='skipped'" in sql:
+ return _Result([{"id": "00000000-0000-0000-0000-000000000010"}])
+ return _Result()
+
+
+class _Context:
+ def __init__(self, conn):
+ self.conn = conn
+
+ def __enter__(self):
+ return self.conn
+
+ def __exit__(self, *_args):
+ return False
+
+
+class _Engine:
+ def __init__(self, conn):
+ self.conn = conn
+
+ def begin(self):
+ return _Context(self.conn)
+
+
+def test_reschedule_customer_call_updates_existing_without_duplicate(monkeypatch):
+ import app.followup_service as followups
+
+ conn = _FollowupConnection(existing=True)
+ monkeypatch.setattr(followups, "engine", _Engine(conn))
+ monkeypatch.setattr(followups, "_load_opportunity_context", lambda _id: {"id": _id, "stage": "QUOTE_SENT"})
+ monkeypatch.setattr(followups, "create_follow_up_task", lambda **kw: (_ for _ in ()).throw(AssertionError("duplicate created")))
+ result = followups.schedule_customer_call(
+ opportunity_id="00000000-0000-0000-0000-000000000001",
+ due_at=NOW + timedelta(days=2), note="Ligar de manhã",
+ )
+ assert result["status"] == "rescheduled"
+ assert sum("UPDATE tasks SET due_at" in sql for sql, _ in conn.sql) == 1
+ assert any("lifecycle_state='scheduled_follow_up'" in sql for sql, _ in conn.sql)
+
+
+def test_cancel_customer_call_removes_scheduled_condition(monkeypatch):
+ import app.followup_service as followups
+
+ conn = _FollowupConnection(existing=True)
+ monkeypatch.setattr(followups, "engine", _Engine(conn))
+ result = followups.cancel_scheduled_customer_call(
+ opportunity_id="00000000-0000-0000-0000-000000000001",
+ )
+ assert result["cancelled"] == 1
+ assert any("status='skipped'" in sql for sql, _ in conn.sql)
+ assert any("next_follow_up_at=NULL" in sql for sql, _ in conn.sql)
+
+
+def test_stale_scheduled_lifecycle_without_active_task_is_ignored():
+ from app.admin_ui.pages.opportunities import _opportunity_lifecycle_state
+
+ assert _opportunity_lifecycle_state({
+ "lifecycle_state": "scheduled_follow_up",
+ "next_follow_up_at": "2099-01-01T10:00:00+00:00",
+ "pending_follow_up_action_code": None,
+ }) == "active"
+
+
+def test_missing_cancelled_or_completed_call_never_produces_scheduled_work():
+ cases = [
+ [],
+ [{"source": "task", "id": "cancelled", "status": "cancelled", "action_code": "CALL_CUSTOMER"}],
+ [{"source": "task", "id": "done", "status": "done", "action_code": "CALL_CUSTOMER"}],
+ ]
+ for refs in cases:
+ result = evaluate_operational_eligibility(item(
+ "CALL_CUSTOMER", refs=refs, due_at=datetime(2099, 1, 1, tzinfo=timezone.utc),
+ opportunity_lifecycle_state="scheduled_follow_up",
+ ), now=NOW)
+ assert (result.eligible, result.queue, result.reason_code) == (
+ False, "not_current", "NO_ACTIVE_SCHEDULED_CALL",
+ )
+
+
+def test_scheduled_call_read_helpers_contain_no_repair_write():
+ from inspect import getsource
+ from app.admin_ui.pages.opportunities import _opportunity_lifecycle_state
+
+ source = getsource(_opportunity_lifecycle_state).upper()
+ assert "UPDATE " not in source and "INSERT " not in source and "DELETE " not in source
+
+
+def test_active_call_task_wins_when_lifecycle_state_is_stale():
+ from app.admin_ui.pages.opportunities import _opportunity_lifecycle_state
+
+ assert _opportunity_lifecycle_state({
+ "lifecycle_state": "active",
+ "pending_follow_up_action_code": "CALL_CUSTOMER",
+ "pending_follow_up_due_at": "2099-01-01T10:00:00+00:00",
+ }) == "scheduled_follow_up"
+ assert _opportunity_lifecycle_state({
+ "lifecycle_state": "active",
+ "pending_follow_up_action_code": "CALL_CUSTOMER",
+ "pending_follow_up_due_at": "2020-01-01T10:00:00+00:00",
+ }) == "follow_up_due"