Multi-Agent Production System · Observabilidade

Observabilidade Distribuída

Um trace_id único atravessando múltiplos agentes ECS — do request inicial até a última chamada A2A. OpenTelemetry com W3C TraceContext, spans por nó LangGraph, e X-Ray para visualização unificada.

OpenTelemetryAWS X-RayW3C TraceContext OTLP gRPCParentBasedLangGraph spans

O problema de visibilidade em multi-agent

Em um sistema single-agent, quando algo dá errado você vai nos logs do processo, encontra o erro, corrige. Em multi-agent com quatro agentes em ECS tasks diferentes, um request do usuário gera logs em quatro lugares distintos. Sem correlação entre eles, você sabe que algo deu errado mas não sabe onde na cadeia aconteceu, qual agente foi, ou qual nó específico do grafo LangGraph estava executando.

O problema se torna ainda mais sutil com falhas intermitentes. Se o agente de busca às vezes demora mais do que o esperado, você precisa saber: é o nó de busca em si? É a chamada HTTP para a API de hotéis? É a latência de rede entre os agentes? Sem spans hierárquicos, você vê o tempo total mas não a distribuição interna.

Sem observabilidade distribuída:
Usuário reclama: "minha reserva está demorando"
  └─► Log booking-agent: "request recebido às 14:00:00"
  └─► Log booking-agent: "request concluído às 14:00:28"
  └─► Onde foram os 28 segundos? Não dá para saber.

Com traces distribuídos:
Trace: 1-abc123 (28.3s total)
  └─► [booking-agent] POST /invoke          28.3s
       ├─► node.router                       0.1s
       ├─► node.search_hotels               12.4s  ← anomalia aqui
       │    └─► HTTP GET api.hotels.com      12.1s  ← API externa lenta
       └─► [search-agent] A2A               15.8s  ← mesmo trace_id!
            └─► node.search_flights          15.5s
                 └─► HTTP GET flights-api     15.2s

Stack e como as peças se conectam

O fluxo de dados é: código Python cria spans → OpenTelemetry SDK exporta via gRPC → X-Ray Daemon (sidecar) recebe e repassa → AWS X-Ray armazena e visualiza. A vantagem de ter o daemon como sidecar no mesmo ECS task é que não há configuração de rede — containers no mesmo task compartilham o namespace de rede e acessam localhost:4317 diretamente.

CamadaComponenteFunção
VisualizaçãoAWS X-RayService map, trace view, p99 por nó, filtros por atributo
ColetaX-Ray Daemon (sidecar)Recebe via OTLP gRPC :4317, repassa para X-Ray API
ExportaçãoOTLPSpanExporterEnvia spans em batch — mais eficiente que um por um
Spans de nó@traced_nodeCria span filho para cada nó LangGraph
Span raizFastAPIInstrumentorCria root span por request HTTP automaticamente
Propagação HTTPHTTPXClientInstrumentorInjeta traceparent em saídas httpx automaticamente

TracerProvider — configuração global

O TracerProvider é configurado uma vez no startup e registrado globalmente. Todos os módulos importam o tracer global — nenhum instancia seu próprio provider. Isso garante que todos os spans do processo compartilham o mesmo exporter, o mesmo sampler e os mesmos atributos de resource.

observability.py python
import os
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace.sampling import ParentBased, TraceIdRatioBased

SERVICE_NAME  = os.getenv("SERVICE_NAME",  "agent-service")
OTLP_ENDPOINT = os.getenv("OTLP_ENDPOINT", "http://localhost:4317")
SAMPLE_RATE   = float(os.getenv("TRACE_SAMPLE_RATE", "1.0"))


def setup_tracing() -> trace.Tracer:
    resource = Resource.create({
        "service.name":          SERVICE_NAME,
        "service.version":       os.getenv("APP_VERSION", "unknown"),
        "deployment.environment": os.getenv("ENVIRONMENT", "dev"),
    })

    # ParentBased: se o pai já foi sampleado, o filho segue — veja seção de sampling.
    sampler  = ParentBased(root=TraceIdRatioBased(SAMPLE_RATE))
    provider = TracerProvider(resource=resource, sampler=sampler)
    exporter = OTLPSpanExporter(endpoint=OTLP_ENDPOINT, insecure=True)

    # BatchSpanProcessor: acumula spans e envia em lotes — muito mais eficiente
    # que SimpleSpanProcessor que faz I/O para cada span individualmente.
    provider.add_span_processor(BatchSpanProcessor(exporter))
    trace.set_tracer_provider(provider)
    return trace.get_tracer(SERVICE_NAME)


