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