656 lines
20 KiB
Python
656 lines
20 KiB
Python
import json
|
|
import os
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
from sqlalchemy import text
|
|
|
|
from app.db import engine
|
|
|
|
|
|
def _json(value: Any) -> str:
|
|
return json.dumps(value or {}, ensure_ascii=False)
|
|
|
|
|
|
def create_outbox_item(
|
|
*,
|
|
business_event_id: str,
|
|
target_system: str,
|
|
action_type: str,
|
|
payload: Dict[str, Any],
|
|
idempotency_key: str,
|
|
) -> Optional[str]:
|
|
sql = text("""
|
|
INSERT INTO integration_outbox (
|
|
business_event_id,
|
|
target_system,
|
|
action_type,
|
|
payload,
|
|
status,
|
|
retry_count,
|
|
idempotency_key
|
|
)
|
|
VALUES (
|
|
CAST(:business_event_id AS UUID),
|
|
:target_system,
|
|
:action_type,
|
|
CAST(:payload AS JSONB),
|
|
'pending',
|
|
0,
|
|
:idempotency_key
|
|
)
|
|
ON CONFLICT (idempotency_key) DO NOTHING
|
|
RETURNING id::text
|
|
""")
|
|
|
|
with engine.begin() as conn:
|
|
row = conn.execute(sql, {
|
|
"business_event_id": business_event_id,
|
|
"target_system": target_system,
|
|
"action_type": action_type,
|
|
"payload": _json(payload),
|
|
"idempotency_key": idempotency_key,
|
|
}).fetchone()
|
|
|
|
return row[0] if row else None
|
|
|
|
|
|
def build_outbox_specs(
|
|
*,
|
|
business_event_id: str,
|
|
event_type: str,
|
|
task: Dict[str, Any],
|
|
payload: Dict[str, Any],
|
|
) -> List[Dict[str, Any]]:
|
|
"""Constrói ações de outbox para integrações externas ativas.
|
|
|
|
A integração CRM externa antiga foi removida. O pipeline comercial passa a viver no
|
|
ClientFlow: oportunidades, produtos, financeiro e encomendas são entidades
|
|
internas. Esta função já não cria novos itens para o CRM externo antigo.
|
|
"""
|
|
base_payload = {
|
|
"business_event_id": business_event_id,
|
|
"event_type": event_type,
|
|
"task_id": task.get("id"),
|
|
"conversation_id": task.get("conversation_id"),
|
|
"contact_id": task.get("contact_id"),
|
|
"action_code": task.get("action_code"),
|
|
"route": task.get("route"),
|
|
"action": task.get("action"),
|
|
"note": task.get("note"),
|
|
"event_payload": payload or {},
|
|
}
|
|
|
|
if event_type == "invoice_sent":
|
|
return [
|
|
{
|
|
"target_system": "chatwoot",
|
|
"action_type": "add_private_note",
|
|
"payload": {**base_payload, "note": "Fatura marcada como enviada no ClientFlow."},
|
|
},
|
|
{
|
|
"target_system": "mautic",
|
|
"action_type": "add_tag",
|
|
"payload": {**base_payload, "tag": "invoice_sent"},
|
|
},
|
|
]
|
|
|
|
if event_type == "payment_confirmed":
|
|
return [
|
|
{
|
|
"target_system": "mautic",
|
|
"action_type": "add_tag",
|
|
"payload": {**base_payload, "tag": "payment_confirmed"},
|
|
},
|
|
{
|
|
"target_system": "mautic",
|
|
"action_type": "remove_tag",
|
|
"payload": {**base_payload, "tag": "proforma_unpaid"},
|
|
},
|
|
]
|
|
|
|
if event_type == "order_prepared":
|
|
return [
|
|
{
|
|
"target_system": "chatwoot",
|
|
"action_type": "add_private_note",
|
|
"payload": {**base_payload, "note": "Encomenda preparada no ClientFlow."},
|
|
},
|
|
]
|
|
|
|
if event_type == "shipment_validated":
|
|
return [
|
|
{
|
|
"target_system": "chatwoot",
|
|
"action_type": "add_private_note",
|
|
"payload": {**base_payload, "note": "Envio validado no ClientFlow."},
|
|
},
|
|
{
|
|
"target_system": "mautic",
|
|
"action_type": "add_tag",
|
|
"payload": {**base_payload, "tag": "shipment_validated"},
|
|
},
|
|
]
|
|
|
|
return []
|
|
|
|
def create_outbox_for_business_event(
|
|
*,
|
|
business_event_id: str,
|
|
event_type: str,
|
|
task: Dict[str, Any],
|
|
payload: Optional[Dict[str, Any]] = None,
|
|
) -> List[str]:
|
|
specs = build_outbox_specs(
|
|
business_event_id=business_event_id,
|
|
event_type=event_type,
|
|
task=task,
|
|
payload=payload or {},
|
|
)
|
|
|
|
created_ids: List[str] = []
|
|
|
|
for spec in specs:
|
|
idempotency_key = ":".join([
|
|
"outbox",
|
|
business_event_id,
|
|
spec["target_system"],
|
|
spec["action_type"],
|
|
])
|
|
|
|
outbox_id = create_outbox_item(
|
|
business_event_id=business_event_id,
|
|
target_system=spec["target_system"],
|
|
action_type=spec["action_type"],
|
|
payload=spec["payload"],
|
|
idempotency_key=idempotency_key,
|
|
)
|
|
|
|
if outbox_id:
|
|
created_ids.append(outbox_id)
|
|
|
|
return created_ids
|
|
|
|
|
|
|
|
def create_email_outbox_item(
|
|
*,
|
|
task_id: str,
|
|
opportunity_id: str = "",
|
|
to_email: str,
|
|
subject: str,
|
|
body: str,
|
|
draft_id: str = "",
|
|
selected_document_ids: Optional[List[str]] = None,
|
|
created_by: str = "operator",
|
|
) -> Optional[str]:
|
|
"""Create a safe email outbox item from an editable ClientFlow draft.
|
|
|
|
This does not send automatically. The existing outbox worker will only
|
|
process it if an email handler is explicitly implemented/enabled; until
|
|
then it remains visible for operator review instead of hiding the draft in
|
|
a task page.
|
|
"""
|
|
import hashlib
|
|
|
|
to_email = str(to_email or "").strip()
|
|
subject = str(subject or "").strip() or "Seguimento do processo"
|
|
body = str(body or "").strip()
|
|
if not to_email:
|
|
raise ValueError("Email do destinatário em falta.")
|
|
if not body:
|
|
raise ValueError("Mensagem vazia.")
|
|
|
|
fingerprint = hashlib.sha256((to_email + "\n" + subject + "\n" + body).encode("utf-8")).hexdigest()[:16]
|
|
idempotency_key = ":".join(["email_draft", str(task_id or ""), str(draft_id or "no_draft"), fingerprint])
|
|
payload = {
|
|
"channel": "email",
|
|
"mode": "operator_review_required",
|
|
"to": to_email,
|
|
"subject": subject,
|
|
"body": body,
|
|
"task_id": task_id,
|
|
"opportunity_id": opportunity_id or None,
|
|
"draft_id": draft_id or None,
|
|
"selected_document_ids": selected_document_ids or [],
|
|
"created_by": created_by,
|
|
"send_automatically": False,
|
|
}
|
|
return create_outbox_item(
|
|
business_event_id=task_id,
|
|
target_system="email",
|
|
action_type="send_email",
|
|
payload=payload,
|
|
idempotency_key=idempotency_key,
|
|
)
|
|
|
|
def list_outbox(
|
|
*,
|
|
status: Optional[str] = None,
|
|
target_system: Optional[str] = None,
|
|
limit: int = 100,
|
|
) -> List[Dict[str, Any]]:
|
|
where = []
|
|
params: Dict[str, Any] = {"limit": limit}
|
|
|
|
if status:
|
|
where.append("status = :status")
|
|
params["status"] = status
|
|
|
|
if target_system:
|
|
where.append("target_system = :target_system")
|
|
params["target_system"] = target_system
|
|
|
|
where_sql = ""
|
|
if where:
|
|
where_sql = "WHERE " + " AND ".join(where)
|
|
|
|
sql = text(f"""
|
|
SELECT
|
|
id::text,
|
|
business_event_id::text,
|
|
target_system,
|
|
action_type,
|
|
payload,
|
|
status,
|
|
retry_count,
|
|
idempotency_key,
|
|
last_error,
|
|
created_at,
|
|
updated_at,
|
|
sent_at,
|
|
locked_at,
|
|
lock_owner,
|
|
ignored_at
|
|
FROM integration_outbox
|
|
{where_sql}
|
|
ORDER BY created_at DESC
|
|
LIMIT :limit
|
|
""")
|
|
|
|
with engine.begin() as conn:
|
|
rows = conn.execute(sql, params).mappings().all()
|
|
|
|
return [dict(row) for row in rows]
|
|
|
|
|
|
def get_outbox_item(outbox_id: str) -> Optional[Dict[str, Any]]:
|
|
sql = text("""
|
|
SELECT
|
|
id::text,
|
|
business_event_id::text,
|
|
target_system,
|
|
action_type,
|
|
payload,
|
|
status,
|
|
retry_count,
|
|
idempotency_key,
|
|
last_error,
|
|
created_at,
|
|
updated_at,
|
|
sent_at,
|
|
locked_at,
|
|
lock_owner,
|
|
ignored_at
|
|
FROM integration_outbox
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
LIMIT 1
|
|
""")
|
|
|
|
with engine.begin() as conn:
|
|
row = conn.execute(sql, {"outbox_id": outbox_id}).mappings().first()
|
|
|
|
return dict(row) if row else None
|
|
|
|
|
|
def list_pending_outbox(
|
|
limit: int = 50,
|
|
target_system: Optional[str] = None,
|
|
) -> List[Dict[str, Any]]:
|
|
return list_outbox(
|
|
status="pending",
|
|
target_system=target_system,
|
|
limit=limit,
|
|
)
|
|
|
|
|
|
def claim_pending_outbox(
|
|
*,
|
|
limit: int = 50,
|
|
target_system: Optional[str] = None,
|
|
lock_owner: str = "worker",
|
|
) -> List[Dict[str, Any]]:
|
|
"""Claim pending outbox rows atomically for one worker.
|
|
|
|
This prevents two systemd timers/workers from processing the same pending
|
|
integration item at the same time. PostgreSQL SKIP LOCKED lets concurrent
|
|
workers take different rows without blocking each other.
|
|
|
|
Build the optional target filter explicitly instead of using
|
|
``(:target_system IS NULL OR target_system = :target_system)``. With
|
|
PostgreSQL + psycopg3, that expression can fail before execution with
|
|
``AmbiguousParameter: could not determine data type of parameter`` because
|
|
the bind parameter is used in an ``IS NULL`` expression.
|
|
"""
|
|
target_system = str(target_system).strip() if target_system else None
|
|
target_filter = "AND target_system = :target_system" if target_system else ""
|
|
sql = text(f"""
|
|
WITH picked AS (
|
|
SELECT id
|
|
FROM integration_outbox
|
|
WHERE status = 'pending'
|
|
{target_filter}
|
|
ORDER BY created_at ASC
|
|
FOR UPDATE SKIP LOCKED
|
|
LIMIT :limit
|
|
)
|
|
UPDATE integration_outbox io
|
|
SET
|
|
status = 'processing',
|
|
locked_at = now(),
|
|
lock_owner = :lock_owner,
|
|
updated_at = now(),
|
|
last_error = NULL
|
|
FROM picked
|
|
WHERE io.id = picked.id
|
|
RETURNING
|
|
io.id::text,
|
|
io.business_event_id::text,
|
|
io.target_system,
|
|
io.action_type,
|
|
io.payload,
|
|
io.status,
|
|
io.retry_count,
|
|
io.idempotency_key,
|
|
io.last_error,
|
|
io.created_at,
|
|
io.updated_at,
|
|
io.sent_at,
|
|
io.locked_at,
|
|
io.lock_owner
|
|
""")
|
|
params = {
|
|
"limit": int(limit),
|
|
"lock_owner": str(lock_owner or "worker")[:120],
|
|
}
|
|
if target_system:
|
|
params["target_system"] = target_system
|
|
|
|
with engine.begin() as conn:
|
|
rows = conn.execute(sql, params).mappings().all()
|
|
return [dict(row) for row in rows]
|
|
|
|
|
|
|
|
def outbox_stale_minutes() -> int:
|
|
"""Configured threshold for stuck processing rows."""
|
|
raw = os.getenv("OUTBOX_STALE_PROCESSING_MINUTES", "30").strip()
|
|
try:
|
|
return max(1, int(raw))
|
|
except ValueError:
|
|
return 30
|
|
|
|
|
|
def recover_stale_processing_outbox(
|
|
*,
|
|
stale_minutes: Optional[int] = None,
|
|
mode: Optional[str] = None,
|
|
limit: int = 100,
|
|
actor: str = "system",
|
|
) -> List[Dict[str, Any]]:
|
|
"""Recover or expose outbox items left in processing too long.
|
|
|
|
Modes:
|
|
- manual_only: mark rows as ``stale`` so the operator can decide;
|
|
- mark_failed: mark rows as ``failed`` with a stale-processing reason;
|
|
- retry_pending: return rows to ``pending`` so the worker retries them.
|
|
"""
|
|
stale_minutes = int(stale_minutes or outbox_stale_minutes())
|
|
mode = str(mode or os.getenv("OUTBOX_STALE_RECOVERY_MODE", "manual_only")).strip().lower()
|
|
if mode not in {"manual_only", "mark_failed", "retry_pending"}:
|
|
mode = "manual_only"
|
|
|
|
if mode == "retry_pending":
|
|
new_status = "pending"
|
|
retry_sql = "retry_count = retry_count + 1,"
|
|
error = f"Processing stale há mais de {stale_minutes} minutos; reposto para pending por {actor}."
|
|
elif mode == "mark_failed":
|
|
new_status = "failed"
|
|
retry_sql = "retry_count = retry_count + 1,"
|
|
error = f"Processing stale há mais de {stale_minutes} minutos; marcado failed por {actor}."
|
|
else:
|
|
new_status = "stale"
|
|
retry_sql = ""
|
|
error = f"Processing stale há mais de {stale_minutes} minutos; requer revisão manual."
|
|
|
|
sql = text(f"""
|
|
WITH picked AS (
|
|
SELECT id
|
|
FROM integration_outbox
|
|
WHERE status = 'processing'
|
|
AND locked_at IS NOT NULL
|
|
AND locked_at < now() - (:stale_minutes * interval '1 minute')
|
|
ORDER BY locked_at ASC
|
|
FOR UPDATE SKIP LOCKED
|
|
LIMIT :limit
|
|
)
|
|
UPDATE integration_outbox io
|
|
SET
|
|
status = :new_status,
|
|
{retry_sql}
|
|
locked_at = NULL,
|
|
lock_owner = NULL,
|
|
last_error = :error,
|
|
updated_at = now()
|
|
FROM picked
|
|
WHERE io.id = picked.id
|
|
RETURNING
|
|
io.id::text,
|
|
io.business_event_id::text,
|
|
io.target_system,
|
|
io.action_type,
|
|
io.payload,
|
|
io.status,
|
|
io.retry_count,
|
|
io.idempotency_key,
|
|
io.last_error,
|
|
io.created_at,
|
|
io.updated_at,
|
|
io.sent_at,
|
|
io.locked_at,
|
|
io.lock_owner,
|
|
io.ignored_at
|
|
""")
|
|
|
|
with engine.begin() as conn:
|
|
rows = conn.execute(sql, {
|
|
"stale_minutes": stale_minutes,
|
|
"limit": int(limit),
|
|
"new_status": new_status,
|
|
"error": error[:2000],
|
|
}).mappings().all()
|
|
|
|
recovered = [dict(row) for row in rows]
|
|
if recovered:
|
|
try:
|
|
from app.operator_audit_service import record_operator_action_best_effort
|
|
for item in recovered:
|
|
record_operator_action_best_effort(
|
|
action="outbox_stale_recovered",
|
|
entity_type="outbox",
|
|
entity_id=item.get("id"),
|
|
actor=actor,
|
|
payload={
|
|
"mode": mode,
|
|
"stale_minutes": stale_minutes,
|
|
"new_status": item.get("status"),
|
|
"target_system": item.get("target_system"),
|
|
"action_type": item.get("action_type"),
|
|
},
|
|
)
|
|
except Exception as exc:
|
|
print(f"ClientFlow stale outbox audit failed: {exc}", flush=True)
|
|
|
|
return recovered
|
|
|
|
def mark_outbox_sent(outbox_id: str) -> None:
|
|
with engine.begin() as conn:
|
|
conn.execute(text("""
|
|
UPDATE integration_outbox
|
|
SET
|
|
status = 'sent',
|
|
sent_at = now(),
|
|
locked_at = NULL,
|
|
lock_owner = NULL,
|
|
updated_at = now(),
|
|
last_error = NULL
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
"""), {"outbox_id": outbox_id})
|
|
|
|
|
|
def mark_outbox_failed(outbox_id: str, error: str) -> None:
|
|
with engine.begin() as conn:
|
|
conn.execute(text("""
|
|
UPDATE integration_outbox
|
|
SET
|
|
status = 'failed',
|
|
retry_count = retry_count + 1,
|
|
locked_at = NULL,
|
|
lock_owner = NULL,
|
|
last_error = :error,
|
|
updated_at = now()
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
"""), {
|
|
"outbox_id": outbox_id,
|
|
"error": error[:2000],
|
|
})
|
|
|
|
|
|
def set_outbox_status(
|
|
*,
|
|
outbox_id: str,
|
|
status: str,
|
|
error: str | None = None,
|
|
) -> None:
|
|
"""Atualiza o estado de um item da outbox com semântica operacional clara."""
|
|
allowed = {"pending", "processing", "sent", "failed", "blocked", "dry_run", "ignored", "cancelled", "stale"}
|
|
if status not in allowed:
|
|
raise ValueError(f"Estado inválido: {status}")
|
|
|
|
status_defaults = {
|
|
"failed": "Marcado manualmente como failed.",
|
|
"blocked": "Bloqueado por configuração ou pré-condição.",
|
|
"dry_run": "Validado em OUTBOX_DRY_RUN=true; nenhuma integração real foi executada.",
|
|
"ignored": "Ignorado manualmente.",
|
|
"cancelled": "Cancelado manualmente.",
|
|
"stale": "Processing preso; requer revisão ou reprocessamento manual.",
|
|
}
|
|
|
|
if status == "sent":
|
|
sql = text("""
|
|
UPDATE integration_outbox
|
|
SET
|
|
status = 'sent',
|
|
sent_at = COALESCE(sent_at, now()),
|
|
locked_at = NULL,
|
|
lock_owner = NULL,
|
|
last_error = NULL,
|
|
updated_at = now()
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
""")
|
|
params = {"outbox_id": outbox_id}
|
|
|
|
elif status == "pending":
|
|
sql = text("""
|
|
UPDATE integration_outbox
|
|
SET
|
|
status = 'pending',
|
|
sent_at = NULL,
|
|
locked_at = NULL,
|
|
lock_owner = NULL,
|
|
ignored_at = NULL,
|
|
last_error = NULL,
|
|
updated_at = now()
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
""")
|
|
params = {"outbox_id": outbox_id}
|
|
|
|
elif status == "processing":
|
|
sql = text("""
|
|
UPDATE integration_outbox
|
|
SET
|
|
status = 'processing',
|
|
locked_at = now(),
|
|
lock_owner = COALESCE(:error, 'worker'),
|
|
updated_at = now()
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
""")
|
|
params = {"outbox_id": outbox_id, "error": error}
|
|
|
|
else:
|
|
ignored_at_sql = "ignored_at = now()," if status == "ignored" else "ignored_at = ignored_at,"
|
|
retry_sql = "retry_count = retry_count + 1," if status == "failed" else ""
|
|
sql = text(f"""
|
|
UPDATE integration_outbox
|
|
SET
|
|
status = :status,
|
|
{retry_sql}
|
|
sent_at = NULL,
|
|
locked_at = NULL,
|
|
lock_owner = NULL,
|
|
{ignored_at_sql}
|
|
last_error = :error,
|
|
updated_at = now()
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
""")
|
|
params = {
|
|
"outbox_id": outbox_id,
|
|
"status": status,
|
|
"error": (error or status_defaults.get(status) or "Estado atualizado.")[:2000],
|
|
}
|
|
|
|
with engine.begin() as conn:
|
|
conn.execute(sql, params)
|
|
|
|
|
|
def mark_outbox_dry_run(outbox_id: str, message: str | None = None) -> None:
|
|
set_outbox_status(
|
|
outbox_id=outbox_id,
|
|
status="dry_run",
|
|
error=message or "OUTBOX_DRY_RUN=true; ação não executada na integração externa.",
|
|
)
|
|
|
|
|
|
def mark_outbox_blocked(outbox_id: str, message: str | None = None) -> None:
|
|
set_outbox_status(
|
|
outbox_id=outbox_id,
|
|
status="blocked",
|
|
error=message or "Integração desativada ou configuração incompleta.",
|
|
)
|
|
|
|
def get_outbox_item(outbox_id: str) -> Optional[Dict[str, Any]]:
|
|
sql = text("""
|
|
SELECT
|
|
id::text,
|
|
business_event_id::text,
|
|
target_system,
|
|
action_type,
|
|
payload,
|
|
status,
|
|
retry_count,
|
|
idempotency_key,
|
|
last_error,
|
|
created_at,
|
|
updated_at,
|
|
sent_at,
|
|
locked_at,
|
|
lock_owner,
|
|
ignored_at
|
|
FROM integration_outbox
|
|
WHERE id = CAST(:outbox_id AS UUID)
|
|
""")
|
|
|
|
with engine.begin() as conn:
|
|
row = conn.execute(sql, {"outbox_id": outbox_id}).mappings().first()
|
|
|
|
return dict(row) if row else None
|