tracer = setup_tracing()  # tracer global — importado pelos outros módulos

Instrumentação automática do FastAPI e HTTPX

Com duas linhas no startup, todos os requests HTTP recebidos e todas as chamadas HTTP enviadas via httpx passam a ter spans automáticos com o contexto de trace correto. Isso significa que chamadas A2A via httpx já propagam o traceparent sem nenhum código adicional por request.

app.py python
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
from opentelemetry.instrumentation.httpx   import HTTPXClientInstrumentor

app = FastAPI()

# FastAPIInstrumentor: cria root span para cada request.
# Se o request tiver traceparent no header, continua aquele trace.
# Se não tiver, cria um novo trace com trace_id fresco.
FastAPIInstrumentor.instrument_app(app)

# HTTPXClientInstrumentor: injeta traceparent em todas as saídas httpx.
# Isso cobre automaticamente chamadas A2A — sem código adicional por endpoint.
HTTPXClientInstrumentor().instrument()

Por que a instrumentação de nós é manual

O LangGraph não tem instrumentação nativa de OpenTelemetry. Não existe hook de "antes/depois de cada nó" que você possa usar sem modificar a biblioteca. A solução é um decorator que envolve cada função de nó — criando um span filho no contexto do span pai do request HTTP.

O contexto de span é herdado automaticamente em asyncio: quando o FastAPIInstrumentor cria um root span para o request, e dentro desse request o grafo LangGraph chama os nós, cada nó que usa tracer.start_as_current_span() herda o root span como pai. A hierarquia emerge do fluxo natural do código, sem precisar passar o contexto explicitamente entre funções.

@traced_node

observability.py — traced_node python
from functools import wraps
from opentelemetry.trace import SpanKind, Status, StatusCode


def traced_node(node_name: str, attributes: dict | None = None):
    def decorator(func):
        @wraps(func)
        async def wrapper(state: dict, config: dict, *args, **kwargs):
            thread_id = config.get("configurable", {}).get("thread_id", "unknown")

            # start_as_current_span: cria o span e o define como current no contexto.
            # Spans filhos criados dentro desta função herdarão este como pai.
            # kind=INTERNAL: não é chamada HTTP — é processamento interno do grafo.
            with tracer.start_as_current_span(
                f"langgraph.node.{node_name}",
                kind=SpanKind.INTERNAL,
            ) as span:
                span.set_attribute("langgraph.node",      node_name)
                span.set_attribute("langgraph.thread_id", thread_id)

                if attributes:
                    for k, v in attributes.items():
                        span.set_attribute(k, str(v))

                try:
                    result = await func(state, config, *args, **kwargs)
                    # Anota campos escalares do resultado — visíveis no trace do X-Ray.
                    # Útil para ver booking_id ou status diretamente no span.
                    if isinstance(result, dict):
                        for k, v in result.items():
                            if isinstance(v, (str, int, float, bool)):
                                span.set_attribute(f"langgraph.output.{k}", v)
                    span.set_status(Status(StatusCode.OK))
                    return result
                except Exception as exc:
                    # record_exception: captura stack trace no span.
                    # Visível no X-Ray como "exceptions" no span — filtráveis.
                    span.record_exception(exc)
                    span.set_status(Status(StatusCode.ERROR, str(exc)))
                    raise
        return wrapper
    return decorator

Ordem dos decorators importa

Quando você empilha @traced_node, @idempotent_node e outros decorators em um nó, a ordem define qual envolve qual. O decorator mais externo na declaração é o mais externo na execução — é ele que mede o tempo total.

@traced_node deve ficar por fora: assim o span captura o tempo total incluindo o check de idempotência. Se houver cache hit (o nó já executou antes), o span vai aparecer com duração próxima de zero — sinal visual imediato de que o resultado veio do cache.

nodes.py python
# Ordem de execução (de fora para dentro na declaração):
# traced_node → idempotent_node → função real
# O span cobre tudo — incluindo o check de cache do idempotent_node.

