エージェント同士の会話で6000万円の請求書が届く理由
TuBrief 편집팀
2026년 9월 13일
0
Computing/Software원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
커뮤니티의 다른 글
댓글 (0)
Log in to leave a comment
아직 작성된 글이 없습니다
원본 영상을 바탕으로 AI의 도움을 받아 작성했습니다. 원본 영상이 기준입니다.
Log in to leave a comment
아직 작성된 글이 없습니다
開発チームのコードレビューボットとサポートチームの問い合わせエージェントを連携させた社内環境において、エンジニアが最も頻繁に直面する惨事は、モデルの知能不足ではありません。非同期キューに入ったパース失敗メッセージの1行、そして2つのエージェントが質問をやり取りした結果、週末中回り続けた無限ループです。
エージェントを自律的な同僚として扱うのは、プロンプトの実験室でのみ通用します。社内インフラに載せた瞬間、エージェントは信頼性の低い入力を吐き出す分散サービスに過ぎません。自然言語のプロンプト内に「答えを知らなければ中断して」と書き留めておく方法は、本番環境では確実に破綻します。プラットフォームエンジニアがキューやゲートウェイレベルに直接設置すべき物理的な制御ラインを解説します。
エージェント同士に自然言語やマークダウンのバックティック(json)を混ぜて通信させると、パースエラーの追跡だけに週に5〜6時間を奪われます。モデルの非決定的な出力のせいでクォーテーションが1つ抜けたり、フィールド名が微かに変わったりしてダウンストリームのコンシューマーがクラッシュしてしまいます。
解決策は、Google A2AのドラフトやJSON-RPC 2.0のように、転送層に厳格なスキーマを強制することです。エージェント間で伝送されるKafkaのペイロードは、ランタイムにおいて必ず以下のフィールドを通過しなければなりません。
| フィールド名 | タイプ | 必須の有無 | 検証の目的 |
|---|---|---|---|
message_id |
UUIDv7 | 必須 | 時系列順のソートが可能なグローバルメッセージID |
task_id |
UUIDv4 | 必須 | 単一のビジネス作業追跡単位 |
context_id |
String | 必須 | 上位の会話セッション識別子 |
sender_id |
String | 必須 | 送信元名前空間 (ドメイン:エージェント名) |
receiver_id |
String | 必須 | 受信先名前空間 (ドメイン:エージェント名) |
hop_count |
Integer | 必須 | エージェント間の伝送累積回数 (初期値: 0) |
max_hops |
Integer | 必須 | 許容最大伝送回数 (推奨値: 5) |
constraints |
Object | 選択 | タイムアウト、トークン予算制限 |
data |
Object | 必須 | 構造化されたビジネスデータ |
コンシューマープロセスがクラッシュする事故を防ぐには、流入地点にPydantic検証インターセプターを置き、規約に違反したメッセージは即座にデッドレターキュー(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)
`
このパターンを適用しておけば、コンシューマーが機能停止するポイズンピル現象が解消されます。パースのデバッグに費やしていたエンジニアリングリソースも即座に回収できます。
エージェントパイプラインにおいて最も危険な瞬間は、2つのエージェントがお互いの出力を検証しながら無限ループに陥った時です。
2026年3月の事後分析レポートで公表された実際の事例があります。外部キー制約を解消できなかったSQL生成ボットと検証ボットが異なるクエリを際限なく生成し続けながらピンポンを繰り返し、アクション自体が毎回異なっていたために単純な「同一アクション」の再試行カウンター50回が機能しませんでした。両者は11日間(264時間)にわたり隔離なしで稼働し、合計47,200ドル(約6,300万円)のAPI費用を消費しました。IAL-Scanがオープンソースのエージェントリポジトリ6,549個を分析した際にも、47個のプロジェクトから68件の致命的な無限ループがそのまま発見されました。
プロンプトの脱出条件を過信してはなりません。APIゲートウェイレベルで物理的なキースイッチをかける必要があります。
| 制御基準 | 制御方式 | 推奨値 | 動作ルール |
|---|---|---|---|
| ホップ数制限 | ヘッダーの呼び出し深度追跡 | 最大5ホップ | 5回を超えてお互いに委任した場合、503を返して中断 |
| 予算上限 | Redisのトランザクション別累積コスト追跡 | タスクあたり $10 | 累積支出が10ドルを超過した場合、429を返して遮断 |
| 再試行制限 | 同一作業ID対象の繰り返し回数 | 最大3回 | 内容が変わっても同一作業が3回失敗した時点で終了 |
ゲートウェイミドル웨어で X-Agent-Hop-Count とRedisを連携させれば、支出の急増を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
`
社内コード分析エージェントに、全社のGitリポジトリ書き込み権限を持つマスターAPIキーを渡すのは危険です。プロンプトインジェクションやモデルのハルシネーションが1度起きるだけで、不正なブランチ削除クエリが実行される恐れがあります。
NIST SP 800-207のゼロトラストの原則に従い、エージェントが長期的な資格情報を持たないように制限すべきです。エージェントが作業を開始する際、IdPからのOAuth 2.0トークン交換(RFC 8693)を通じて、わずか300秒(5分間)だけ持続する狭いスコープのScoped JWTを発行させるようにします。
ゲートウェイのサイドカーにOpen Policy Agent(OPA)を組み込み、以下のRegoポリシーをデプロイすれば、読み取り専用の分析エージェントが悪意のある変更クエリを送信したとしても、インフラ側で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)
}
`
複数のエージェントが結びついたパイプラインにおいて下流のエージェントが停止した場合、キューはメッセージ全体を再試行します。このとき、上流のエージェントが全く同じ高コストなプロンプトを再呼び出ししないよう防ぐ必要があります。
OpenTelemetry GenAI Semantic Conventionsに準拠し、各エージェントの実行をスパンとしてまとめ、Redisの分散ロックをかけて同一ステップの再実行を大元から遮断します。
| セマンティック属性キー | タイプ | 例文・値 | 観測の目的 |
|---|---|---|---|
gen_ai.operation.name |
String | invoke_agent |
エージェントの作業タイプの区分 |
gen_ai.provider.name |
String | openai |
プロバイダー別のレイテンシとエラー率 |
gen_ai.request.model |
String | gpt-4o |
使用モデルの識別 |
gen_ai.usage.input_tokens |
Integer | 2048 |
ステップ別の入力コスト精算 |
gen_ai.usage.output_tokens |
Integer | 512 |
生成完了トークンの追跡 |
gen_ai.conversation.id |
String | task-session-9821 |
エージェントセッション全体の追跡 |
agent.prompt.hash |
String | sha256:7f83b165... |
入力プロンプトのバージョン管理 |
`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)
`
社内のコードレビューボットやデプロイ問い合わせボットは、表現が少し変わっただけの同じ技術標準に関する質問を繰り返し送信します。単純な文字列一致によるキャッシュは、語順が変わっただけで役に立たなくなるため、外部APIの呼び出しがそのまま発生してしまいます。
RedisVLベースのセマンティックキャッシュをフロントエンドに配置し、社内技術ハンドブックデータの歪みを防ぐためにコサイン距離を0.1(類似度0.95以上)にタイトに絞り込む必要があります。
| コサイン類似度 | コサイン距離 | キャッシュヒット率 | 意味の歪みリスク | 推奨用途 |
|---|---|---|---|---|
| 0.95 以上 | 0.1 以下 | 30% ~ 40% | 0.1% 未満 | 社内規程、API仕様、コードアシスタント |
| 0.85 ~ 0.94 | 0.1 ~ 0.2 | 50% ~ 70% | 中間レベル | 一般案内および社内向け問い合わせ |
| 0.80 未満 | 0.2 超過 | 75% 以上 | 非常に高い | 本番環境での使用不可 |
元のドキュメントが修正された際に古いキャッシュを返してしまう問題は、DebeziumのCDC(Change Data Capture)によりDBの変更イベントを検知してキャッシュを即座に削除することで解決できます。
`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"
`
距離の閾値を0.1としたセマンティックキャッシュレイヤーを配置すれば、外部LLM APIの呼び出し回数を約35%削減でき、応答遅延時間も数秒単位に短縮されます。プロンプト内の脱出条件に頼って事態の好転を祈るよりも、ゲートウェイやキューのレベルに物理的な遮断コードを組み込む方がはるかに安全です。