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_typebase64 no JSON — ignorado se texto presente) - Suporta arquivos (
file_data+file_mime_type+file_namebase64 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.voiceoumessage.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):
- Adapter detecta mensagem de áudio →
MessageType.AUDIO+audio_data/audio_mime_typeno metadata - Processor (step 1b): se
MessageType.AUDIOe sem texto → chamatranscribe_audio(bytes, mime) - Transcrição substitui
message.text; processamento continua normalmente - 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):
- Adapter detecta arquivo →
MessageType.FILE+file_data/file_mime_type/file_nameno metadata - Processor (step 2):
Part.from_bytes(data, mime_type)+Part.from_text(text) - Se não há texto, prompt padrão: "Analise este arquivo."
- 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 doIncomingMessageestão preenchidos (incluindomessage_type,metadata) - [ ] Ao adicionar suporte a áudio: use
MessageType.AUDIO+audio_data/audio_mime_typeno metadata - [ ] Ao adicionar suporte a arquivos: use
MessageType.FILE+file_data/file_mime_type/file_nameno metadata - [ ] Regra: se áudio E texto presentes na mesma mensagem → áudio é ignorado (use
message.textcheck) - [ ] 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.yamlao adicionar novo adapter
Anterior: ← O Agente iFriend · Próximo: Sessão, Memória e Auth →