Skip to content

4. Adaptadores de Mensageria

Visão geral

O sistema de mensageria permite que o bot converse em múltiplas plataformas usando uma interface unificada. Ele usa três design patterns:

flowchart TD
    subgraph "Design Patterns"
        direction LR
        FP["🏭 Factory Pattern<br/><i>AdapterFactory cria adapters</i>"]
        AP["🔌 Adapter Pattern<br/><i>Interface uniforme para todas as plataformas</i>"]
        SP["📐 Strategy Pattern<br/><i>Cada plataforma é uma estratégia</i>"]
    end

Estrutura de arquivos

runtime/messaging/
├── base.py               ← Interface abstrata (MessagingAdapter + IncomingMessage + OutgoingMessage)
├── factory.py            ← AdapterFactory (registro e resolução)
├── processor.py          ← ConversationProcessor (orquestra tudo — audio+file)
├── loop_guard.py         ← LoopGuard (anti-loop bot-to-bot, 4 camadas)
├── audio/
│   └── transcriber.py    ← Transcrição de áudio via Gemini 2.5 Flash (Vertex AI)
├── adapters/
│   ├── slack_adapter.py              ← Slack (Bolt SDK)
│   ├── whatsapp_evolution_adapter.py ← WhatsApp via Evolution API
│   ├── whatsapp_official_adapter.py  ← WhatsApp Business API (Meta)
│   ├── telegram_adapter.py           ← Telegram Bot API
│   ├── webchat_adapter.py            ← WebChat (HTTP REST)
│   └── sse_adapter.py                ← Server-Sent Events
└── sse/
    └── stream_manager.py             ← Gerencia streams SSE

A interface base: MessagingAdapter

Todos os adapters implementam esta interface:

classDiagram
    class MessagingAdapter {
        <<abstract>>
        +platform_name: str
        +webhook_path: str
        +_validate_config()
        +setup() async
        +parse_webhook_request(request) Optional~IncomingMessage~ async
        +send_message(message: OutgoingMessage) bool async
        +send_typing_indicator(channel_id, thread_id) async
        +update_message(message_id, new_text, channel_id) bool async
        +delete_message(message_id, channel_id) bool async
        +format_text(markdown_text) str
        +is_bot_message(message_data) bool async
        +cleanup() async
    }

    class IncomingMessage {
        +user_id: str
        +channel_id: str
        +thread_id: str
        +text: str
        +platform: str
        +message_type: MessageType
        +is_direct_message: bool
        +metadata: dict
    }

    class OutgoingMessage {
        +channel_id: str
        +thread_id: str
        +text: str
        +message_type: MessageType
        +metadata: dict
    }

    class MessageType {
        <<enum>>
        TEXT
        IMAGE
        FILE
        AUDIO
        VIDEO
    }

    MessagingAdapter ..> IncomingMessage : produz
    MessagingAdapter ..> OutgoingMessage : consome
    IncomingMessage *-- MessageType
    OutgoingMessage *-- MessageType

    class SlackAdapter { }
    class WhatsAppEvolutionAdapter { }
    class WhatsAppOfficialAdapter { }
    class TelegramAdapter { }
    class WebChatAdapter { }
    class SSEAdapter { }

    MessagingAdapter <|-- SlackAdapter
    MessagingAdapter <|-- WhatsAppEvolutionAdapter
    MessagingAdapter <|-- WhatsAppOfficialAdapter
    MessagingAdapter <|-- TelegramAdapter
    MessagingAdapter <|-- WebChatAdapter
    MessagingAdapter <|-- SSEAdapter

Métodos obrigatórios

Método Responsabilidade
parse_webhook_request() Recebe o request HTTP cru da plataforma e extrai: quem mandou, o que disse, token JWT, metadata
send_message() Envia a resposta do agente de volta para a plataforma
send_typing_indicator() Mostra indicador "digitando..." na plataforma
format_text() Converte Markdown genérico para formato da plataforma (ex: Slack usa *bold* em vez de **bold**)

Properties obrigatórias

Property Exemplo
platform_name "slack", "whatsapp", "webchat"
webhook_path "/slack/events", "/whatsapp/webhook"

Adapters disponíveis

