O motivo pelo qual a fatura chega a 60 milhões de wons quando agentes conversam entre si
TuBrief 편집팀
2026년 9월 13일
0
Computing/Software원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
커뮤니티의 다른 글
댓글 (0)
Log in to leave a comment
아직 작성된 글이 없습니다
원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
Log in to leave a comment
아직 작성된 글이 없습니다
Em um ambiente interno onde o bot de revisão de código da equipe de desenvolvimento e o agente de consultas da equipe de suporte estão integrados, o desastre que os engenheiros mais enfrentam não é a falta de inteligência do modelo. É uma única linha de mensagem de falha de parsing que entrou na fila assíncrona e um loop infinito em que dois agentes trocam perguntas e ficam rodando o fim de semana inteiro.
Tratar agentes como colegas autônomos só funciona em laboratórios de prompt. No momento em que você os coloca na infraestrutura interna, um agente é apenas um serviço distribuído que expele entradas não confiáveis. A abordagem de escrever "pare se não souber a resposta" dentro de um prompt em linguagem natural certamente falha em produção. Vamos examinar as linhas de controle físico que os engenheiros de plataforma precisam aplicar diretamente nos níveis de fila e gateway.
Se você permitir que os agentes comuniquem-se usando linguagem natural ou crases de Markdown misturadas (json), você perderá cinco ou seis horas por semana apenas rastreando erros de parsing. Devido à saída não determinística do modelo, falta uma única aspa ou o nome do campo muda ligeiramente, fazendo com que o consumidor downstream trave.
A solução é impor um esquema rigoroso na camada de transporte, assim como o rascunho do Google A2A e o JSON-RPC 2.0. Os payloads do Kafka transmitidos entre agentes devem obrigatoriamente passar pelos seguintes campos em tempo de execução.
| Nome do Campo | Tipo | Obrigatório | Propósito da Validação |
|---|---|---|---|
message_id |
UUIDv7 | Obrigatório | ID de mensagem global ordenável cronologicamente |
task_id |
UUIDv4 | Obrigatório | Unidade de rastreamento de tarefa de negócio única |
context_id |
String | Obrigatório | Identificador da sessão de conversas superior |
sender_id |
String | Obrigatório | Namespace do remetente (domínio:nome_do_agente) |
receiver_id |
String | Obrigatório | Namespace do destinatário (domínio:nome_do_agente) |
hop_count |
Integer | Obrigatório | Número acumulado de transferências entre agentes (valor inicial: 0) |
max_hops |
Integer | Obrigatório | Número máximo permitido de transferências (valor recomendado: 5) |
constraints |
Object | Opcional | Timeout, limite de orçamento de tokens |
data |
Object | Obrigatório | Dados de negócios estruturados |
Para evitar acidentes em que o processo do consumidor quebra, você deve colocar um interceptor de validação Pydantic no ponto de entrada e enviar imediatamente as mensagens que violam a especificação para uma Dead Letter Queue (DLQ).
`python
import uuid
from typing import Any, Dict
from pydantic import BaseModel, Field, ValidationError
from confluent_kafka import Consumer, Producer, KafkaError
class TaskConstraints(BaseModel):
timeout_ms: int = Field(default=30000, ge=1000, le=300000)
token_budget: int = Field(default=8000, ge=500, le=128000)
allow_delegation: bool = Field(default=True)
class A2AMessagePayload(BaseModel):
message_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
task_id: str = Field(..., description="ID de tarefa de negócio exclusivo")
context_id: str = Field(..., description="Identificador da sessão de transação")
sender_id: str = Field(..., pattern=r"^[a-z0-9_-]+:[a-z0-9_-]+")
hop_count: int = Field(default=0, ge=0)
max_hops: int = Field(default=5, ge=1, le=10)
constraints: TaskConstraints = Field(default_factory=TaskConstraints)
data: Dict[str, Any] = Field(..., description="Payload de negócios")
class ResilientAgentConsumer:
def init(self, kafka_conf: dict, main_topic: str, dlq_topic: str):
self.consumer = Consumer(kafka_conf)
self.producer = Producer({"bootstrap.servers": kafka_conf["bootstrap.servers"]})
self.main_topic = main_topic
self.dlq_topic = dlq_topic
self.consumer.subscribe([self.main_topic])
def route_to_dlq(self, raw_bytes: bytes, reason: str):
headers = [("dlq_error", reason.encode("utf-8")), ("origin_topic", self.main_topic.encode("utf-8"))]
self.producer.produce(topic=self.dlq_topic, value=raw_bytes, headers=headers)
self.producer.flush()
def process_events(self, dispatch_fn):
msg = self.consumer.poll(timeout=1.0)
if msg is None:
return
if msg.error():
if msg.error().code() != KafkaError._PARTITION_EOF:
self.route_to_dlq(msg.value() or b"", str(msg.error()))
return
try:
validated = A2AMessagePayload.model_validate_json(msg.value().decode("utf-8"))
except (ValidationError, UnicodeDecodeError) as err:
self.route_to_dlq(msg.value(), f"SCHEMA_VALIDATION_ERROR: {str(err)}")
self.consumer.commit(msg)
return
try:
dispatch_fn(validated)
self.consumer.commit(msg)
except Exception as exec_err:
self.route_to_dlq(msg.value(), f"EXECUTION_ERROR: {str(exec_err)}")
self.consumer.commit(msg)
`
Se você configurar esse padrão, o fenômeno de poison pill que trava o consumidor desaparecerá. Os recursos de engenharia gastos na depuração de parsing também podem ser recuperados imediatamente.
O momento mais perigoso em um pipeline de agentes é quando dois agentes entram em um loop infinito validando as saídas um do outro.
Há um caso real relatado em um relatório de pós-análise de março de 2026. Um bot gerador de SQL que não conseguiu resolver as restrições de chave estrangeira e um bot validador jogaram pingue-pongue gerando consultas diferentes sem parar, e como a ação em si era diferente a cada vez, o contador simples de tentativas de "mesma ação" de 50 vezes não funcionou. Os dois operaram sem isolamento por 11 dias (264 horas) e queimaram um total de 47.200 dólares (cerca de 63 milhões de wons) em custos de API. Quando o IAL-Scan analisou 6.549 repositórios de agentes open-source, 68 loops infinitos críticos foram descobertos exatamente assim em 47 projetos.
Você não deve confiar nas condições de escape do prompt. Um interruptor de parada físico deve ser aplicado no nível do API gateway.
| Critério de Controle | Método de Controle | Valor Recomendado | Regra de Operação |
|---|---|---|---|
| Limite de hops | Rastreamento da profundidade de chamadas no header | Máximo de 5 hops | Retorna 503 e para se delegarem um ao outro mais de 5 vezes |
| Teto de orçamento | Rastreamento de custo acumulado por transação no Redis | $10 por tarefa | Retorna 429 e bloqueia quando o gasto acumulado exceder 10 dólares |
| Limite de tentativas | Número de repetições para o mesmo ID de tarefa | Máximo de 3 vezes | Encerra se a mesma tarefa falhar 3 vezes, mesmo que o conteúdo mude |
Ao conectar X-Agent-Hop-Count e Redis no middleware do gateway, você pode controlar a explosão de gastos na faixa de $10.
`python
from fastapi import FastAPI, Request, Response, status
from fastapi.responses import JSONResponse
from starlette.middleware.base import BaseHTTPMiddleware
import redis.asyncio as redis
import logging
logger = logging.getLogger("AgentCircuitBreaker")
class AgentGovernanceMiddleware(BaseHTTPMiddleware):
def init(self, app: FastAPI, redis_pool: redis.Redis, max_hops: int = 5, cost_limit_usd: float = 10.0):
super().init(app)
self.redis = redis_pool
self.max_hops = max_hops
self.cost_limit_usd = cost_limit_usd
self.token_cost_ratio = 0.000015 # Cálculo baseado em $0.015 por 1.000 tokens
async def dispatch(self, request: Request, call_next) -> Response:
trace_id = request.headers.get("X-Trace-ID") or request.headers.get("traceparent", "trace-root")
current_hops = int(request.headers.get("X-Agent-Hop-Count", "0"))
if current_hops >= self.max_hops:
logger.error(f"Bloqueio por excesso de limite de hops: trace_id={trace_id}, hops={current_hops}")
return JSONResponse(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
content={"error": "CIRCUIT_BREAKER_HOP_LIMIT_EXCEEDED", "trace_id": trace_id}
)
cost_key = f"governance:cost:{trace_id}"
spent_cost_raw = await self.redis.get(cost_key)
accumulated_cost = float(spent_cost_raw.decode("utf-8")) if spent_cost_raw else 0.0
if accumulated_cost >= self.cost_limit_usd:
logger.error(f"Bloqueio por excesso de orçamento: trace_id={trace_id}, spent=${accumulated_cost}")
return JSONResponse(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
content={"error": "CIRCUIT_BREAKER_BUDGET_EXHAUSTED", "trace_id": trace_id}
)
custom_headers = dict(request.scope["headers"])
custom_headers[b"x-agent-hop-count"] = str(current_hops + 1).encode("utf-8")
request.scope["headers"] = list(custom_headers.items())
response = await call_next(request)
consumed_tokens_hdr = response.headers.get("X-LLM-Tokens-Consumed")
if consumed_tokens_hdr:
incremental_cost = int(consumed_tokens_hdr) * self.token_cost_ratio
await self.redis.incrbyfloat(cost_key, incremental_cost)
await self.redis.expire(cost_key, 3600)
return response
`
É perigoso entregar a chave de API mestra com permissão de gravação no repositório Git de toda a empresa a um agente de análise de código interno. Um único ataque de injeção de prompt ou alucinação do modelo pode disparar uma consulta incorreta de exclusão de branch.
Seguindo o princípio de Zero Trust do NIST SP 800-207, você deve impedir que os agentes tenham credenciais de longo prazo. Quando o agente inicia uma tarefa, faça-o emitir um JWT com escopo restrito (Scoped JWT) válido por exatamente 300 segundos (5 minutos) por meio de troca de token OAuth 2.0 (RFC 8693) no IdP.
Se você anexar o Open Policy Agent (OPA) ao sidecar do gateway e implantar a seguinte política Rego, mesmo que um agente de análise somente leitura envie uma consulta de alteração maliciosa, a infraestrutura a rejeitará com um 403.
`rego
package agent.authz
import future.keywords.in
default allow = false
required_perm_map := {
"GET": "read",
"HEAD": "read",
"POST": "write",
"PUT": "write",
"PATCH": "write",
"DELETE": "admin"
}
allow {
input.token.payload.exp > time.now_ns() / 1000000000
startswith(input.token.payload.sub, "agent:")
input.token.payload.aud == "enterprise-internal-api"
required_perm := required_perm_map[input.http_method]
expected_scope := sprintf("%s:%s", [input.resource_type, required_perm])
expected_scope in input.token.payload.scopes
not is_forbidden_mutation(input.token.payload.role, input.path)
}
is_forbidden_mutation(role, path) {
role == "readonly_sweeper"
regex.match("^/.*/(mutate|delete|drop|update|write)$", path)
}
`
Em um pipeline onde vários agentes estão encadeados, se um agente downstream travar, a fila retenta toda a mensagem. Nesse caso, você deve impedir que o agente upstream chame novamente o mesmo prompt custoso.
Em conformidade com as OpenTelemetry GenAI Semantic Conventions, agrupe cada execução de agente em um span e bloqueie a reexecução da mesma etapa pela raiz aplicando um lock distribuído no Redis.
| Chave de Atributo Semântico | Tipo | Valor de Exemplo | Propósito de Observabilidade |
|---|---|---|---|
gen_ai.operation.name |
String | invoke_agent |
Distinguir o tipo de tarefa do agente |
gen_ai.provider.name |
String | openai |
Latência e taxa de erro por provedor |
gen_ai.request.model |
String | gpt-4o |
Identificar o modelo usado |
gen_ai.usage.input_tokens |
Integer | 2048 |
Cálculo de custo de entrada por etapa |
gen_ai.usage.output_tokens |
Integer | 512 |
Rastrear tokens gerados na conclusão |
gen_ai.conversation.id |
String | task-session-9821 |
Rastreamento de sessão de agente inteira |
agent.prompt.hash |
String | sha256:7f83b165... |
Gerenciamento de versão do prompt de entrada |
`python
import hashlib
from opentelemetry import trace
from opentelemetry.trace import Status, StatusCode
import redis.asyncio as redis
tracer = trace.get_tracer("agent.pipeline.worker", "1.0.0")
async def execute_agent_step_idempotent(task_payload: dict, redis_conn: redis.Redis, llm_gateway_client) -> dict:
task_id = task_payload["task_id"]
step_id = task_payload["step_id"]
prompt = task_payload["prompt"]
idempotency_key = f"step:result:{task_id}:{step_id}"
lock_key = f"lock:step:{task_id}:{step_id}"
acquired = await redis_conn.set(lock_key, "processing", nx=True, ex=120)
if not acquired:
raise RuntimeError(f"Step {step_id} for Task {task_id} is already in progress.")
try:
cached_result = await redis_conn.get(idempotency_key)
if cached_result:
return {"status": "CACHED", "output": cached_result.decode("utf-8")}
prompt_hash = hashlib.sha256(prompt.encode("utf-8")).hexdigest()
with tracer.start_as_current_span(f"step_{step_id}") as span:
span.set_attribute("gen_ai.operation.name", "invoke_agent")
span.set_attribute("gen_ai.provider.name", "openai")
span.set_attribute("gen_ai.request.model", "gpt-4o")
span.set_attribute("gen_ai.conversation.id", task_payload["context_id"])
span.set_attribute("agent.prompt.hash", prompt_hash)
try:
inference_resp = await llm_gateway_client.generate(prompt)
span.set_attribute("gen_ai.usage.input_tokens", inference_resp.prompt_tokens)
span.set_attribute("gen_ai.usage.output_tokens", inference_resp.completion_tokens)
span.set_status(Status(StatusCode.OK))
await redis_conn.set(idempotency_key, inference_resp.content, ex=86400)
return {"status": "SUCCESS", "output": inference_resp.content}
except Exception as exc:
span.record_exception(exc)
span.set_status(Status(StatusCode.ERROR, str(exc)))
raise exc
finally:
await redis_conn.delete(lock_key)
`
O bot de revisão de código interno e o bot de consulta de implantação continuam fazendo as mesmas perguntas de padrão técnico com apenas pequenas mudanças de expressão. Um cache de correspondência de strings simples se torna inútil se a ordem das palavras mudar, fazendo com que as chamadas de API externa ocorram exatamente da mesma forma.
Você deve colocar um cache semântico baseado em RedisVL na frente e vincular rigidamente a distância de cosseno a 0.1 (similaridade de 0.95 ou superior) para evitar distorções nos dados do manual técnico interno.
| Similaridade de Cosseno | Distância de Cosseno | Taxa de Acerto de Cache | Risco de Distorção de Significado | Uso Recomendado |
|---|---|---|---|---|
| 0.95 ou superior | 0.1 ou inferior | 30% ~ 40% | Menos de 0.1% | Regulamentos internos, especificações de API, assistente de código |
| 0.85 ~ 0.94 | 0.1 ~ 0.2 | 50% ~ 70% | Nível intermediário | Orientações gerais e consultas de conveniência interna |
| Abaixo de 0.80 | Acima de 0.2 | 75% ou superior | Muito alto | Não utilizável em ambiente de produção |
O problema de emitir um cache incorreto quando o documento original é modificado pode ser resolvido detectando eventos de alteração no banco de dados com Debezium CDC e limpando o cache imediatamente.
`python
from fastapi import FastAPI, BackgroundTasks
from redisvl.extensions.cache.llm import SemanticCache
from redisvl.utils.vectorize import OpenAITextVectorizer
app = FastAPI()
vectorizer = OpenAITextVectorizer(model="text-embedding-3-small")
semantic_cache = SemanticCache(
redis_url="redis://localhost:6379",
distance_threshold=0.1, # Apenas acerta com similaridade de cosseno de 0.95 ou superior
vectorizer=vectorizer,
ttl=86400
)
@app.post("/v1/agent/query")
async def execute_agent_query(payload: dict):
query_text = payload["query"]
hit = semantic_cache.check(prompt=query_text)
if hit:
return {
"source": "SEMANTIC_CACHE",
"distance": hit[0].get("vector_distance"),
"response": hit[0]["response"]
}
llm_result = await call_upstream_llm(query_text)
semantic_cache.store(
prompt=query_text,
response=llm_result,
metadata={"domain": payload.get("domain", "general")}
)
return {"source": "LLM_GENERATED", "response": llm_result}
@app.post("/v1/cache/invalidate")
async def handle_cdc_invalidation(event: dict, background_tasks: BackgroundTasks):
table = event.get("source", {}).get("table")
op = event.get("op")
if table == "engineering_handbook" and op in ["u", "d"]:
background_tasks.add_task(purge_cache_index)
return {"status": "INVALIDATION_TRIGGERED"}
async def purge_cache_index():
semantic_cache.clear()
async def call_upstream_llm(prompt: str) -> str:
return "LLM Inference Result"
`
Se você mantiver uma camada de cache semântica com base no limite de distância de 0.1, o número de chamadas de API LLM externas será reduzido em cerca de 35% e a latência de resposta será encurtada para a faixa de segundos. Em vez de rezar implorando por condições de escape no prompt, é muito mais seguro colocar códigos de bloqueio físico nos níveis de gateway e fila.