Agente Mediador via SSE
Migração de WebSocket para Server-Sent Events no front do operador — streaming de tokens, histórico em Redis Streams com dual-write em DynamoDB, e SQS FIFO como espinha dorsal assíncrona. ECS Fargate atrás de ALB interno, API Gateway com VPC Link para as rotas síncronas.
Visão geral
O desenho anterior usava um WebSocket com rotas connect/disconnect/init/message/likedislike, mas o "assíncrono" era ilusório: o front travava esperando resposta mesmo com todo o backend rodando atrás de SQS. Essa doc registra o redesenho completo — por que SSE substitui o WS aqui, onde a arquitetura original tinha corridas escondidas, e como cada peça (Redis, SQS FIFO, DynamoDB, ALB, API Gateway) se encaixa.
O cliente nunca precisa empurrar dado em tempo real pela mesma conexão que recebe — ele manda mensagem e decisão via POST comum. Só o sentido servidor→cliente precisa ser "always-on". Isso é exatamente o que SSE resolve nativamente, sem o overhead operacional de gerenciar conexão bidirecional (tabela de connection_id, $connect/$disconnect, post_to_connection).
Diagrama completo
Numeração 1–8 segue a execução real do fluxo de /message. Os passos 3 e 5 acontecem em paralelo — geração/streaming e persistência não dependem um do outro; nenhum bloqueia o outro.
Peça por peça
| Peça | Função |
|---|---|
| CloudFront | Único caminho pro /stream. Fica na frente do ALB especificamente para fugir do timeout de 29s do API Gateway em conexões long-lived. |
| API Gateway HTTP API + VPC Link | Único caminho pras rotas síncronas (/message, /warmup, /feedback, /stream/token) — todas respondem em milissegundos, então o cap de 29s nunca é um problema aqui. |
| ALB interno | Roteia pro target group do ECS. Não tem sticky session — é premissa de design que qualquer task pode atender qualquer request. |
| ECS Fargate (tasks) | Mesma imagem, mesmo código — o papel de "segura SSE" ou "processa POST" é só função de qual request cada task recebeu, não uma distinção de deployment. |
| ElastiCache Redis | Quatro estruturas de dado diferentes na mesma instância: Streams (histórico ordenado), Pub/Sub (fan-out entre tasks), string com TTL (warm:{id}), e dedupe key (SET NX) pra idempotência. |
| SQS FIFO | Espinha dorsal de persistência assíncrona. MessageGroupId = session_id garante ordem dentro da conversa sem sacrificar paralelismo entre sessões. |
| Worker | Consumer do SQS. Faz o dual-write condicional em Redis Stream + DynamoDB. Pode viver na mesma ECS service ou numa separada — desacoplado do handler HTTP de propósito. |
| DynamoDB | Fonte de verdade durável desde a primeira mensagem — não é mais "migração pós-TTL", é escrita em paralelo ao Redis. |
| LangGraph Mediator / Bedrock AgentCore Runtime | Invocado pela task que recebeu o POST, via astream_events — o generator que alimenta o Pub/Sub chunk a chunk. |
| Azure EntraID | Autoridade de identidade. Validado uma vez no token exchange, nunca diretamente no /stream. |
Auth & token exchange
EventSource nativo do browser não permite setar header Authorization — só query string ou cookie. Isso quebra o padrão de bearer herdado do Azure indo em header, que você já usa em todo o resto do sistema.
Solução: token exchange de curta duração. O front chama POST /stream/token com o bearer real do Azure no header (request normal, sem limitação nenhuma), o backend valida contra o EntraID com o mesmo padrão fail-closed do seu ACL atual, e devolve um JWT de ~60 segundos escopado pra aquele session_id. O front abre GET /stream?session_id=X&token=<curto> com esse token na query string.
@app.post("/stream/token")
async def issue_stream_token(request: Request, authorization: str = Header(None)):
claims = validate_azure_bearer(authorization)
now = int(time.time())
short_token = jwt.encode(
{
"session_id": claims["session_id"],
"client_id": claims["client_id"],
"operator_id": claims["operator_id"],
"scope": "stream",
"iat": now,
"exp": now + STREAM_TOKEN_TTL_S, # 60s
},
STREAM_SIGNING_KEY,
algorithm="HS256",
)
return {"token": short_token, "expires_in": STREAM_TOKEN_TTL_S}
A lib fetch-event-source permite header customizado direto, mas você perde o reconnect automático nativo do EventSource — teria que reimplementar backoff e Last-Event-ID na mão. O token curto na query string expira antes de qualquer replay de log valer a pena, então o ganho de segurança do header não compensa a perda de reconexão de graça.
O limite de 29s do API Gateway
HTTP API com integração privada via VPC Link tem timeout hard-capped entre 50ms e 29 segundos — não configurável pra cima. Só REST API regional/privada pode pedir quota increase, e mesmo assim o teto é 300s (5 minutos). Se a conexão SSE passar por API Gateway HTTP API, ela morre em 29s independente de atividade — inclusive durante um interrupt() de HITL, que por natureza pode ficar minutos em silêncio.
É por isso que o diagrama tem dois caminhos de entrada separados: /stream vai direto no ALB via CloudFront, nunca toca o API Gateway. As rotas síncronas (/message, /warmup, /feedback, /stream/token) respondem em milissegundos e continuam atrás do API Gateway + VPC Link normalmente — o cap de 29s nunca chega perto de ser um problema pra elas.
Fan-out entre tasks do ECS
Sem sticky session no ALB, a task que segura a conexão SSE de um session_id não é necessariamente a task que recebe o POST /message daquele mesmo session_id. Sem um mecanismo de fan-out, o evento gerado pelo agente nunca chega no cliente certo.
Redis Pub/Sub resolve isso desacoplando completamente onde a geração acontece de onde a conexão está aberta: a task que processa a mensagem publica em session:{session_id}:events; qualquer task que esteja segurando aquele SSE está subscrita no canal e repassa pro EventSource do front.
Fluxo de /message
| 1 | Task recebe POST /message, responde 202 imediatamente — nunca bloqueia o front. |
| 2 | Em paralelo: invoca o agente via astream_events, publicando cada chunk no Pub/Sub conforme gera. |
| 3 | Também em paralelo, manda a mensagem crua pro SQS FIFO — path de persistência totalmente desacoplado da geração. |
| 4 | Task com a conexão SSE recebe do Pub/Sub e escreve no stream HTTP como event: token / event: done. |
| 5 | Ao terminar, a resposta completa do assistente também vai pro SQS pra persistência. |
@app.post("/message")
async def post_message(request: Request):
body = await request.json()
session_id = body["session_id"]
prompt = body["prompt"]
message_id = str(uuid.uuid4()) # gerado na origem — chave de idempotência ponta a ponta
# dois caminhos em paralelo, nenhum bloqueia o outro
asyncio.create_task(_generate_and_stream(session_id, message_id, prompt))
await _enqueue_for_persistence(session_id, message_id, role="user", content=prompt)
return {"status": "accepted", "message_id": message_id}
Dual-write + idempotência
Dual-write elimina a corrida do TTL do Redis (migração pós-expiração pode perder dado se o job atrasar), mas introduz um problema novo: consistência entre dois writes independentes. Se o worker cair entre gravar no Redis e gravar no DynamoDB, e o SQS reentregar a mensagem, nenhum dos dois pode duplicar.
Resolvido com o mesmo padrão de idempotência que você já usa em produção — SET NX como dedupe key — mais ConditionExpression do lado do Dynamo:
def process_message(body):
session_id = body["session_id"]
message_id = body["message_id"]
dedupe_key = f"processed:{message_id}"
if not r.set(dedupe_key, "1", nx=True, ex=86400):
return # redelivery do SQS — já processado, sai sem duplicar
r.xadd(f"session:{session_id}:history", {
"message_id": message_id, "role": body["role"], "content": body["content"],
})
try:
table.put_item(
Item={
"PK": f"SESSION#{session_id}",
"SK": f"MSG#{body['ts']}#{message_id}",
"role": body["role"], "content": body["content"],
},
ConditionExpression="attribute_not_exists(SK)",
)
except ClientError as e:
if e.response["Error"]["Code"] != "ConditionalCheckFailedException":
raise # erro real — deixa o SQS re-tentar
Se cair entre XADD e put_item, o retry pula o XADD (dedupe key já setada) mas ainda tenta o put_item, seguro pela ConditionExpression. O inverso — Dynamo gravou, Redis não — não acontece nessa ordem de execução. Se quiser blindar de vez, inverte: grava Dynamo primeiro (fonte de verdade), Redis depois.
Schema do DynamoDB
SESSION#{session_id}MSG#{iso_timestamp}#{message_id}SK prefixado por timestamp garante ordenação lexicográfica nativa em Query, sem índice extra. Feedback é item próprio (SK: FEEDBACK#{message_id}) referenciando a mensagem original — nunca UpdateItem nela, pra manter Redis (append-only) e Dynamo simétricos e preservar trilha de auditoria de quando o feedback foi dado.
Ordem — por que FIFO e não standard
XADD sem ID customizado usa ordem de chegada no Redis, não ordem de criação da mensagem. Com SQS standard e múltiplos workers concorrentes, nada garante que o turno N seja processado antes do N+1 se ambos estiverem na fila ao mesmo tempo — você pode acabar com a resposta do assistente persistida antes da pergunta do usuário.
SQS FIFO com MessageGroupId = session_id garante ordem estrita dentro de uma sessão enquanto ainda permite paralelismo real entre sessões diferentes — cada session_id vira um grupo processado sequencialmente, grupos distintos escalam em paralelo entre os consumers. Só essa fila específica precisa ser FIFO; warmup e outras filas isoladas do sistema continuam standard.
sqs_client.send_message(
QueueUrl=SQS_QUEUE_URL,
MessageBody=json.dumps({...}),
MessageGroupId=session_id, # ordem garantida dentro da sessão
MessageDeduplicationId=message_id, # dedupe nativo do FIFO como camada extra
)
Warmup & a corrida escondida
Ordem correta: abre o SSE primeiro, dispara o warmup depois — nunca em paralelo. Mesmo assim, Pub/Sub é fire-and-forget: se o evento ready for publicado um milissegundo antes da task se inscrever no canal (deploy rolando, scale-out, latência de rede), a mensagem simplesmente deixa de existir pra qualquer um ler.
Além de publicar no Pub/Sub, grava SET warm:{session_id} 1 EX 300. A task do /stream, ao se inscrever, primeiro faz GET warm:{session_id} — se já for true, emite ready na hora sem esperar pub/sub nenhum. Cobre tanto "warmup chegou antes da inscrição" quanto "front reconectou depois de já estar warm".
Warmup não gera XADD nem write no Dynamo — é sinalização de runtime, não conteúdo de conversa. Fica fora do histórico.
Reconexão & heartbeat
- Idle timeout do ALB (60s default) e do CloudFront (origin response timeout, também perto de 30s por padrão): comentário
: keep-alive\n\na cada 15–20s independente de haver evento real — silêncio durante uminterrupt()longo não pode significar conexão morta. Last-Event-IDreal, não decorativo: usa o ID do Redis Stream comoid:de cada evento SSE. Em reconexão, o browser mandaLast-Event-IDautomaticamente — o servidor fazXRANGEa partir dali pra reenviar só o que perdeu, depois volta a assinar o pub/sub ao vivo.
if last_event_id:
missed = await redis_client.xrange(history_key, min=f"({last_event_id}", max="+")
for entry_id, fields in missed:
yield _format_sse(entry_id, fields.get("role", "assistant"), fields)
Deploy sem cortar sessão ativa
ECS rolling deploy manda SIGTERM antes de matar a task antiga. Sem tratamento, qualquer SSE aberto ali morre abrupto — o cliente só percebe pelo timeout de leitura, alguns segundos de silêncio percebidos como travamento.
def handle_sigterm(signum, frame):
logger.info("SIGTERM recebido, avisando %d conexões ativas", len(active_connections))
shutdown_event.set()
# respiro pro ALB parar de rotear novas requests antes de sair de vez
# alinhar com deregistration_delay do target group
No loop do /stream, quando shutdown_event está setado, emite event: reconnect explícito antes de fechar. O front, ao receber, chama o token exchange de novo e reabre — cai numa task saudável via ALB sem o usuário perceber gap.
Trace propagation através do SQS
SQS não propaga traceparent sozinho como HTTP propaga entre serviços via ADOT — precisa injetar manualmente no publish e extrair no consumer, senão o worker aparece como span órfão no X-Ray.
# publish
carrier = {}
propagate.inject(carrier)
sqs.send_message(
...,
MessageAttributes={"traceparent": {"DataType": "String", "StringValue": carrier["traceparent"]}},
)
# consumer
ctx = propagate.extract({"traceparent": msg["MessageAttributes"]["traceparent"]["StringValue"]})
with tracer.start_as_current_span("process_message", context=ctx):
...
session_id como atributo em todo span — SSE connection, POST handler, mensagem SQS, worker, write no Redis/Dynamo — é o fio que segue uma conversa inteira ponta a ponta no X-Ray mesmo cruzando o boundary assíncrono do SQS.
Tabela de decisões
| Decisão | Por quê |
|---|---|
| SSE em vez de WS | Canal servidor→cliente é o único que precisa ser always-on; comandos do cliente vão por POST comum. |
CloudFront + ALB direto pro /stream, fora do API Gateway | HTTP API + VPC Link tem timeout hard-capped em 29s — mata SSE de sessão longa. |
| Redis Pub/Sub pro fan-out | ALB não garante afinidade entre a task que segura o SSE e a que processa o POST da mesma sessão. |
| Redis Streams em vez de List | ID monotônico serve direto como Last-Event-ID; XRANGE permite replay sem reprocessar tudo. |
| Dual-write Redis + Dynamo | Elimina a corrida do TTL — job de "migração pós-expiração" pode perder dado se atrasar. |
message_id + SET NX + ConditionExpression | Idempotência ponta a ponta contra redelivery do SQS, sem duplicar em nenhum dos dois stores. |
SQS FIFO com MessageGroupId=session_id | SQS standard não garante ordem — histórico podia gravar resposta antes da pergunta. |
| Token exchange de curta duração | EventSource nativo não permite header Authorization customizado. |
Estado (warm:{id}) além de Pub/Sub no warmup | Pub/Sub é fire-and-forget — corrida entre publish e subscribe perde o evento sem estado de apoio. |
| Trace propagation manual no SQS | ADOT propaga automático em HTTP, não em mensageria — sem isso o worker vira span órfão. |
Handler completo
Código integral do handler — sse_handler.py — entregue junto com esta doc. Pontos de integração marcados com NotImplementedError são onde plugar o client real do LangGraph/Bedrock AgentCore Runtime e a validação do JWT do Azure contra o seu ACL existente.