flowchart LR
    subgraph "Produção ✅"
        SLACK["Slack<br/>/slack/events<br/><i>Bot Token + Signing Secret</i>"]
        WAEVO["WhatsApp Evolution<br/>/whatsapp/webhook<br/><i>Evolution API (self-hosted)</i>"]
        WAOFF["WhatsApp Official<br/>/whatsapp-official/webhook<br/><i>Meta Business API</i>"]
        WC["WebChat<br/>/webchat/message<br/><i>HTTP REST</i>"]
        SSE["SSE<br/>/sse/chat + /sse/stream/{id}<br/><i>Server-Sent Events</i>"]
    end

    subgraph "Exemplo ⚠️"
        TG["Telegram<br/>/telegram/webhook<br/><i>Bot API</i>"]
    end

Detalhes por adapter

Slack (SlackAdapter)

  • Usa Slack Bolt SDK
  • Responde em threads
  • Formata Markdown para Slack (mrkdwn)
  • Suporta atualização de mensagens (feedback visual)
  • Env vars: SLACK_BOT_TOKEN, SLACK_SIGNING_SECRET

WhatsApp Evolution (WhatsAppEvolutionAdapter)

  • Integra com Evolution API (instância self-hosted)
  • Suporta áudio (messageType:"audioMessage"/"ptt" → download via /message/getMedia)
  • Suporta documentos (messageType:"documentMessage" → download via Evolution API)
  • Suporta JWT Lookup via Redis (comando /ifriend-auth)
  • Env vars: WHATSAPP_EVOLUTION_API_URL, WHATSAPP_EVOLUTION_API_KEY, WHATSAPP_EVOLUTION_INSTANCE

WhatsApp Official (WhatsAppOfficialAdapter)

  • Integra com Meta WhatsApp Business API
  • Verifica webhook token no GET
  • Suporta áudio (type:"audio" → 2 chamadas: GET media info → GET download URL)
  • Suporta documentos (type:"document" → mesmo fluxo de download de media)
  • Envia/recebe via Graph API
  • Env vars: WHATSAPP_OFFICIAL_ACCESS_TOKEN, WHATSAPP_OFFICIAL_PHONE_NUMBER_ID, WHATSAPP_OFFICIAL_VERIFY_TOKEN

WebChat (WebChatAdapter)

  • REST puro — client envia POST, recebe resposta no body
  • Ideal para integração em websites
  • JWT via header Authorization: Bearer
  • Suporta áudio (audio_data + audio_mime_type base64 no JSON — ignorado se texto presente)
  • Suporta arquivos (file_data + file_mime_type + file_name base64 no JSON, max 20MB)
  • Env vars: WEBCHAT_ENABLED, WEBCHAT_ALLOWED_ORIGINS

SSE (SSEAdapter)

  • Processa em background, emite eventos em tempo real
  • Client faz POST em /sse/chat, depois abre stream em /sse/stream/{request_id}
  • Eventos: stream_start, tool_call_start, text_chunk, stream_complete
  • Env vars: SSE_ENABLED

Telegram (TelegramAdapter)

  • Implementação de exemplo/referência
  • Usa python-telegram-bot
  • Suporta áudio (message.voice ou message.audio — só baixa se não houver texto)
  • Suporta documentos (message.document — mantém texto junto)
  • Env vars: TELEGRAM_BOT_TOKEN

AdapterFactory

A factory gerencia o registro e resolução de adapters:

flowchart TD
    ENV["Env Vars:<br/>MESSAGING_PLATFORMS=slack,whatsapp<br/>SLACK_BOT_TOKEN=xoxb-...<br/>WHATSAPP_EVOLUTION_API_KEY=..."]

    ENV --> REG["AdapterFactory.register_from_env()"]

    REG --> DICT["Registry interno:<br/>/slack/events → SlackAdapter<br/>/whatsapp/webhook → WhatsAppEvolutionAdapter"]

    REQ["Request HTTP<br/>POST /slack/events"] --> RES["AdapterFactory.resolve(path)"]

    DICT --> RES
    RES --> ADAPTER["SlackAdapter instance"]

Como funciona: 1. No startup, register_from_env() lê as env vars e instancia os adapters configurados 2. Cada adapter registra seu webhook_path na factory 3. Quando chega um request, a factory resolve o adapter pelo path da URL 4. Se MESSAGING_PLATFORMS não está definido, detecta automaticamente pelos tokens presentes

ConversationProcessor