@traced_node("book_hotel", attributes={"agent.capability": "booking"})
@idempotent_node(redis, "book_hotel", key_fields=["hotel_id", "checkin"])
async def book_hotel_node(state: AgentState, config: dict) -> dict:
    return await with_retry(
        hotel_api.create_booking, hotel_id=state["hotel_id"],
        config=RETRY_CONFIGS["hotel_api"],
        circuit_breaker=circuit_registry.get("hotel_api"),
        service_name="hotel_api",
    )

W3C TraceContext — um trace_id para toda a cadeia

O W3C TraceContext é um padrão que define o formato do header traceparent para propagar contexto de trace entre serviços. O formato é 00-{trace_id}-{span_id}-{flags} onde trace_id é um identificador de 128 bits único para toda a cadeia de chamadas.

Quando o Agente A chama o Agente B com esse header, o Agente B cria seus spans como filhos do span do Agente A — mesmo trace_id, span_id diferente. No X-Ray, isso aparece como uma árvore unificada: você vê todos os spans de todos os agentes em uma única visualização hierárquica, como se fosse um único processo.

Por que W3C TraceContext e não o header nativo do X-Ray

O X-Ray tem seu próprio formato de header: X-Amzn-Trace-Id. O problema de usá-lo diretamente é o lock-in: se você quiser migrar de X-Ray para Grafana Tempo, Jaeger ou Zipkin no futuro, precisa mudar o código de propagação em todos os serviços. O W3C TraceContext é um padrão aberto suportado por todos esses sistemas.

O OpenTelemetry SDK com OTLPSpanExporter faz a tradução automaticamente: você escreve código W3C, e o X-Ray Daemon recebe o formato que o X-Ray entende. A migração de backend de observabilidade vira uma mudança de configuração, não de código.

inject() e extract() — a mecânica da propagação

O HTTPXClientInstrumentor faz o inject() automaticamente em saídas httpx. Mas em chamadas A2A onde você também precisa adicionar o JWT de autenticação e o budget de timeout, é útil entender o mecanismo para garantir que os headers não se sobrescrevam.

a2a_client.py python
import httpx
from opentelemetry.propagate import inject
from cascade_timeout import remaining_budget
from token_manager import token_manager


async def call_search_agent(endpoint: str, payload: dict) -> dict:
    token  = await token_manager.get_token()
    budget = remaining_budget()

    headers = {
        "Authorization":   f"Bearer {token}",
        "X-Timeout-Budget": str(budget),
        "Content-Type":    "application/json",
    }
    # inject() adiciona traceparent (e tracestate) ao dict de headers.
    # Modifica in-place — os headers do JWT e budget já estão presentes e
    # não são sobrescritos porque inject() só adiciona chaves de trace.
    inject(headers)

    async with httpx.AsyncClient(timeout=budget) as client:
        resp = await client.post(endpoint, json=payload, headers=headers)
        resp.raise_for_status()
        return resp.json()

Recebimento no Agente B

O FastAPIInstrumentor faz o extract() automaticamente — o span raiz do Agente B já herda o trace_id do Agente A sem código adicional. Para visualizar explicitamente o que acontece:

app.py — Agente B python
from opentelemetry.propagate import extract
from opentelemetry.trace import SpanKind

# FastAPIInstrumentor faz isso automaticamente.
# Mostrado aqui para clareza sobre o mecanismo interno:
@app.post("/search/invoke")
async def search_invoke(request: Request, body: SearchRequest,
                         claims=Depends(require_scope("search:invoke"))):
    # extract() lê o traceparent do header e cria um SpanContext com o
    # trace_id do Agente A. O span criado aqui é filho daquele span.
    ctx = extract(dict(request.headers))

    with tracer.start_as_current_span("search_agent.invoke",
                                       context=ctx, kind=SpanKind.SERVER) as span:
        span.set_attribute("agent.name",         "search-agent")
        span.set_attribute("langgraph.thread_id", body.thread_id)
        return await run_search_graph(body)

ParentBased — a decisão de sampling na cadeia A2A

Sampling em multi-agent tem um requisito que não existe em single-agent: consistência na cadeia. Se o Agente A decide samplear um trace (10% do tráfego), o Agente B também precisa samplear aquele trace — caso contrário você tem metade de um trace no X-Ray, o que é inútil para debugging.

