Streaming & Persistência

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.

Premissa central

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.

Front (Browser) EventSource + fetch CloudFront distribution → ALB (/stream) API Gateway — HTTP API + VPC Link (rotas síncronas) Azure EntraID valida bearer / claims ALB interno target group: ECS Fargate ECS FARGATE SERVICE Task A conexão SSE ativa Task B processa POST /message ElastiCache Redis Streams (histórico) · Pub/Sub (fan-out) warm:{id} · dedupe SET NX SQS FIFO MessageGroupId = session_id LangGraph Mediator Bedrock AgentCore Runtime Worker consumer SQS FIFO DynamoDB histórico durável + feedback Todos os spans (Task A/B, Worker) carregam session_id + traceparent injetado manualmente no SQS → correlação ponta a ponta no X-Ray via ADOT GET /stream (abre conexão) eventos SSE (token / done) 1 POST /message /warmup /feedback /stream/token valida bearer VPC Link 2 3 PUBLISH tokens 4 SUBSCRIBE (fan-out entre tasks) 5 SendMessage FIFO invoke / astream_events 6 poll 7 XADD + SET NX (dedupe) 8 PutItem (Condition)
Linhas azuis = caminho crítico do fan-out (Pub/Sub entre Task B e Task A) — é a peça que existe especificamente porque o ALB não garante que a mesma task atenda a conexão SSE e o POST da mesma sessão.

Peça por peça

PeçaFunçã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 internoRoteia 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 RedisQuatro 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 FIFOEspinha dorsal de persistência assíncrona. MessageGroupId = session_id garante ordem dentro da conversa sem sacrificar paralelismo entre sessões.
WorkerConsumer 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.
DynamoDBFonte 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 RuntimeInvocado pela task que recebeu o POST, via astream_events — o generator que alimenta o Pub/Sub chunk a chunk.
Azure EntraIDAutoridade 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.

PYTHON
@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}
Por que não trocar pra fetch-event-source

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

Falha em produção se ignorado

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

1Task recebe POST /message, responde 202 imediatamente — nunca bloqueia o front.
2Em paralelo: invoca o agente via astream_events, publicando cada chunk no Pub/Sub conforme gera.
3Também em paralelo, manda a mensagem crua pro SQS FIFO — path de persistência totalmente desacoplado da geração.
4Task com a conexão SSE recebe do Pub/Sub e escreve no stream HTTP como event: token / event: done.
5Ao terminar, a resposta completa do assistente também vai pro SQS pra persistência.
PYTHON
@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:

PYTHON
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
Ordem dos writes importa

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

PK
SESSION#{session_id}
SK
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

Falha em produção se ignorado

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.

PYTHON
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.

Estado, não só notificação

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\n a cada 15–20s independente de haver evento real — silêncio durante um interrupt() longo não pode significar conexão morta.
  • Last-Event-ID real, não decorativo: usa o ID do Redis Stream como id: de cada evento SSE. Em reconexão, o browser manda Last-Event-ID automaticamente — o servidor faz XRANGE a partir dali pra reenviar só o que perdeu, depois volta a assinar o pub/sub ao vivo.
PYTHON
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.

PYTHON
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

Quebra a correlação automática

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.

PYTHON
# 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ãoPor quê
SSE em vez de WSCanal 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 GatewayHTTP API + VPC Link tem timeout hard-capped em 29s — mata SSE de sessão longa.
Redis Pub/Sub pro fan-outALB 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 ListID monotônico serve direto como Last-Event-ID; XRANGE permite replay sem reprocessar tudo.
Dual-write Redis + DynamoElimina a corrida do TTL — job de "migração pós-expiração" pode perder dado se atrasar.
message_id + SET NX + ConditionExpressionIdempotência ponta a ponta contra redelivery do SQS, sem duplicar em nenhum dos dois stores.
SQS FIFO com MessageGroupId=session_idSQS standard não garante ordem — histórico podia gravar resposta antes da pergunta.
Token exchange de curta duraçãoEventSource nativo não permite header Authorization customizado.
Estado (warm:{id}) além de Pub/Sub no warmupPub/Sub é fire-and-forget — corrida entre publish e subscribe perde o evento sem estado de apoio.
Trace propagation manual no SQSADOT 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.