O ConversationProcessor é o orquestrador central. Ele conecta adapter → session → ADK:

flowchart TD
    IM[IncomingMessage] --> CP[ConversationProcessor]

    CP --> LG["0️⃣ LoopGuard<br/>(anti-loop bot-to-bot)"]

    LG --> TYPING["1️⃣ typing indicator<br/>(re-envio a cada 20s)"]

    TYPING --> AUDIO{"1b️⃣ É AUDIO<br/>e sem texto?"}
    AUDIO -->|sim| TRANS["Transcreve via<br/>Gemini 2.5 Flash<br/>(Vertex AI)"]
    TRANS --> REPLACE["transcript → message.text"]

    AUDIO -->|não| FILE{"2️⃣ Tem FILE<br/>no metadata?"}
    REPLACE --> FILE

    FILE -->|sim| PART["Part.from_bytes(file)<br/>+ Part.from_text(text)"]
    FILE -->|não| PARTT["Part.from_text(text)"]
    PART --> JWT
    PARTT --> JWT

    JWT["3️⃣ Decodifica JWT<br/>(se presente)"] --> SESSION["4️⃣ Obtém/cria Session<br/>(CloudSQL)"]
    SESSION --> STATE["5️⃣ Injeta jwt_context +<br/>message_metadata<br/>no session.state"]
    STATE --> ADK["6️⃣ runner.run_async()<br/>(ADK processa)"]

    ADK --> FEEDBACK["↩️ Feedback visual<br/>(atualiza msg com nome da tool)"]
    FEEDBACK --> ADK

    ADK --> EXTRACT["7️⃣ Extrai texto da resposta"]
    EXTRACT --> DELETE["🗑️ Deleta status message"]
    DELETE --> SEND["8️⃣ adapter.send_message()"]
    SEND --> USER[Resposta formatada para o usuário]

Feedback visual: durante o processamento, o Processor atualiza a mensagem de status na plataforma com o nome da tool sendo executada (ex: "🔍 Buscando produtos...").

Suporte a Áudio e Arquivos

Áudio (Transcrição via Gemini)

O pipeline de áudio usa Gemini 2.5 Flash (Vertex AI) em vez de GCP Speech-to-Text (~14x mais barato):

  1. Adapter detecta mensagem de áudio → MessageType.AUDIO + audio_data/audio_mime_type no metadata
  2. Processor (step 1b): se MessageType.AUDIO e sem texto → chama transcribe_audio(bytes, mime)
  3. Transcrição substitui message.text; processamento continua normalmente
  4. Se áudio E texto na mesma mensagem → áudio é ignorado (texto tem prioridade)

Arquivo: runtime/messaging/audio/transcriber.py
Max: 10MB (AUDIO_MAX_BYTES)
Modelo: gemini-2.5-flash (AUDIO_TRANSCRIPTION_MODEL)

Arquivos/Documentos

Arquivos são passados diretamente ao Gemini (que processa PDF, imagens, CSV, DOCX, TXT nativamente):

  1. Adapter detecta arquivo → MessageType.FILE + file_data/file_mime_type/file_name no metadata
  2. Processor (step 2): Part.from_bytes(data, mime_type) + Part.from_text(text)
  3. Se não há texto, prompt padrão: "Analise este arquivo."
  4. Diferente de áudio: arquivo mantém o texto (que é a pergunta sobre o arquivo)

MessageType Enum

Valor Metadata keys (IncomingMessage)
TEXT —
AUDIO audio_data (bytes), audio_mime_type (str)
FILE file_data (bytes), file_mime_type (str), file_name (str)
IMAGE (futuro)
VIDEO (futuro)

SSE: como funciona o streaming

O SSE é especial porque usa dois endpoints e processamento assíncrono:

sequenceDiagram
    actor Client as Frontend
    participant API as FastAPI
    participant SSE as SSEAdapter
    participant SM as StreamManager
    participant ADK as Runner ADK

    Client->>API: POST /sse/chat {message, jwt}
    API->>SSE: parse_webhook_request()
    SSE->>SM: create_stream(request_id)
    SSE-->>Client: {request_id: "abc123"}

    Note over SSE,ADK: Processamento em background (asyncio.create_task)
    SSE->>ADK: runner.run_async()

    Client->>API: GET /sse/stream/abc123
    API->>SM: subscribe(request_id)

    loop Eventos em tempo real
        ADK-->>SM: emit("tool_call_start", {tool: "busca_produtos"})
        SM-->>Client: event: tool_call_start\ndata: {...}

        ADK-->>SM: emit("text_chunk", {text: "Encontrei 3..."})
        SM-->>Client: event: text_chunk\ndata: {...}
    end

    ADK-->>SM: emit("stream_complete", {text: "resposta final"})
    SM-->>Client: event: stream_complete\ndata: {...}