O ParentBased resolve exatamente isso. Quando o Agente B recebe uma chamada com traceparent e o flag de sampling ativo, ele sempre samplea — independente do seu ratio local configurado. Se o pai decidiu samplear, todos os filhos na cadeia seguem. Se o pai decidiu não samplear, os filhos também não amostram.

AmbienteSamplerRatioComportamento
devTraceIdRatioBased1.0100% — toda chamada gera trace, máxima visibilidade
stagingParentBased(Ratio)0.550% com consistência na cadeia A2A
prod baixo volumeParentBased(Ratio)0.110% — custo de armazenamento aceitável
prod alto volumeErrorAwareSampler0.011% normal + 100% quando há erro

ErrorAwareSampler — sempre captura erros

Em produção de alto volume, samplear 1% é necessário para controlar custo de armazenamento. Mas você nunca quer perder traces de erros — são exatamente os que você precisa para debugging. O ErrorAwareSampler combina os dois: ratio baixo para o tráfego normal, 100% para traces com erro.

observability.py — ErrorAwareSampler python
from opentelemetry.sdk.trace.sampling import (
    Sampler, SamplingResult, Decision, TraceIdRatioBased
)


class ErrorAwareSampler(Sampler):
    def __init__(self, ratio: float):
        self._ratio = TraceIdRatioBased(ratio)

    def should_sample(self, parent_context, trace_id, name, kind, attributes, links):
        # Se vier de um pai já sampleado (cadeia A2A), segue a decisão do pai.
        # Garante consistência na cadeia — mesmo comportamento do ParentBased.
        parent_span = trace.get_current_span(parent_context)
        if parent_span and parent_span.is_recording():
            return SamplingResult(Decision.RECORD_AND_SAMPLE)

        # Atributo "error" definido antes do sampling — samplea sempre que há erro.
        if attributes and attributes.get("error"):
            return SamplingResult(Decision.RECORD_AND_SAMPLE)

        return self._ratio.should_sample(
            parent_context, trace_id, name, kind, attributes, links)

    def get_description(self) -> str:
        return "ErrorAwareSampler"

X-Ray Daemon como sidecar

O X-Ray Daemon roda como container sidecar no mesmo ECS task. Containers no mesmo task compartilham o namespace de rede — o agente acessa o daemon via localhost:4317 sem configuração de security group ou service discovery.

O sidecar é marcado como essential=false: se o daemon travar ou reiniciar, o container principal do agente não é afetado. Você perde alguns spans enquanto o daemon reinicia, mas o agente continua processando requests normalmente.

terraform/task_definition.tf — sidecar hcl
{
  name      = "xray-daemon"
  image     = "amazon/aws-xray-daemon:latest"
  essential = false  # falha do sidecar não derruba o agente
  cpu       = 32
  memory    = 256

  portMappings = [
    { containerPort = 4317, protocol = "tcp" },  # OTLP gRPC — recebe do agente
    { containerPort = 2000, protocol = "udp" },  # X-Ray nativo — SDKs legados
  ]
  # Sem hostPort: só acessível via localhost dentro do task
}

Logging estruturado para CloudWatch Logs Insights

Além de traces, o QueueMonitor emite métricas operacionais como JSON. O CloudWatch Logs Insights indexa campos JSON automaticamente — você cria filtros, dashboards e alarmes sem instrumentação adicional. O padrão é emitir uma linha de log por ciclo com todos os campos relevantes.

recovery_worker.py — QueueMonitor python
class QueueMonitor:
    async def run(self):
        while True:
            try:
                metrics = {
                    "recovery_queue_size":        await self.redis.llen("recovery:queue"),
                    "recovery_queue_utilization": await self._utilization(),
                    "dlq_size":                   await self.redis.llen("recovery:dlq"),
                    "threads_in_progress":        await self.redis.scard("threads:in_progress"),
                }
                # JSON no log — Logs Insights parseia e indexa cada campo.
                # Query: fields @timestamp, recovery_queue_size, dlq_size
                #        | filter ispresent(recovery_queue_size)
                #        | sort @timestamp desc
                logger.info("METRICS %s", json.dumps(metrics))

                if metrics["dlq_size"] > 0:
                    logger.error("DLQ contém %d item(ns). Intervenção necessária.",
                                 metrics["dlq_size"])
            except Exception:
                logger.exception("Erro no QueueMonitor")
            await asyncio.sleep(30)