Por qué una conversación entre agentes puede generar una factura de 60 mil dólares
TuBrief 편집팀
2026년 9월 13일
0
Computing/Software원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
커뮤니티의 다른 글
댓글 (0)
Log in to leave a comment
아직 작성된 글이 없습니다
원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
Log in to leave a comment
아직 작성된 글이 없습니다
En un entorno interno donde se conectan el bot de revisión de código del equipo de desarrollo y el agente de consultas del equipo de soporte, el peor desastre que un ingeniero experimenta con mayor frecuencia no es la falta de inteligencia del modelo. Es una sola línea de mensaje con error de análisis que entró en una cola asíncrona, y un bucle infinito en el que dos agentes se pasaron preguntas el fin de semana entero.
Tratar a los agentes como colegas autónomos solo funciona en un laboratorio de prompts. En el momento en que se despliegan en la infraestructura interna, los agentes no son más que un servicio distribuido que arroja entradas poco confiables. El enfoque de escribir "detente si no sabes la respuesta" dentro de un prompt en lenguaje natural falla inevitablemente en producción. A continuación, revisamos las líneas de control físico que los ingenieros de plataforma deben implementar directamente a nivel de cola y pasarela.
Si dejas que los agentes se comuniquen usando lenguaje natural o mezclando bloques de código en Markdown (json), perderás de cinco a seis horas a la semana solo rastreando errores de análisis. Debido al resultado no determindeterminado del modelo, la falta de una sola comilla o un cambio mínimo en el nombre de un campo hará que el consumidor aguas abajo falle.
La solución es aplicar un esquema estricto en la capa de transporte, al igual que el borrador de Google A2A y JSON-RPC 2.0. La carga útil de Kafka transferida entre agentes debe pasar obligatoriamente por los siguientes campos en tiempo de ejecución:
| Nombre del campo | Tipo | Requerido | Propósito de validación |
|---|---|---|---|
message_id |
UUIDv7 | Sí | ID de mensaje global ordenable cronológicamente |
task_id |
UUIDv4 | Sí | Unidad de seguimiento de tarea de negocio única |
context_id |
String | Sí | Identificador de sesión de conversación superior |
sender_id |
String | Sí | Espacio de nombres del remitente (dominio:nombre_de_agente) |
receiver_id |
String | Sí | Espacio de nombres del receptor (dominio:nombre_de_agente) |
hop_count |
Integer | Sí | Número acumulado de transferencias entre agentes (valor inicial: 0) |
max_hops |
Integer | Sí | Número máximo de transferencias permitidas (valor recomendado: 5) |
constraints |
Object | Opcional | Tiempo de espera, límite de presupuesto de tokens |
data |
Object | Sí | Datos de negocio estructurados |
Para evitar accidentes donde los procesos consumidores colapsen, debes colocar un interceptor de validación Pydantic en el punto de entrada y enviar inmediatamente los mensajes que violen la especificación a una cola de mensajes muertos (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")
context_id: str = Field(..., description="트랜잭션 세션 식별자")
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="비즈니스 페이로드")
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)
`
Al implementar este patrón, el fenómeno de mensajes envenenados (poison pill) que bloquea al consumidor desaparece. También puedes recuperar de inmediato los recursos de ingeniería gastados en depurar análisis.
El momento más peligroso en una canalización de agentes es cuando dos agentes entran en un bucle infinito validando las salidas del otro.
Existe un caso real documentado en un informe de análisis posterior de marzo de 2026. Un bot generador de SQL que no pudo resolver restricciones de clave foránea y un bot de validación crearon consultas diferentes sin parar y jugaron al ping-pong; como la acción en sí era diferente cada vez, el contador simple de reintentos de "misma acción" de 50 veces no funcionó. Ambos operaron sin aislamiento durante 11 días (264 horas) y consumieron un total de 47,200 dólares en costos de API. Cuando IAL-Scan analizó 6,549 repositorios de agentes de código abierto, se descubrieron directamente 68 bucles infinitos críticos en 47 proyectos.
No debes confiar en las condiciones de salida de los prompts. Debes aplicar un interruptor de apagado físico a nivel de la pasarela de API.
| Criterio de control | Método de control | Valor recomendado | Regla de operación |
|---|---|---|---|
| Límite de saltos | Seguimiento de profundidad de llamadas en cabecera | Máximo 5 saltos | Devuelve 503 y se detiene si se delegan mutuamente más de 5 veces |
| Límite de presupuesto | Seguimiento de costo acumulado por transacción en Redis | $10 por tarea | Devuelve 429 y bloquea si el gasto acumulado supera los 10 dólares |
| Límite de reintentos | Número de repeticiones para el mismo ID de tarea | Máximo 3 veces | Finaliza tras 3 fallos de la misma tarea aunque el contenido cambie |
Si vinculas X-Agent-Hop-Count y Redis en el middleware de la pasarela, puedes controlar la explosión de gastos en un límite 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 # 1,000토큰당 $0.015 기준 계산
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"홉 한도 초과 차단: 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"예산 초과 차단: 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
`
Entregar una clave API maestra con permisos de escritura en el repositorio Git corporativo a un agente de análisis de código interno es peligroso. Una sola inyección de prompt o alucinación del modelo podría enviar una consulta de eliminación de rama errónea.
Siguiendo los principios de Confianza Cero de NIST SP 800-207, se debe evitar que los agentes posean credenciales a largo plazo. Haz que, al iniciar una tarea, el agente obtenga un JWT con ámbito limitado (Scoped JWT) válido por exactamente 300 segundos (5 minutos) mediante el intercambio de tokens OAuth 2.0 (RFC 8693) en el IdP.
Si adjuntas Open Policy Agent (OPA) al sidecar de la pasarela y despliegas la siguiente política Rego, incluso si un agente de análisis de solo lectura lanza una consulta de modificación maliciosa, la infraestructura la rechazará con un error 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)
}
`
En una canalización con múltiples agentes interconectados, si un agente inferior falla, la cola reintentará todo el mensaje. En este caso, debes evitar que el agente superior vuelva a invocar exactamente el mismo y costoso prompt.
Cumpliendo con las convenciones semánticas de OpenTelemetry GenAI (OpenTelemetry GenAI Semantic Conventions), agrupa cada ejecución de agente en un span y aplica un bloqueo distribuido en Redis para bloquear desde la raíz la reejecución de la misma fase.
| Clave de atributo semántico | Tipo | Valor de ejemplo | Propósito de observación |
|---|---|---|---|
gen_ai.operation.name |
String | invoke_agent |
Distinguir el tipo de tarea del agente |
gen_ai.provider.name |
String | openai |
Latencia y tasa de error por proveedor |
gen_ai.request.model |
String | gpt-4o |
Identificación del modelo utilizado |
gen_ai.usage.input_tokens |
Integer | 2048 |
Liquidación de costo de entrada por fase |
gen_ai.usage.output_tokens |
Integer | 512 |
Seguimiento de tokens generados de salida |
gen_ai.conversation.id |
String | task-session-9821 |
Seguimiento de la sesión general del agente |
agent.prompt.hash |
String | sha256:7f83b165... |
Gestión de versiones del 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)
`
El bot de revisión de código interno y el bot de consultas de despliegue envían continuamente las mismas consultas de estándares técnicos con ligeras variaciones en la redacción. Una caché de coincidencia de cadenas simple resulta inútil si cambia el orden de las palabras, lo que provoca que las llamadas a la API externa ocurran tal cual.
Debes colocar una caché semántica basada en RedisVL al frente y ajustar rígidamente la distancia coseno a 0.1 (similitud de 0.95 o superior) para evitar la distorsión de los datos del manual técnico interno.
| Similitud coseno | Distancia coseno | Tasa de aciertos de caché | Riesgo de distorsión de significado | Uso recomendado |
|---|---|---|---|---|
| 0.95 o más | 0.1 o menos | 30% ~ 40% | Menos de 0.1% | Normas internas, especificaciones de API, asistente de código |
| 0.85 ~ 0.94 | 0.1 ~ 0.2 | 50% ~ 70% | Nivel medio | Guías generales y consultas de conveniencia interna |
| Menos de 0.80 | Más de 0.2 | 75% o más | Muy alto | No apto para entorno de producción |
El problema de arrojar una caché errónea cuando se modifica el documento original se resuelve detectando los eventos de cambio de base de datos con Debezium CDC para borrar la caché de inmediato.
`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, # 코사인 유사도 0.95 이상만 적중
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"
`
Establecer una capa de caché semántica basada en un umbral de distancia de 0.1 reduce las llamadas a la API de LLM externa en aproximadamente un 35% y también acorta la latencia de respuesta a nivel de unos pocos segundos. En lugar de rezar confiando en las condiciones de salida en los prompts, es mucho más seguro subir código de bloqueo físico a los niveles de pasarela y cola.