Como criar um novo adapter

1. Crie o arquivo do adapter

# runtime/messaging/adapters/meu_adapter.py

from ..base import MessagingAdapter, IncomingMessage, OutgoingMessage

class MeuAdapter(MessagingAdapter):

    @property
    def platform_name(self) -> str:
        return "minha_plataforma"

    @property
    def webhook_path(self) -> str:
        return "/minha-plataforma/webhook"

    async def parse_webhook_request(self, request) -> IncomingMessage:
        body = await request.json()
        return IncomingMessage(
            user_id=body["sender_id"],
            channel_id=body["chat_id"],
            text=body["message"],
            platform=self.platform_name,
        )

    async def send_message(self, message: OutgoingMessage) -> bool:
        # Envia resposta via API da plataforma
        async with aiohttp.ClientSession() as session:
            await session.post(
                "https://api.minhaplataforma.com/send",
                json={"chat_id": message.channel_id, "text": message.text}
            )
        return True

    def format_text(self, text: str) -> str:
        # Converte markdown se necessário
        return text

2. Registre na factory

Em runtime/messaging/factory.py, adicione a detecção:

# Na função register_from_env()
if os.getenv("MINHA_PLATAFORMA_TOKEN"):
    adapter = MeuAdapter(token=os.getenv("MINHA_PLATAFORMA_TOKEN"))
    self.register(adapter)

3. Adicione a route no FastAPI

Em unified_bot.py:

@app.post("/minha-plataforma/webhook")
async def minha_plataforma_webhook(request: Request):
    adapter = adapter_factory.resolve("/minha-plataforma/webhook")
    message = await adapter.parse_webhook_request(request)
    response = await processor.process_message(message, adapter)
    return response

4. Adicione testes

Em tests/test_meu_adapter.py.

LoopGuard (Proteção Anti-Loop)

Proteção contra loops infinitos bot-to-bot, implementada em runtime/messaging/loop_guard.py com 4 camadas:

Camada Nome Trigger Resultado
1 Message ID Dedup message_id duplicado (webhook re-delivery) DUPLICATE — ignora silencioso
2 Circuit Breaker > N msgs em X segundos por usuário CIRCUIT_OPEN — cooldown msg 1x
3 Repetição Textual ≥ N msgs idênticas em janela Trip circuit breaker
4 Hard Turn Limit > N turnos por sessão TURN_LIMIT — escalação humana

Env vars: LOOP_GUARD_MAX_MESSAGES (5), LOOP_GUARD_WINDOW_SECONDS (10), LOOP_GUARD_COOLDOWN_SECONDS (30), LOOP_GUARD_MAX_TURNS_PER_SESSION (200), LOOP_GUARD_REPEAT_THRESHOLD (3), LOOP_GUARD_REPEAT_WINDOW_SECONDS (60)

Checklist para manutenção de adapters

  • [ ] Ao alterar parse_webhook_request(), valide que todos os campos do IncomingMessage estão preenchidos (incluindo message_type, metadata)
  • [ ] Ao adicionar suporte a áudio: use MessageType.AUDIO + audio_data/audio_mime_type no metadata
  • [ ] Ao adicionar suporte a arquivos: use MessageType.FILE + file_data/file_mime_type/file_name no metadata
  • [ ] Regra: se áudio E texto presentes na mesma mensagem → áudio é ignorado (use message.text check)
  • [ ] Regra: se arquivo E texto presentes → mantenha ambos (texto é pergunta sobre o arquivo)
  • [ ] Ao alterar send_message(), teste com mensagens longas (truncamento)
  • [ ] format_text() deve lidar com markdown, emojis e caracteres especiais
  • [ ] Verifique se send_typing_indicator() está implementado (UX)
  • [ ] Confira as env vars no cloudbuild.yaml ao adicionar novo adapter

Anterior: ← O Agente iFriend · Próximo: Sessão, Memória e Auth →