Pourquoi une discussion entre agents se solde par une facture de 60 millions de wons
TuBrief 편집팀
2026년 9월 13일
0
Computing/Software원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
커뮤니티의 다른 글
댓글 (0)
Log in to leave a comment
아직 작성된 글이 없습니다
원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
Log in to leave a comment
아직 작성된 글이 없습니다
Dans un environnement interne combinant le bot de revue de code de l'équipe de développement et l'agent de questions de l'équipe de support, le pire cauchemar des ingénieurs n'est pas le manque d'intelligence du modèle. C'est une simple ligne de message d'échec d'analyse entrée dans la file d'attente asynchrone, et une boucle infinie où deux agents s'envoient des questions pendant tout le week-end.
Traiter les agents comme des collaborateurs autonomes ne fonctionne que dans un laboratoire de prompts. Dès qu'ils sont déployés sur l'infrastructure interne, les agents ne sont plus que des services distribués crachant des entrées non fiables. Écrire « arrête-toi si tu ne connais pas la réponse » dans un prompt en langage naturel échouera inévitablement en production. Examinons les lignes de contrôle physiques que les ingénieurs plateforme doivent placer directement au niveau des files d'attente et des passerelles.
Si vous laissez les agents communiquer en mélangeant du langage naturel et des accents Markdown (json), vous perdrez cinq à six heures par semaine rien qu'à traquer les erreurs d'analyse. En raison de la sortie non déterministe du modèle, un simple guillemet manquant ou un nom de champ légèrement modifié peut faire planter le consommateur en aval.
La solution consiste à imposer un schéma strict sur la couche de transport, à l'instar du projet Google A2A et de JSON-RPC 2.0. Les charges utiles Kafka transmises entre agents doivent obligatoirement valider les champs suivants à l'exécution.
| Nom du champ | Type | Obligatoire | Objectif de validation |
|---|---|---|---|
message_id |
UUIDv7 | Obligatoire | ID de message global triable par ordre chronologique |
task_id |
UUIDv4 | Obligatoire | Unité de suivi des tâches métier uniques |
context_id |
String | Obligatoire | Identifiant de session de conversation parent |
sender_id |
String | Obligatoire | Espace de noms de l'émetteur (domaine:nom_agent) |
receiver_id |
String | Obligatoire | Espace de noms du destinataire (domaine:nom_agent) |
hop_count |
Integer | Obligatoire | Nombre cumulé de transferts entre agents (valeur initiale : 0) |
max_hops |
Integer | Obligatoire | Nombre maximal de transferts autorisés (valeur recommandée : 5) |
constraints |
Object | Optionnel | Limite de délai d'attente et de budget de jetons |
data |
Object | Obligatoire | Données métier structurées |
Pour éviter les pannes du processus consommateur, placez un intercepteur de validation Pydantic au point d'entrée et envoyez immédiatement les messages non conformes vers une file de messages morts (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)
`
L'application de ce pattern élimine le phénomène de pilule empoisonnée qui bloque les consommateurs. Les ressources d'ingénierie consacrées au débogage de l'analyse peuvent ainsi être récupérées immédiatement.
Le moment le plus dangereux dans un pipeline d'agents survient lorsque deux agents bouclent indéfiniment en validant mutuellement leurs sorties respectives.
Il existe un cas réel documenté dans un rapport post-mortem de mars 2026. Un bot de génération SQL incapable de résoudre une contrainte de clé étrangère et un bot de validation se sont renvoyé la balle en générant sans fin des requêtes différentes. Comme les actions elles-mêmes changeaient à chaque fois, le simple compteur de nouvelles tentatives de « même action » (réglé à 50) n'a servi à rien. Les deux bots ont fonctionné sans isolation pendant 11 jours (264 heures) et ont englouti un total de 47 200 $ (environ 63 millions de wons) en frais d'API. Lorsque IAL-Scan a analysé 6 549 dépôts d'agents open source, 68 boucles infinies critiques ont été découvertes telles quelles dans 47 projets.
Il ne faut pas compter sur les conditions de sortie par prompt. Un interrupteur d'arrêt physique doit être configuré au niveau de la passerelle API.
| Critère de contrôle | Méthode de contrôle | Valeur recommandée | Règle de fonctionnement |
|---|---|---|---|
| Limite du nombre de sauts | Suivi de la profondeur d'appel dans l'en-tête | Max 5 sauts | Si délégué mutuellement plus de 5 fois, renvoie 503 et arrête |
| Plafond budgétaire | Suivi des coûts cumulés par transaction Redis | 10 $ par tâche | En cas de dépassement de 10 $ de dépenses cumulées, renvoie 429 et bloque |
| Limite de nouvelles tentatives | Nombre de répétitions pour le même ID de tâche | Max 3 fois | Même si le contenu change, arrête après 3 échecs de la même tâche |
L'association de X-Agent-Hop-Count et de Redis dans les middlewares de la passerelle permet de maîtriser l'explosion des dépenses autour de la limite des 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
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
`
Confier une clé API principale disposant de droits d'écriture sur l'ensemble des dépôts Git de l'entreprise à un agent d'analyse de code interne est dangereux. Une seule injection de prompt ou une seule hallucination du modèle pourrait déclencher une requête de suppression de branche erronée.
Conformément aux principes de zéro trust de la norme NIST SP 800-207, les agents doivent être empêchés de posséder des identifiants à long terme. Lorsqu'un agent initie une tâche, il doit obtenir via l'échange de jetons OAuth 2.0 (RFC 8693) un JWT délimité valable pendant exactement 300 secondes (5 minutes).
En attachant Open Policy Agent (OPA) au sidecar de la passerelle et en déployant la politique Rego suivante, même si un agent d'analyse en lecture seule émet une requête de modification malveillante, l'infrastructure la rejettera avec une erreur 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)
}
`
Dans un pipeline combinant plusieurs agents, si un agent en aval plante, la file d'attente réessaie le message entier. Il faut alors empêcher l'agent en amont de relancer le même prompt coûteux.
Conformément aux conventions sémantiques OpenTelemetry GenAI, chaque exécution d'agent est regroupée dans un span, et un verrou distribué Redis est appliqué pour bloquer définitivement la réexécution d'une même étape.
| Clé d'attribut sémantique | Type | Exemple de valeur | Objectif d'observabilité |
|---|---|---|---|
gen_ai.operation.name |
String | invoke_agent |
Distinction du type de tâche de l'agent |
gen_ai.provider.name |
String | openai |
Latence et taux d'erreur par fournisseur |
gen_ai.request.model |
String | gpt-4o |
Identification du modèle utilisé |
gen_ai.usage.input_tokens |
Integer | 2048 |
Facturation des coûts d'entrée par étape |
gen_ai.usage.output_tokens |
Integer | 512 |
Suivi des jetons de génération de complétion |
gen_ai.conversation.id |
String | task-session-9821 |
Suivi de la session d'agent globale |
agent.prompt.hash |
String | sha256:7f83b165... |
Gestion des versions du prompt d'entrée |
`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)
`
Le bot de revue de code interne et le bot de questions sur le déploiement posent continuellement les mêmes questions de normes techniques, formulées avec de légères variations. Un cache de correspondance de chaînes simple devient inutile dès que l'ordre des mots change, ce qui déclenche des appels d'API externes inchangés.
Il est nécessaire de placer un cache sémantique basé sur RedisVL en amont et de resserrer la distance cosinus à 0,1 (similarité de 0,95 ou plus) pour éviter la distorsion des données du manuel technique interne.
| Similarité cosinus | Distance cosinus | Taux de réussite du cache | Risque de distorsion sémantique | Utilisation recommandée |
|---|---|---|---|---|
| 0,95 ou plus | 0,1 ou moins | 30% à 40% | Moins de 0,1 % | Règlements internes, spécifications d'API, assistant de code |
| 0,85 ~ 0,94 | 0,1 ~ 0,2 | 50% à 70% | Niveau moyen | Requêtes d'information générale et de commodité interne |
| Moins de 0,80 | Plus de 0,2 | 75 % ou plus | Très élevé | Inutilisable en environnement de production |
Le problème du renvoi d'un cache erroné lorsque le document d'origine est modifié peut être résolu en détectant les événements de modification de base de données via Debezium CDC pour purger immédiatement le cache.
`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,
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"
`
L'implémentation d'une couche de cache sémantique basée sur un seuil de distance de 0,1 permet de réduire d'environ 35 % le nombre d'appels d'API LLM externes et de raccourcir le délai de réponse de plusieurs secondes. Plutôt que de compter sur les conditions de sortie écrites dans les prompts, il est beaucoup plus sûr de placer du code de blocage physique au niveau de la passerelle et de la file d'attente.