← Notes

障害情報の分析を自律化してみた(環境詳細)

2026-09-13
目次

元記事

zenn に投稿した元記事はこちらです。

障害情報の調査と通知を自律化してみたzenn.dev

環境

実行環境の概要です。

#アプリ役割バージョン
1job-api / slot-manager監視対象(自作)。FastAPI + OTel Python SDKPython 3.13.15、fastapi 0.141.1、opentelemetry-sdk 1.44.0、instrumentation 0.65b0
2OpenTelemetry CollectorOTLP の振り分けcontrib 0.160.0
3Prometheusメトリクス保存とアラート評価3.14.0
4Lokiログ保存、ruler でログ由来のアラート3.7.7
5Tempoトレース保存、MCP サーバ(/api/mcp)3.0.0
6Alertmanagerwebhook 転送0.28.1
7Grafana可視化と annotation13.2.1
8mcp-grafanaMCP サーバ(読み取り専用。Tempo ツールも中継)1.4.1
9Claude Codeエージェント(Sonnet 5)2.1.268
10llama.cppローカル LLM(Qwen3.6-35B-A3B UD-IQ3_XXS)ghcr.io/ggml-org/llama.cpp:server-cuda、ctx 65536、--jinja
11LiteLLMClaude Code からローカル LLM への変換ghcr.io/berriai/litellm:main-stable
12local_agent.pyローカル LLM 用の最小エージェント(自作)Python 3.14.0、mcp 2.21.0、openai SDK
13Mattermost通知先(社内チャットの代わり)mattermost-team-edition 11.11.0 + postgres:16-alpine
14Docker Desktop(Windows)2〜8、11、13 のコンテナを動かすEngine 29.7.2、Compose 5.5.0
15ハードウェアGPURTX 5070 Ti 16 GB

図1: docker compose の 12 コンテナと docker ホストのプロセス、ポート 図1: docker compose の 12 コンテナと docker ホストのプロセス、ポート

監視対象(自作サービス)

監視対象のサービスとして、計算ジョブの投入システムを 2 サービスで作りました。 app/ に保存します。

図2: 監視対象の 2 サービスと故障注入の入口 図2: 監視対象の 2 サービスと故障注入の入口

サービスポートエンドポイント役割
job-api8100POST /jobs、POST /admin/faultジョブ(種別 mip / lp / simulation / ml、size 1〜3)を受け、1 件ごとに slot-manager へ枠の確保を頼む
slot-manager8101POST /reserve、POST /admin/faultDB を SELECT → UPDATE(各 40 ms を sleep で模擬、接続プール 4 本)して枠を確保する。無ければ 409

計装は OTel Python SDKを利用します。 opentelemetry-instrumentation-fastapi と -httpx の自動計装に、標準 logging の LoggingHandler を足しました(app/telemetry.py)。 OTEL_SEMCONV_STABILITY_OPT_IN=http で http.server.request.duration(秒)が出力されます。

出力は次のとおりです。

出力情報内容
メトリクスjobs_total{outcome = ok / failed}、jobs_duration_seconds(ヒストグラム)、自動計装の http_server_request_duration_seconds など。種別のラベルは付けない
ログjob received / job submitted / job failed: ...(job-api)、reserved: mip x2 / failed to reserve: no free slot: mip / slots running low: mip(slot-manager)。属性 job.type、slots.remaining は Loki の構造化メタデータになる
トレースPOST /jobs → POST(httpx)→ POST /reserve → db.pool.acquire → db.query(db.operation.name = SELECT slots / UPDATE slots)

異常の発生は POST /admin/fault で行います。

curl -s -X POST localhost:8101/admin/fault -d '{"latency_ms": 300}'      # DB を遅くする
curl -s -X POST localhost:8101/admin/fault -d '{"slots": {"mip": 0}}'     # mip の枠を 0 に
curl -s -X POST localhost:8100/admin/fault -d '{"duplicate_calls": 8}'   # 1 ジョブで 8 回 reserve する(二重投入のバグ)
curl -s -X POST localhost:8101/admin/fault -d '{"reset": true}'          # 戻す(8100 も同様)

Grafana 周辺環境

docker-compose.yml のサービスとポート。既存の Grafana などと衝突しないようずらしてある。

サービスイメージホスト側ポート備考
otel-collectorotel/opentelemetry-collector-contrib:0.160.04317、4318、8889OTLP 受信。トレース → Tempo、ログ → Loki(/otlp)、メトリクス → :8889 を Prometheus が scrape
prometheusprom/prometheus:v3.14.09092alert_rules.yml を評価し Alertmanager へ
alertmanagerprom/alertmanager:v0.28.19094host.docker.internal:9095/alert へ webhook
lokigrafana/loki:3.7.73101ruler 有効。loki/rules/fake/job-system.yml を同じ Alertmanager へ
tempografana/tempo:3.0.03200query_frontend.mcp_server.enabled: true
grafanagrafana/grafana:13.2.13082データソース 3 つと Job system ダッシュボードを provisioning
mcp-grafanagrafana/mcp-grafana:1.4.18097--disable-write と --disable-* で 27 ツールに絞る
litellmghcr.io/berriai/litellm:main-stable4000ローカル LLM 用
mattermostmattermost/mattermost-team-edition:11.11.08065通知先。#alerts に Alertmanager の発火通知、#sonnet / #qwen にモデルごとの AI 調査レポート。DB は postgres:16-alpine(公開しない)

アラートルール(Prometheus)。

alert式for
JobApiLatencyHighhistogram_quantile(0.95, sum(rate(jobs_duration_seconds_bucket[2m])) by (le)) > 0.51m
JobErrorRateHighsum(rate(jobs_total{outcome="failed"}[2m])) / sum(rate(jobs_total[2m])) > 0.21m
ScrapeTargetDownup == 01m

Loki ruler のルール(ログ由来)。

- alert: SlotsRunningLow
  expr: sum(count_over_time({service_name="slot-manager"} |= "slots running low" [5m])) > 3
  for: 1m

mcp-grafana は --disable-write に加えて、調査に使わない分類(admin、dashboard、search、oncall、incident、sift、pyroscope、docs など 28 個)を --disable-* で落とした。 残るのは alerting / annotations / datasource / prometheus / loki / proxied(Tempo)の 27 ツール。 ツール定義のトークン量を減らすためで、Claude Code のシステムプロンプト約 18k に全 64 ツールを足すと 36k になり、ローカル LLM の 32k を初回で超えた。

Tempo は query_frontend.mcp_server.enabled: true で http://tempo:3200/api/mcp に MCP サーバを立てる。 mcp-grafana 1.4.1 はこれを tempo_traceql-search / tempo_get-trace / tempo_get-attribute-names など 7 つのツールとして中継する。 Claude Code の mcp.json は mcp-grafana 1 本でよい。

AI エージェント周辺環境

図3: 発火から調査、書き戻しまでの部品 図3: 発火から調査、書き戻しまでの部品

sequenceDiagram
    participant P as Prometheus → Alertmanager
    participant R as investigate.py
    participant C as エージェント
    participant M as mcp-grafana / Grafana
    participant MM as Mattermost
    P->>MM: 発火通知(#35;alerts)
    P->>R: POST /alert(webhook)
    R->>R: 重複を除外(300 秒)
    R->>C: PROMPT.md で起動
    loop メトリクス → 仮説 → ログ・トレース
        C->>M: ツール呼び出し
        M-->>C: 結果
    end
    C-->>R: レポート
    R->>R: reports/ に保存
    R->>M: annotation(先頭 1,500 文字)
    R->>MM: レポート全文(#35;sonnet / #35;qwen)

図4: 発火から調査、通知までの流れ

受け口(investigate.py)

標準ライブラリだけの HTTP サーバ(ポート 9095)。

  1. POST /alert で Alertmanager の webhook を受け、status == firing のアラートだけを対象にする

  2. 同じ fingerprint は DEDUP_SECONDS(実験中は 300)の間は再調査しない

  3. PROMPT.md の {alert_json} をアラートの JSON で置き換え、claude -p --model sonnet --output-format json --mcp-config mcp.json --strict-mcp-config --allowedTools "mcp__grafana__*" --disallowedTools <組み込みツール全部> --no-session-persistence を起動する

  4. 結果を reports/<UTC 時刻>_<alertname>.md に保存し、所要時間・ターン数・トークン使用量を reports/results.jsonl に追記する

  5. レポートの先頭 1,500 文字を Grafana の annotation(tags: ["ai-investigation", alertname])に書き戻す。認証は実験なので admin の基本認証

  6. レポートの全文(Markdown)を Mattermost の #alerts に incoming webhook で投稿する(上限 16,000 文字)。投稿先は環境変数 MATTERMOST_CHANNEL で変えられる。実験では Sonnet 5 の受け口を #sonnet、Qwen の受け口を #qwen に向けて、モデルごとにレポートを分けた

CLAUDE_BIN を差し替えると別のエージェントを同じ受け口で動かせる。

通知先(Mattermost)

プロキシ越しに Slack や Teams へ出られないオンプレを想定し、通知先は同じ compose に立てた Mattermost(Team Edition)にした。 Alertmanager の slack_configs は Mattermost の incoming webhook をそのまま受け付けるので、発火と解消は Alertmanager から、調査レポートは受け口から、同じ #alerts に届く。

docker compose up -d mattermost      # mattermost-db(PostgreSQL)も上がる
python setup_mattermost.py           # 管理者 / チーム ops / チャンネル alerts / webhook を作り、.env と alertmanager.yml を生成
docker compose restart alertmanager

alertmanager/alertmanager.yml は alertmanager.tmpl.yml の {{MATTERMOST_HOOK_ID}} を埋めて生成する。画面は http://localhost:8065(初期ユーザーは setup_mattermost.py の引数)。

受け口の投稿先は環境変数 MATTERMOST_CHANNEL(既定 alerts)で変えられる。incoming webhook は投稿時にチャンネル名を指定できるので、webhook は 1 本のままでよい。 実験ではモデルごとにチャンネルを分け、Sonnet 5 の受け口は #sonnet、Qwen の受け口は #qwen に投稿させた。

python setup_mattermost.py --channel sonnet   # チャンネルだけ足す(webhook は既存を使い回す)
python setup_mattermost.py --channel qwen
MATTERMOST_CHANNEL=sonnet python investigate.py

図5: #sonnet に届いた調査レポート 図5: #sonnet に届いた調査レポート

プロンプト(PROMPT.md)

アラートの JSON、監視対象 2 サービスとデータソース 3 つの説明、クエリの書き方(PromQL の窓と stepSeconds、LogQL のフィルタ、TraceQL の例)、手順(メトリクス → 仮説 → ログとトレース → annotation)、出力の見出し構成(症状 / 仮説と検証 / 結論 / 次の一手)。 v2 で「調べ方の原則」を 3 行足した。

  • エラーの原因は、エラーを返した側(下流)のログから引く。呼び出し側の HTTP ステータスは結果であって原因ではない
  • 失敗が一部に偏っていないか(種別・対象・ユーザー)をログの属性で集計し、偏りの値を結論に書く
  • 過去の調査結果(annotation)は参考に留め、今回の時間帯のデータで確認できたことだけを根拠にする

ローカル LLM

llama.cpp は既存のスタック(ghcr.io/ggml-org/llama.cpp:server-cuda)を使う。関係する起動オプションは -hf unsloth/Qwen3.6-35B-A3B-GGUF:UD-IQ3_XXS、-c 65536(既定 32768 から上げた)、-fa on、--cache-type-k q8_0 --cache-type-v q8_0、-np 1、--jinja(tool calling に必要)、--reasoning。 16 GB の GPU で ctx 65536 にすると空き VRAM は約 1.6 GB。

Claude Code のモデルだけを差し替える経路は LiteLLM を使う。litellm/config.yaml で qwen3.6-35b-a3b を openai/qwen3.6-35b-a3b(api_base: http://host.docker.internal:8099/v1)に対応づけ、Claude Code 側は環境変数で向きを変える。

$env:ANTHROPIC_BASE_URL = "http://localhost:4000"
$env:ANTHROPIC_AUTH_TOKEN = "sk-local"
$env:ANTHROPIC_MODEL = "qwen3.6-35b-a3b"
$env:ANTHROPIC_SMALL_FAST_MODEL = "qwen3.6-35b-a3b"
$env:CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC = "1"
$env:CLAUDE_MODEL = "qwen3.6-35b-a3b"

この経路は 32k でも 64k でもコンテキストを超えたので、記事の計測では自作の local_agent.py を使った。 mcp SDK(Streamable HTTP)で mcp-grafana のツール一覧を取り、OpenAI 互換 API に function calling で回す約 200 行。 受け口の CLAUDE_BIN を python local_agent.py にすると差し替わる。既定値は次のとおり(環境変数で変えられる)。

項目値理由
ツール呼び出しの予算20 回(LOCAL_MAX_TOOL_CALLS)各結果の末尾に残り回数を付け、使い切ったら報告を書かせる
ツール結果の上限6,000 文字(LOCAL_MAX_TOOL_CHARS)Tempo の検索結果とトレースは表に整形してから収める
圧縮を始める prompt_tokens52,000(LOCAL_COMPACT_AT)古いツール結果を 300 文字に畳む
サンプリング温度 0.6、top_p 0.95、presence_penalty 1.0Qwen3 の推奨値。温度 0.2 では同じ文の繰り返しに入った
同じ引数の再呼び出し実行せず注意文を返す
本文に漏れた思考</think> より前を落とすllama.cpp が分離し損ねたとき用

ツール呼び出しの回し方は次のとおり。

sequenceDiagram
    participant R as investigate.py
    participant L as local_agent.py
    participant M as mcp-grafana(MCP)
    participant Q as llama.cpp(Qwen3.6-35B-A3B)
    R->>L: プロンプト(標準入力)
    L->>M: tools/list(27 ツール)
    L->>Q: system + プロンプト + ツール定義
    loop 最大 20 回のツール呼び出し
        Q-->>L: tool_calls
        alt 同じ引数の再呼び出し
            L-->>Q: 実行せず注意文を返す
        else 新しい呼び出し
            L->>M: tools/call
            M-->>L: 結果
            L->>L: Tempo の結果は表に整形、6,000 文字で切る、残り回数を付ける
            L-->>Q: tool result
        end
        opt prompt_tokens が 52,000 を超えた
            L->>L: 古いツール結果を 300 文字に畳む
        end
    end
    Q-->>L: 最終レポート(思考タグより前を落とす)
    L-->>R: JSON(result、num_turns、duration_ms)

図6: local_agent.py の 1 回の調査

32k で使うなら 14 回 / 3,000 文字 / 24,000 に下げる。

docker compose 起動手順

  1. docker compose up -d --build
  2. Mattermost の初期設定: python setup_mattermost.py のあと docker compose restart alertmanager(初回だけ)
  3. 平常時の負荷を流す: python loadgen.py --rps 5
  4. 受け口を起動する: python investigate.py(GRAFANA_URL の既定は http://localhost:3082)。ローカル LLM なら CLAUDE_BIN="python local_agent.py" CLAUDE_MODEL=qwen3.6-35b-a3b を付ける。投稿先を分けるなら MATTERMOST_CHANNEL=sonnet のように付ける
  5. docker compose ps で 12 コンテナが Up、Grafana(3082)で Prometheus / Loki / Tempo の 3 データソースが緑、claude mcp list で grafana が connected、であれば準備完了

docker ホストが Windows で、別の PC から ssh で操作するときの注意が 2 つある。

  • 公開鍵認証の ssh セッションでは Docker の資格情報ヘルパー(wincred)が「A specified logon session does not exist」で落ち、docker pull と compose build が失敗する。PATH から Docker のディレクトリを外して docker.exe をフルパスで呼び、DOCKER_CONFIG を空の config.json に向けると匿名で pull できる
  • ssh セッションから Start-Process した常駐プロセス(受け口、loadgen)はセッションが切れると死ぬ。Invoke-CimMethod -ClassName Win32_Process -MethodName Create で起動するとセッションから切り離される

実験手順

共通: 前の故障を戻し、アラートが解消してから 5 分空けて次を注入する。故障スクリプトは注入 → 指定秒数待つ → 戻す、まで行う。 ただしプロンプトの検索窓は発火の前後 30 分なので、5 分では前の実験のトレースが窓に残る。実験 3 の Qwen は 20 分前の実験 1 のトレースを開いて誤答した。前の実験の痕跡を完全に消すなら 30 分空ける。

sequenceDiagram
    participant E as 実験者(faults/*.py)
    participant S as job-api / slot-manager
    participant P as Prometheus → Alertmanager
    participant R as investigate.py → Mattermost
    E->>S: POST /admin/fault(reset)
    Note over S,P: 解消を待ち、さらに 5 分空ける
    E->>S: POST /admin/fault(注入。600 秒で自動復旧)
    S-->>P: メトリクスが変わる(投入は 5 rps のまま)
    P->>R: firing → webhook(注入の 1〜3 分後)
    R-->>R: 調査(Sonnet 5 約 2 分、Qwen 約 40 秒)
    Note over E,R: reports/ を読み、故障を名指しできたかを ○△× で採点

図7: 1 回の実験の手順

実験1(投入の急増)

python faults/load_spike.py --rps 80 --seconds 600 --max-inflight 60

平常 5 rps に 80 rps を上乗せする。DB の接続プール 4 本 × 80 ms/ジョブ = 約 50 rps を超えるのでプール待ちが伸び、JobApiLatencyHigh が注入の約 1 分 20 秒後に発火する。 --max-inflight 60 で同時実行を頭打ちにする。頭打ちなしだと待ち行列が無限に伸び、job-api の httpx タイムアウトが 502 になって JobErrorRateHigh まで鳴る。

実験2(特定種別の枠切れ)

python faults/slot_out.py --job mip --seconds 600

mip の枠を 0 にする。mip の投入(全体の約 25%)だけが 409 で失敗し、JobErrorRateHigh が約 2 分 40 秒後に発火する。 --slots 45 --seconds 25 にすると slots running low の警告だけが出て、Loki ruler の SlotsRunningLow が鳴る(2b)。

実験3(DB の遅延)

python faults/slow_db.py --latency-ms 300 --seconds 600

slot-manager の DB 問い合わせに 300 ms を足す。1 ジョブで 2 回問い合わせるので投入 1 件が約 0.7 秒になり、実験 1 と同じ JobApiLatencyHigh が約 1 分 20 秒後に発火する。 python faults/retry_storm.py --calls 8 は job-api が 1 ジョブで 8 回 reserve するバグ(3b)。

判定

reports/<時刻>_<alertname>.md を読み、注入した故障を原因として名指しできたかを○△×で採点する。所要時間とターン数は reports/results.jsonl。 ローカル LLM の場合はツール呼び出しの記録(呼んだツール、引数、結果の文字数)が reports/<時刻>_<alertname>.log に残る。 Claude Code の場合は --output-format json の出力に最終レポートしか入らないので、どのクエリを実行したかはレポートの「検証に使ったクエリ」列から追う。呼び出しまで残すなら --output-format stream-json にして受け口で保存する。 レポートの主張は Loki や Tempo に同じクエリを直接投げて突き合わせる。実験 2 なら sum by (job_type) (count_over_time({service_name="slot-manager"} |= "failed to reserve" [10m])) で失敗が 1 種別に偏っているかを数える。

コード

実験に使ったコードと設定の全文。zenn 側のリポジトリは非公開なので、ここに置く。ファイル名をクリックすると開く。

スタック

docker-compose.yml 12 コンテナの定義。ポートは既存環境とずらしてある
# アラート起点の自律調査(レベル1〜3)の実験スタック。
# 監視対象は自作の計算ジョブ投入システム(job-api / slot-manager)。メトリクス・ログ・トレースを
# OTLP で Collector に送り、Prometheus / Loki / Tempo に振り分ける。
# 受け口(investigate.py)は claude CLI が必要なので docker ホストで動かす。
#   Grafana      : http://localhost:3082 (admin / admin)
#   Prometheus   : http://localhost:9092
#   Alertmanager : http://localhost:9094
#   mcp-grafana  : http://localhost:8097/mcp
#   Tempo MCP    : http://localhost:3200/api/mcp
#   Loki         : http://localhost:3101
#   job-api      : http://localhost:8100
#   slot-manager : http://localhost:8101
# 既存の Grafana / Prometheus / Loki と衝突しないようポートをずらしてある。
name: grafana-alert

x-app: &app
  build: ./app
  restart: unless-stopped
  environment:
    OTEL_EXPORTER_OTLP_ENDPOINT: http://otel-collector:4317
    OTEL_EXPORTER_OTLP_INSECURE: "true"
    OTEL_SEMCONV_STABILITY_OPT_IN: http
    PYTHONUNBUFFERED: "1"
  depends_on:
    - otel-collector

services:
  job-api:
    <<: *app
    command: uvicorn job_api:app --host 0.0.0.0 --port 8000 --no-access-log
    ports:
      - "8100:8000"
    environment:
      OTEL_EXPORTER_OTLP_ENDPOINT: http://otel-collector:4317
      OTEL_EXPORTER_OTLP_INSECURE: "true"
      OTEL_SEMCONV_STABILITY_OPT_IN: http
      PYTHONUNBUFFERED: "1"
      SLOT_MANAGER_URL: http://slot-manager:8000
    depends_on:
      - otel-collector
      - slot-manager

  slot-manager:
    <<: *app
    command: uvicorn slot_manager:app --host 0.0.0.0 --port 8000 --no-access-log
    ports:
      - "8101:8000"

  otel-collector:
    image: otel/opentelemetry-collector-contrib:0.160.0
    volumes:
      - ./otel-collector/config.yaml:/etc/otelcol-contrib/config.yaml:ro
    ports:
      - "4317:4317"
      - "4318:4318"
      - "8889:8889"
    depends_on:
      - loki
      - tempo
    restart: unless-stopped

  loki:
    image: grafana/loki:3.7.7
    ports:
      - "3101:3100"
    volumes:
      - ./loki/config.yaml:/etc/loki/config.yaml:ro
      - ./loki/rules:/loki/rules:ro
      - loki-data:/loki
    command:
      - -config.file=/etc/loki/config.yaml
    restart: unless-stopped

  tempo:
    image: grafana/tempo:3.0.0
    user: "0"
    ports:
      - "3200:3200"
    volumes:
      - ./tempo/config.yaml:/etc/tempo/config.yaml:ro
      - tempo-data:/var/tempo
    command:
      - -config.file=/etc/tempo/config.yaml
    restart: unless-stopped

  prometheus:
    image: prom/prometheus:v3.14.0
    ports:
      - "9092:9090"
    volumes:
      - ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml:ro
      - ./prometheus/alert_rules.yml:/etc/prometheus/alert_rules.yml:ro
      - prometheus-data:/prometheus
    command:
      - --config.file=/etc/prometheus/prometheus.yml
      - --storage.tsdb.retention.time=15d
    restart: unless-stopped

  alertmanager:
    image: prom/alertmanager:v0.28.1
    ports:
      - "9094:9093"
    volumes:
      - ./alertmanager/alertmanager.yml:/etc/alertmanager/alertmanager.yml:ro
    command:
      - --config.file=/etc/alertmanager/alertmanager.yml
    extra_hosts:
      - "host.docker.internal:host-gateway"
    restart: unless-stopped

  grafana:
    image: grafana/grafana:13.2.1
    ports:
      - "3082:3000"
    environment:
      GF_SECURITY_ADMIN_USER: admin
      GF_SECURITY_ADMIN_PASSWORD: admin
      GF_USERS_ALLOW_SIGN_UP: "false"
      # スクリーンショット用にログインなしの閲覧を許可(自宅 PC 内に閉じている前提)
      GF_AUTH_ANONYMOUS_ENABLED: "true"
      GF_AUTH_ANONYMOUS_ORG_ROLE: Viewer
    volumes:
      - ./grafana/provisioning:/etc/grafana/provisioning:ro
      - ./grafana/dashboards:/var/lib/grafana/dashboards:ro
      - grafana-data:/var/lib/grafana
    depends_on:
      - prometheus
      - loki
      - tempo
    restart: unless-stopped

  mcp-grafana:
    image: grafana/mcp-grafana:1.4.1
    ports:
      - "8097:8000"
    environment:
      GRAFANA_URL: http://grafana:3000
      # 実験なので admin の基本認証。本番はサービスアカウントのトークンにする
      GRAFANA_USERNAME: admin
      GRAFANA_PASSWORD: admin
    command:
      - -t
      - streamable-http
      - --address
      - 0.0.0.0:8000
      # エージェントは読み取り専用。annotation の書き戻しは受け口が REST API で行う
      - --disable-write
      # 調査に使わないツール群を落として、ツール定義のトークン量を減らす(64 個 → 20 個弱)。
      # 残すのは alerting / annotations / datasource / prometheus / loki / proxied(Tempo の MCP)
      - --disable-admin
      - --disable-agento11y
      - --disable-api
      - --disable-asserts
      - --disable-assistant
      - --disable-cloudwatch
      - --disable-config
      - --disable-dashboard
      - --disable-docs
      - --disable-elasticsearch
      - --disable-examples
      - --disable-folder
      - --disable-graphite
      - --disable-incident
      - --disable-influxdb
      - --disable-navigation
      - --disable-oncall
      - --disable-plugin
      - --disable-provisioning
      - --disable-pyroscope
      - --disable-quickwit
      - --disable-rendering
      - --disable-runpanelquery
      - --disable-search
      - --disable-sift
      - --disable-snapshot
      - --disable-snowflake
      - --disable-sql
      - --disable-user
      - --allowed-hosts
      - ${MCP_ALLOWED_HOSTS:-*}
      - --allowed-origins
      - ${MCP_ALLOWED_ORIGINS:-*}
    depends_on:
      - grafana
    restart: unless-stopped

  # ローカル LLM 用。Claude Code の Anthropic 形式を llama.cpp(OpenAI 互換)へ変換する
  litellm:
    image: ghcr.io/berriai/litellm:main-stable
    ports:
      - "4000:4000"
    volumes:
      - ./litellm/config.yaml:/app/config.yaml:ro
    command:
      - --config
      - /app/config.yaml
      - --port
      - "4000"
    extra_hosts:
      - "host.docker.internal:host-gateway"
    restart: unless-stopped

  # 社内チャットの代わり。Alertmanager の発火通知と受け口の調査レポートを #alerts に投稿する
  mattermost-db:
    image: postgres:16-alpine
    environment:
      POSTGRES_USER: mmuser
      POSTGRES_PASSWORD: mmpass
      POSTGRES_DB: mattermost
    volumes:
      - mattermost-db:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U mmuser -d mattermost"]
      interval: 5s
      timeout: 3s
      retries: 20
    restart: unless-stopped

  mattermost:
    image: mattermost/mattermost-team-edition:11.11.0
    ports:
      - "8065:8065"
    environment:
      MM_SQLSETTINGS_DRIVERNAME: postgres
      MM_SQLSETTINGS_DATASOURCE: postgres://mmuser:mmpass@mattermost-db:5432/mattermost?sslmode=disable&connect_timeout=10
      MM_SERVICESETTINGS_SITEURL: ${MM_SITEURL:-http://localhost:8065}
      MM_SERVICESETTINGS_ENABLEINCOMINGWEBHOOKS: "true"
      MM_SERVICESETTINGS_ENABLEPOSTUSERNAMEOVERRIDE: "true"
      MM_SERVICESETTINGS_ENABLEPOSTICONOVERRIDE: "true"
      MM_EMAILSETTINGS_REQUIREEMAILVERIFICATION: "false"
      MM_PLUGINSETTINGS_ENABLEMARKETPLACE: "false"
      MM_LOGSETTINGS_CONSOLELEVEL: ERROR
    volumes:
      - mattermost-config:/mattermost/config
      - mattermost-data:/mattermost/data
      - mattermost-logs:/mattermost/logs
      - mattermost-plugins:/mattermost/plugins
      - mattermost-client-plugins:/mattermost/client/plugins
    depends_on:
      mattermost-db:
        condition: service_healthy
    restart: unless-stopped

volumes:
  mattermost-db:
  mattermost-config:
  mattermost-data:
  mattermost-logs:
  mattermost-plugins:
  mattermost-client-plugins:
  prometheus-data:
  grafana-data:
  loki-data:
  tempo-data:
otel-collector/config.yaml OTLP を Tempo / Loki / Prometheus に振り分ける
# アプリからの OTLP を受けて 3 つの保存先に振り分ける。
#   トレース   -> Tempo(OTLP gRPC)
#   ログ       -> Loki(OTLP HTTP。service.name などがインデックスのラベルになる)
#   メトリクス -> :8889 で Prometheus 形式で公開し、Prometheus がスクレイプする
receivers:
  otlp:
    protocols:
      grpc:
        endpoint: 0.0.0.0:4317
      http:
        endpoint: 0.0.0.0:4318

processors:
  batch:
    timeout: 2s

exporters:
  otlp/tempo:
    endpoint: tempo:4317
    tls:
      insecure: true
  otlphttp/loki:
    endpoint: http://loki:3100/otlp
  prometheus:
    endpoint: 0.0.0.0:8889
    # service.name などのリソース属性を service_name のようなラベルにも展開する
    resource_to_telemetry_conversion:
      enabled: true

service:
  pipelines:
    traces:
      receivers: [otlp]
      processors: [batch]
      exporters: [otlp/tempo]
    metrics:
      receivers: [otlp]
      processors: [batch]
      exporters: [prometheus]
    logs:
      receivers: [otlp]
      processors: [batch]
      exporters: [otlphttp/loki]
prometheus/prometheus.yml Collector の :8889 をスクレイプ。honor_labels で service_name を残す
global:
  scrape_interval: 15s
  evaluation_interval: 15s

rule_files:
  - /etc/prometheus/alert_rules.yml

alerting:
  alertmanagers:
    - static_configs:
        - targets: ["alertmanager:9093"]

scrape_configs:
  - job_name: prometheus
    static_configs:
      - targets: ["localhost:9090"]

  # アプリのメトリクスは OTLP で Collector に届き、Collector が Prometheus 形式で公開する。
  # honor_labels: Collector が付ける job / instance(= service.name / service.instance.id)を優先する
  - job_name: otel-collector
    honor_labels: true
    static_configs:
      - targets: ["otel-collector:8889"]
prometheus/alert_rules.yml 実験 1〜3 のアラートルール
# 実験 1〜3 に対応するアラートルール。
# しきい値は人が書く。ラベルに原因のヒント(scenario など)を入れない。
# 実験 1(負荷)と実験 3(下流の DB の遅延)は同じ JobApiLatencyHigh を鳴らし、原因だけが違う。
groups:
  - name: job-system
    rules:
      - alert: JobApiLatencyHigh
        expr: histogram_quantile(0.95, sum(rate(jobs_duration_seconds_bucket[2m])) by (le)) > 0.5
        for: 1m
        labels:
          severity: warning
        annotations:
          summary: "ジョブ投入 API の p95 レイテンシが 1 分以上 0.5 秒を超えている"
          description: "直近 2 分の p95 = {{ $value | humanizeDuration }}"

      - alert: JobErrorRateHigh
        expr: sum(rate(jobs_total{outcome="failed"}[2m])) / sum(rate(jobs_total[2m])) > 0.2
        for: 1m
        labels:
          severity: warning
        annotations:
          summary: "ジョブ投入の失敗率が 1 分以上 20% を超えている"
          description: "直近 2 分の失敗率 = {{ $value | humanizePercentage }}"

      - alert: ScrapeTargetDown
        expr: up == 0
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: "スクレイプ対象 {{ $labels.job }} が応答しない"
          description: "{{ $labels.instance }} の up が 1 分以上 0"
alertmanager/alertmanager.tmpl.yml 受け口への webhook と Mattermost への通知(setup_mattermost.py が alertmanager.yml を生成)
# alertmanager.yml のひな型。setup_mattermost.py が {{MATTERMOST_HOOK_ID}} を埋めて alertmanager.yml を書く。
# 発火したアラートは 2 か所へ送る:
#   1. docker ホストの受け口(investigate.py、ポート 9095)。AI 調査を起動する
#   2. Mattermost の #alerts(Slack 互換の incoming webhook)。運用者への通知
route:
  receiver: ai-investigator
  group_by: ["alertname"]
  group_wait: 10s
  group_interval: 1m
  repeat_interval: 4h

receivers:
  - name: ai-investigator
    webhook_configs:
      - url: http://host.docker.internal:9095/alert
        send_resolved: false
    slack_configs:
      - api_url: http://mattermost:8065/hooks/{{MATTERMOST_HOOK_ID}}
        channel: "#alerts"
        username: alertmanager
        send_resolved: true
        title: "{{ .CommonLabels.alertname }} [{{ .Status | toUpper }}]"
        text: '{{ range .Alerts }}{{ .Annotations.summary }}{{ if .Annotations.description }}: {{ .Annotations.description }}{{ end }}{{ "\n" }}{{ end }}'
        color: '{{ if eq .Status "firing" }}danger{{ else }}good{{ end }}'
loki/config.yaml 単一ノードの Loki。ruler を Alertmanager に向ける
# 単一ノードの Loki。OTLP で受けたログをローカルディスクに保存する。
# ruler を有効にして、ログ由来のアラート(rules/fake/*.yml)を Alertmanager に送る。
auth_enabled: false

server:
  http_listen_port: 3100
  grpc_listen_port: 9096

common:
  path_prefix: /loki
  storage:
    filesystem:
      chunks_directory: /loki/chunks
      rules_directory: /loki/rules
  replication_factor: 1
  ring:
    kvstore:
      store: inmemory

schema_config:
  configs:
    - from: "2024-04-01"
      store: tsdb
      object_store: filesystem
      schema: v13
      index:
        prefix: index_
        period: 24h

limits_config:
  allow_structured_metadata: true
  volume_enabled: true

ruler:
  storage:
    type: local
    local:
      directory: /loki/rules
  rule_path: /tmp/loki-rules
  alertmanager_url: http://alertmanager:9093
  ring:
    kvstore:
      store: inmemory
  enable_api: true
loki/rules/fake/job-system.yml ログ由来のアラート(実験 2b)
# ログ由来のアラート(実験 2b)。メトリクスには何も出ず、ログの行数だけで鳴る。
# "fake" は auth_enabled: false のときのテナント名。
groups:
  - name: job-system
    rules:
      - alert: SlotsRunningLow
        expr: sum(count_over_time({service_name="slot-manager"} |= "slots running low" [5m])) > 3
        for: 1m
        labels:
          severity: warning
        annotations:
          summary: "計算枠が少ないという警告ログが 5 分で 3 行を超えた"
tempo/config.yaml MCP サーバを有効にした Tempo
# 単一ノードの Tempo。OTLP で受けたトレースをローカルディスクに保存する。
# query_frontend.mcp_server を有効にすると http://tempo:3200/api/mcp に
# MCP サーバ(traceql-search / get-trace など)が立つ。mcp-grafana には Tempo ツールが無いので併用する。
server:
  http_listen_port: 3200

distributor:
  receivers:
    otlp:
      protocols:
        grpc:
          endpoint: 0.0.0.0:4317
        http:
          endpoint: 0.0.0.0:4318

storage:
  trace:
    backend: local
    local:
      path: /var/tempo/blocks
    wal:
      path: /var/tempo/wal

query_frontend:
  mcp_server:
    enabled: true
grafana/provisioning/datasources/datasources.yml Prometheus / Loki / Tempo のデータソース
apiVersion: 1
datasources:
  - name: Prometheus
    type: prometheus
    uid: prometheus
    access: proxy
    url: http://prometheus:9090
    isDefault: true

  - name: Loki
    type: loki
    uid: loki
    access: proxy
    url: http://loki:3100

  - name: Tempo
    type: tempo
    uid: tempo
    access: proxy
    url: http://tempo:3200
    jsonData:
      tracesToLogsV2:
        datasourceUid: loki
        spanStartTimeShift: "-5m"
        spanEndTimeShift: "5m"
        filterByTraceID: true
      serviceMap:
        datasourceUid: prometheus
grafana/provisioning/dashboards/dashboards.yml ダッシュボードの自動登録
apiVersion: 1
providers:
  - name: grafana-alert
    folder: ""
    type: file
    options:
      path: /var/lib/grafana/dashboards
grafana/dashboards/job-system.json ダッシュボード(p95、投入レート、失敗率、ログ、up。annotation を重ねる)
{
  "uid": "job-system",
  "title": "Job system",
  "tags": ["grafana-alert"],
  "timezone": "browser",
  "schemaVersion": 39,
  "refresh": "10s",
  "time": { "from": "now-30m", "to": "now" },
  "annotations": {
    "list": [
      {
        "name": "AI investigation",
        "datasource": { "type": "grafana", "uid": "-- Grafana --" },
        "enable": true,
        "iconColor": "orange",
        "target": {
          "type": "tags",
          "tags": ["ai-investigation"],
          "limit": 100,
          "matchAny": true
        }
      }
    ]
  },
  "panels": [
    {
      "id": 1,
      "type": "timeseries",
      "title": "ジョブ投入 p95 レイテンシ(アラート式)",
      "gridPos": { "x": 0, "y": 0, "w": 12, "h": 8 },
      "datasource": { "type": "prometheus", "uid": "prometheus" },
      "fieldConfig": {
        "defaults": {
          "unit": "s",
          "thresholds": {
            "mode": "absolute",
            "steps": [
              { "color": "green", "value": null },
              { "color": "red", "value": 0.5 }
            ]
          },
          "custom": { "thresholdsStyle": { "mode": "line" } }
        },
        "overrides": []
      },
      "targets": [
        {
          "refId": "A",
          "expr": "histogram_quantile(0.95, sum(rate(jobs_duration_seconds_bucket[2m])) by (le))",
          "legendFormat": "p95"
        }
      ]
    },
    {
      "id": 2,
      "type": "timeseries",
      "title": "ジョブ投入レート(結末別)",
      "gridPos": { "x": 12, "y": 0, "w": 12, "h": 8 },
      "datasource": { "type": "prometheus", "uid": "prometheus" },
      "fieldConfig": { "defaults": { "unit": "reqps" }, "overrides": [] },
      "targets": [
        {
          "refId": "A",
          "expr": "sum(rate(jobs_total[1m])) by (outcome)",
          "legendFormat": "{{outcome}}"
        }
      ]
    },
    {
      "id": 3,
      "type": "timeseries",
      "title": "ジョブ投入の失敗率(アラート式)",
      "gridPos": { "x": 0, "y": 8, "w": 12, "h": 8 },
      "datasource": { "type": "prometheus", "uid": "prometheus" },
      "fieldConfig": {
        "defaults": {
          "unit": "percentunit",
          "min": 0,
          "max": 1,
          "thresholds": {
            "mode": "absolute",
            "steps": [
              { "color": "green", "value": null },
              { "color": "red", "value": 0.2 }
            ]
          },
          "custom": { "thresholdsStyle": { "mode": "line" } }
        },
        "overrides": []
      },
      "targets": [
        {
          "refId": "A",
          "expr": "sum(rate(jobs_total{outcome=\"failed\"}[2m])) / sum(rate(jobs_total[2m]))",
          "legendFormat": "failed"
        }
      ]
    },
    {
      "id": 4,
      "type": "timeseries",
      "title": "サービス別 HTTP p95(自動計装)",
      "gridPos": { "x": 12, "y": 8, "w": 12, "h": 8 },
      "datasource": { "type": "prometheus", "uid": "prometheus" },
      "fieldConfig": { "defaults": { "unit": "s" }, "overrides": [] },
      "targets": [
        {
          "refId": "A",
          "expr": "histogram_quantile(0.95, sum(rate(http_server_request_duration_seconds_bucket[2m])) by (le, service_name))",
          "legendFormat": "{{service_name}}"
        }
      ]
    },
    {
      "id": 5,
      "type": "logs",
      "title": "WARNING / ERROR ログ",
      "gridPos": { "x": 0, "y": 16, "w": 16, "h": 9 },
      "datasource": { "type": "loki", "uid": "loki" },
      "options": {
        "showTime": true,
        "sortOrder": "Descending",
        "wrapLogMessage": true
      },
      "targets": [
        {
          "refId": "A",
          "expr": "{service_name=~\"job-api|slot-manager\"} | detected_level=~\"warn|error\"",
          "queryType": "range"
        }
      ]
    },
    {
      "id": 6,
      "type": "stat",
      "title": "スクレイプ対象 up",
      "gridPos": { "x": 16, "y": 16, "w": 8, "h": 9 },
      "datasource": { "type": "prometheus", "uid": "prometheus" },
      "fieldConfig": {
        "defaults": {
          "thresholds": {
            "mode": "absolute",
            "steps": [
              { "color": "red", "value": null },
              { "color": "green", "value": 1 }
            ]
          }
        },
        "overrides": []
      },
      "options": {
        "reduceOptions": { "calcs": ["lastNotNull"] },
        "textMode": "value_and_name"
      },
      "targets": [{ "refId": "A", "expr": "up", "legendFormat": "{{job}}" }]
    }
  ]
}
litellm/config.yaml Claude Code からローカル LLM への変換(記事の計測では未使用)
# ローカル LLM を Claude Code から使うための変換プロキシ。
# Claude Code は Anthropic Messages API(/v1/messages)を話すので、LiteLLM で受けて
# OpenAI 互換の llama.cpp(既存スタックの :8099)へ流す。
#   claude 側: ANTHROPIC_BASE_URL=http://localhost:4000  ANTHROPIC_AUTH_TOKEN=sk-local  --model qwen3.6-35b-a3b
model_list:
  - model_name: qwen3.6-35b-a3b
    litellm_params:
      model: openai/qwen3.6-35b-a3b
      api_base: http://host.docker.internal:8099/v1
      api_key: none
      timeout: 600

litellm_settings:
  drop_params: true
  set_verbose: false

general_settings:
  master_key: sk-local

監視対象(app/)

app/Dockerfile
FROM python:3.13-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
app/requirements.txt
fastapi
uvicorn
httpx
opentelemetry-sdk
opentelemetry-exporter-otlp-proto-grpc
opentelemetry-instrumentation-fastapi
opentelemetry-instrumentation-httpx
app/telemetry.py OTLP の配線。トレース・メトリクス・ログをまとめて Collector へ
"""OTel の配線。トレース・メトリクス・ログをまとめて OTLP(gRPC)で Collector に送る。

送り先は環境変数 OTEL_EXPORTER_OTLP_ENDPOINT(compose では otel-collector:4317)。
ログは標準 logging に LoggingHandler を足すだけで、trace_id / span_id が自動で付く。
"""

from __future__ import annotations

import logging
import os
import socket

from opentelemetry import metrics, trace
from opentelemetry._logs import set_logger_provider
from opentelemetry.exporter.otlp.proto.grpc._log_exporter import OTLPLogExporter
from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor


def setup(service_name: str) -> tuple[trace.Tracer, metrics.Meter]:
    resource = Resource.create({
        "service.name": service_name,
        "service.version": os.environ.get("APP_VERSION", "0.1.0"),
        "service.instance.id": socket.gethostname(),
    })

    tp = TracerProvider(resource=resource)
    tp.add_span_processor(BatchSpanProcessor(OTLPSpanExporter()))
    trace.set_tracer_provider(tp)

    mp = MeterProvider(
        resource=resource,
        metric_readers=[PeriodicExportingMetricReader(OTLPMetricExporter(), export_interval_millis=5000)],
    )
    metrics.set_meter_provider(mp)

    lp = LoggerProvider(resource=resource)
    lp.add_log_record_processor(BatchLogRecordProcessor(OTLPLogExporter()))
    set_logger_provider(lp)

    # 標準出力(docker logs)にも同じ行を出す
    logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s")
    logging.getLogger().addHandler(LoggingHandler(level=logging.NOTSET, logger_provider=lp))

    return trace.get_tracer(service_name), metrics.get_meter(service_name)
app/job_api.py ジョブの受付。slot-manager を呼ぶ
"""job-api サービス。計算ジョブの投入を受けて slot-manager に計算枠の確保を頼む。

- POST /jobs {"user", "job", "size"}: slot-manager の /reserve を呼び、結果を返す。
  メトリクス jobs(結末別の件数)と jobs.duration(所要秒)を記録する
- POST /admin/fault: 故障注入。{"duplicate_calls": 8} で 1 ジョブあたり /reserve を 8 回呼ぶ(冪等でない再送のバグ)。
  {"reset": true} で戻す
- GET /admin/fault、GET /health

メトリクスのラベルは outcome(ok / failed)だけ。ジョブ種別やエラー内容はログとスパン属性にだけ出る。
"""

from __future__ import annotations

import itertools
import logging
import os
import time

import httpx
from fastapi import FastAPI
from fastapi.responses import JSONResponse
from opentelemetry import trace
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor

import telemetry

tracer, meter = telemetry.setup("job-api")
log = logging.getLogger("job-api")

app = FastAPI(title="job-api")
FastAPIInstrumentor.instrument_app(app, excluded_urls="health,admin")
HTTPXClientInstrumentor().instrument()

SLOT_MANAGER_URL = os.environ.get("SLOT_MANAGER_URL", "http://slot-manager:8000").rstrip("/")
FAULT: dict[str, int] = {"duplicate_calls": 1}
# 種別ごとの見積もり時間(秒)。ログに載せるだけ
ESTIMATE_S = {"mip": 600, "lp": 60, "simulation": 300, "ml": 1800}

jobs_counter = meter.create_counter("jobs", description="ジョブ投入の結末別の件数")
jobs_duration = meter.create_histogram(
    "jobs.duration", unit="s", description="ジョブ投入 1 件の所要時間(枠の確保まで)",
    explicit_bucket_boundaries_advisory=[0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10],
)
_seq = itertools.count(1000)
_client = httpx.AsyncClient(timeout=30.0, limits=httpx.Limits(max_connections=200, max_keepalive_connections=50))


@app.post("/jobs")
async def submit_job(body: dict):
    t0 = time.perf_counter()
    job_id = f"J-{next(_seq)}"
    job = body.get("job", "")
    size = int(body.get("size", 1))
    user = body.get("user", "anonymous")
    attrs = {"job.id": job_id, "user.id": user, "job.type": job, "job.size": size}
    span = trace.get_current_span()
    for k, v in attrs.items():
        span.set_attribute(k, v)
    log.info("job received", extra=attrs)

    outcome, status, detail = "ok", 200, None
    try:
        for _ in range(max(1, FAULT["duplicate_calls"])):
            r = await _client.post(f"{SLOT_MANAGER_URL}/reserve", json={"job": job, "size": size})
        if r.status_code != 200:
            outcome, status = "failed", r.status_code
            try:
                detail = r.json().get("error", r.text)
            except ValueError:
                detail = r.text
    except httpx.HTTPError as e:
        outcome, status, detail = "failed", 502, f"slot-manager unreachable: {type(e).__name__}: {e}"

    elapsed = time.perf_counter() - t0
    jobs_counter.add(1, {"outcome": outcome})
    jobs_duration.record(elapsed, {"outcome": outcome})
    span.set_attribute("job.outcome", outcome)

    if outcome == "ok":
        log.info("job submitted", extra={**attrs, "job.estimate_s": ESTIMATE_S.get(job, 0)})
        return {"job_id": job_id, "job": job, "size": size, "estimate_s": ESTIMATE_S.get(job, 0)}
    log.error("job failed: %s", detail, extra={**attrs, "http.response.status_code": status})
    return JSONResponse({"job_id": job_id, "error": detail}, status_code=status)


@app.get("/health")
async def health():
    return {"status": "ok"}


@app.get("/admin/fault")
async def get_fault():
    return {"fault": FAULT, "slot_manager_url": SLOT_MANAGER_URL}


@app.post("/admin/fault")
async def set_fault(body: dict):
    if body.get("reset"):
        FAULT["duplicate_calls"] = 1
    if "duplicate_calls" in body:
        FAULT["duplicate_calls"] = int(body["duplicate_calls"])
    return await get_fault()
app/slot_manager.py 計算枠 DB アクセスを sleep と接続プールで模擬
"""slot-manager サービス。計算ジョブごとに計算枠(ライセンス席、GPU/CPU スロット)を確保する。job-api の下流。

- POST /reserve {"job", "size"}: 枠を確認して確保する。 DB へのアクセスは sleep で模擬し、
  接続プール(DB_POOL 本)を取り合う。プール待ちは db.pool.acquire、問い合わせは db.query のスパンになる
- POST /admin/fault: 故障注入。{"latency_ms": 300} で DB を遅くする、{"slots": {"mip": 0}} で枠数を書き換える、
  {"error_rate": 0.5} で DB エラーを混ぜる。{"reset": true} で全部戻す
- GET /admin/fault: 現在の故障設定
- GET /health

枠切れは 409 とエラーログ「failed to reserve: no free slot: <job>」で表す。
メトリクスにはジョブ種別のラベルを付けない(どの種別かはログにしか出ない)。
"""

from __future__ import annotations

import asyncio
import logging
import os
import random

from fastapi import FastAPI
from fastapi.responses import JSONResponse
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor

import telemetry

tracer, meter = telemetry.setup("slot-manager")
log = logging.getLogger("slot-manager")

app = FastAPI(title="slot-manager")
FastAPIInstrumentor.instrument_app(app, excluded_urls="health,admin")

BIG = 10**9
JOB_TYPES = ("mip", "lp", "simulation", "ml")
SLOTS: dict[str, int] = {job: BIG for job in JOB_TYPES}
FAULT: dict[str, float] = {"latency_ms": 0, "error_rate": 0.0}
DB_BASE_MS = int(os.environ.get("DB_BASE_MS", "40"))
DB_POOL = int(os.environ.get("DB_POOL", "4"))
LOW_WATERMARK = 50

_pool = asyncio.Semaphore(DB_POOL)


async def db_query(operation: str) -> None:
    """DB 呼び出しの模擬。プール待ちと問い合わせ本体を別のスパンにする。"""
    with tracer.start_as_current_span("db.pool.acquire"):
        await _pool.acquire()
    try:
        with tracer.start_as_current_span("db.query") as span:
            span.set_attribute("db.system.name", "sqlite")
            span.set_attribute("db.operation.name", operation)
            await asyncio.sleep((DB_BASE_MS + FAULT["latency_ms"]) / 1000)
    finally:
        _pool.release()


@app.post("/reserve")
async def reserve(body: dict):
    job = body.get("job", "")
    size = int(body.get("size", 1))
    if job not in SLOTS:
        return JSONResponse({"error": f"unknown job type: {job}"}, status_code=400)

    await db_query("SELECT slots")
    if random.random() < FAULT["error_rate"]:
        log.error("db connection reset while reserving", extra={"job.type": job})
        return JSONResponse({"error": "db error"}, status_code=500)
    if SLOTS[job] < size:
        log.error(
            "failed to reserve: no free slot: %s", job,
            extra={"job.type": job, "job.size": size, "slots.remaining": SLOTS[job]},
        )
        return JSONResponse({"error": "no free slot"}, status_code=409)

    SLOTS[job] -= size
    await db_query("UPDATE slots")
    log.info("reserved: %s x%d", job, size, extra={"job.type": job, "job.size": size, "slots.remaining": SLOTS[job]})
    if SLOTS[job] < LOW_WATERMARK:
        log.warning(
            "slots running low: %s", job,
            extra={"job.type": job, "slots.remaining": SLOTS[job]},
        )
    return {"job": job, "remaining": SLOTS[job]}


@app.get("/health")
async def health():
    return {"status": "ok"}


@app.get("/admin/fault")
async def get_fault():
    return {"fault": FAULT, "slots": {k: v for k, v in SLOTS.items()}, "db_base_ms": DB_BASE_MS, "db_pool": DB_POOL}


@app.post("/admin/fault")
async def set_fault(body: dict):
    if body.get("reset"):
        FAULT.update({"latency_ms": 0, "error_rate": 0.0})
        for job in JOB_TYPES:
            SLOTS[job] = BIG
    if "latency_ms" in body:
        FAULT["latency_ms"] = float(body["latency_ms"])
    if "error_rate" in body:
        FAULT["error_rate"] = float(body["error_rate"])
    for job, n in (body.get("slots") or {}).items():
        if job in SLOTS:
            SLOTS[job] = int(n)
    return await get_fault()

受け口とエージェント

investigate.py Alertmanager の webhook を受けて claude -p(または local_agent.py)を起動し、レポート・annotation・Mattermost 投稿を行う
"""Alertmanager の webhook を受けて Claude Code(headless)に調査させる受け口。

- POST /alert に Alertmanager の webhook ペイロードが届く
- firing のアラートごとに `claude -p` を起動し、Grafana MCP(読み取り専用。Prometheus / Loki / Tempo)で調べさせる
- 結果を reports/ に Markdown で保存し、Grafana の annotation に書き戻す

自作するのはこのファイルだけ。標準ライブラリのみで動く。
"""

from __future__ import annotations

import base64
import json
import os
import shlex
import subprocess
import sys
import threading
import time
import urllib.error
import urllib.request
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path

HERE = Path(__file__).resolve().parent
PORT = int(os.environ.get("PORT", "9095"))
CLAUDE_BIN = os.environ.get("CLAUDE_BIN", "claude")
CLAUDE_MODEL = os.environ.get("CLAUDE_MODEL", "sonnet")
MCP_CONFIG = os.environ.get("MCP_CONFIG", str(HERE / "mcp.json"))
PROMPT_FILE = os.environ.get("PROMPT_FILE", str(HERE / "PROMPT.md"))
REPORT_DIR = Path(os.environ.get("REPORT_DIR", str(HERE / "reports")))
GRAFANA_URL = os.environ.get("GRAFANA_URL", "http://localhost:3082").rstrip("/")
GRAFANA_SA_TOKEN = os.environ.get("GRAFANA_SA_TOKEN", "")
# トークンが無ければ基本認証(実験用。既定は admin:admin)
GRAFANA_BASIC_AUTH = os.environ.get("GRAFANA_BASIC_AUTH", "admin:admin")
# エージェントに許すツール。mcp-grafana の MCP だけ(Tempo のツールも mcp-grafana が tempo_* として中継する)
ALLOWED_TOOLS = os.environ.get("ALLOWED_TOOLS", "mcp__grafana__*")
# 同じアラート(fingerprint)はこの秒数の間は再調査しない
DEDUP_SECONDS = int(os.environ.get("DEDUP_SECONDS", "1800"))
# Mattermost の incoming webhook。未設定なら .env の MATTERMOST_HOOK_ID から組み立てる(setup_mattermost.py が書く)
MATTERMOST_URL = os.environ.get("MATTERMOST_URL", "http://localhost:8065").rstrip("/")
MATTERMOST_WEBHOOK_URL = os.environ.get("MATTERMOST_WEBHOOK_URL", "")
# 投稿先チャンネル。モデルごとに分けるときは受け口の起動時に変える(例: sonnet / qwen)
MATTERMOST_CHANNEL = os.environ.get("MATTERMOST_CHANNEL", "alerts")
MATTERMOST_MAX_CHARS = 16000  # Mattermost の投稿上限は 16,383 文字
# Claude Code の組み込みツール。エージェントには MCP ツールだけを残す
BUILTIN_TOOLS = "Bash,PowerShell,Edit,Write,MultiEdit,NotebookEdit,Read,Glob,Grep,LS,WebFetch,WebSearch,Agent,TodoWrite,Skill"

_recent: dict[str, float] = {}
_lock = threading.Lock()


def log(msg: str) -> None:
    print(f"[{datetime.now().strftime('%H:%M:%S')}] {msg}", flush=True)


def build_prompt(alert: dict) -> str:
    template = Path(PROMPT_FILE).read_text(encoding="utf-8")
    return template.replace("{alert_json}", json.dumps(alert, ensure_ascii=False, indent=2))


def run_claude(prompt: str) -> dict:
    """claude -p を起動して JSON 結果を返す。組み込みツールは拒否し、MCP ツールだけ許可する。"""
    cmd = [
        *shlex.split(CLAUDE_BIN, posix=(os.name != "nt")),
        "-p",
        "--model", CLAUDE_MODEL,
        "--output-format", "json",
        "--mcp-config", MCP_CONFIG,
        "--strict-mcp-config",
        "--allowedTools", ALLOWED_TOOLS,
        # 組み込みツールは名前で全部拒否する(--tools "" だと MCP ツールの許可まで消える)
        "--disallowedTools", BUILTIN_TOOLS,
        "--no-session-persistence",
    ]
    started = time.monotonic()
    proc = subprocess.run(cmd, input=prompt, capture_output=True, text=True, encoding="utf-8")
    elapsed = time.monotonic() - started
    try:
        data = json.loads(proc.stdout)
    except json.JSONDecodeError:
        data = None
    if proc.returncode != 0 and not (isinstance(data, dict) and data.get("result")):
        raise RuntimeError(f"claude exited {proc.returncode}: {proc.stderr[-1500:]} {proc.stdout[-1500:]}")
    if proc.returncode != 0:
        # API エラーなどで途中終了しても、結果(エラー文)を残す。ローカルモデルのコンテキスト超過の記録用
        log(f"claude exited {proc.returncode}(結果は保存する): {str(data.get('result'))[:200]}")
        data["is_error"] = True
    data.setdefault("duration_ms", int(elapsed * 1000))
    data["_stderr"] = proc.stderr
    return data


def grafana_auth_header() -> str:
    if GRAFANA_SA_TOKEN:
        return f"Bearer {GRAFANA_SA_TOKEN}"
    return "Basic " + base64.b64encode(GRAFANA_BASIC_AUTH.encode("utf-8")).decode("ascii")


def post_annotation(alert: dict, report_text: str) -> None:
    labels = alert.get("labels", {})
    started = alert.get("startsAt", "")
    try:
        t = int(datetime.fromisoformat(started.replace("Z", "+00:00")).timestamp() * 1000)
    except ValueError:
        t = int(time.time() * 1000)
    body = {
        "time": t,
        "tags": ["ai-investigation", labels.get("alertname", "unknown")],
        "text": f"[{labels.get('alertname')}] AI 調査\n\n" + report_text[:1500],
    }
    req = urllib.request.Request(
        f"{GRAFANA_URL}/api/annotations",
        data=json.dumps(body).encode("utf-8"),
        headers={
            "Content-Type": "application/json",
            "Authorization": grafana_auth_header(),
        },
        method="POST",
    )
    try:
        with urllib.request.urlopen(req, timeout=10) as resp:
            log(f"annotation 作成: HTTP {resp.status}")
    except urllib.error.URLError as e:
        log(f"annotation 作成に失敗: {e}")


def mattermost_webhook_url() -> str:
    if MATTERMOST_WEBHOOK_URL:
        return MATTERMOST_WEBHOOK_URL
    env = HERE / ".env"
    if env.exists():
        for line in env.read_text(encoding="utf-8").splitlines():
            if line.startswith("MATTERMOST_HOOK_ID="):
                return f"{MATTERMOST_URL}/hooks/{line.split('=', 1)[1].strip()}"
    return ""


def post_mattermost(alert: dict, report_text: str, result: dict) -> None:
    """調査レポートの全文を Mattermost の #alerts に投稿する(Markdown はそのまま描画される)"""
    url = mattermost_webhook_url()
    if not url:
        log("Mattermost の webhook が未設定なので投稿しない")
        return
    labels = alert.get("labels", {})
    head = (f"### [{labels.get('alertname')}] AI 調査レポート\n"
            f"発火 {alert.get('startsAt', '')} / モデル {CLAUDE_MODEL} / "
            f"所要 {result.get('duration_ms', 0) / 1000:.0f} 秒、{result.get('num_turns')} ターン\n\n")
    text = head + report_text
    if len(text) > MATTERMOST_MAX_CHARS:
        text = text[:MATTERMOST_MAX_CHARS] + "\n\n…(上限のため省略)"
    body = {"channel": MATTERMOST_CHANNEL, "username": f"ai-investigator ({CLAUDE_MODEL})", "text": text}
    req = urllib.request.Request(url, data=json.dumps(body, ensure_ascii=False).encode("utf-8"),
                                 headers={"Content-Type": "application/json"}, method="POST")
    try:
        with urllib.request.urlopen(req, timeout=10) as resp:
            log(f"Mattermost 投稿: HTTP {resp.status}")
    except (urllib.error.URLError, OSError) as e:
        log(f"Mattermost 投稿に失敗: {e}")


def investigate(alert: dict) -> None:
    labels = alert.get("labels", {})
    name = labels.get("alertname", "unknown")
    stamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ")
    log(f"調査開始: {name} {labels}")
    try:
        result = run_claude(build_prompt(alert))
    except Exception as e:  # noqa: BLE001
        log(f"調査失敗: {name}: {e}")
        return
    report = result.get("result", "")
    REPORT_DIR.mkdir(parents=True, exist_ok=True)
    path = REPORT_DIR / f"{stamp}_{name}.md"
    header = (
        f"# {name}\n\n"
        f"- 開始: {alert.get('startsAt')}\n"
        f"- ラベル: `{json.dumps(labels, ensure_ascii=False)}`\n"
        f"- モデル: {CLAUDE_MODEL}\n"
        f"- 所要: {result.get('duration_ms', 0) / 1000:.1f} 秒、ターン数: {result.get('num_turns')}\n\n"
    )
    path.write_text(header + report + "\n", encoding="utf-8")
    if result.get("_stderr"):
        # エージェントの標準エラー(local_agent.py ならツール呼び出しの記録)を並べて残す
        path.with_suffix(".log").write_text(result["_stderr"], encoding="utf-8")
    with (REPORT_DIR / "results.jsonl").open("a", encoding="utf-8") as f:
        f.write(json.dumps({
            "time": stamp,
            "alertname": name,
            "model": CLAUDE_MODEL,
            "duration_ms": result.get("duration_ms"),
            "num_turns": result.get("num_turns"),
            "is_error": bool(result.get("is_error")),
            "usage": result.get("usage"),
            "report": str(path.name),
        }, ensure_ascii=False) + "\n")
    log(f"レポート保存: {path}")
    post_annotation(alert, report)
    post_mattermost(alert, report, result)


class Handler(BaseHTTPRequestHandler):
    def do_POST(self) -> None:  # noqa: N802
        if self.path != "/alert":
            self.send_error(404)
            return
        length = int(self.headers.get("Content-Length", "0"))
        try:
            payload = json.loads(self.rfile.read(length) or b"{}")
        except json.JSONDecodeError:
            self.send_error(400, "invalid json")
            return
        self.send_response(200)
        self.end_headers()
        now = time.time()
        for alert in payload.get("alerts", []):
            if alert.get("status") != "firing":
                continue
            key = alert.get("fingerprint") or json.dumps(alert.get("labels", {}), sort_keys=True)
            with _lock:
                if now - _recent.get(key, 0) < DEDUP_SECONDS:
                    log(f"スキップ(直近に調査済み): {alert.get('labels', {}).get('alertname')}")
                    continue
                _recent[key] = now
            threading.Thread(target=investigate, args=(alert,), daemon=True).start()

    def log_message(self, fmt: str, *args) -> None:  # noqa: D102
        return


def main() -> None:
    if hasattr(sys.stdout, "reconfigure"):
        sys.stdout.reconfigure(encoding="utf-8")  # Windows のコンソールでも日本語ログを崩さない
    server = ThreadingHTTPServer(("0.0.0.0", PORT), Handler)
    log(f"受け口を起動: http://0.0.0.0:{PORT}/alert (claude={CLAUDE_BIN}, model={CLAUDE_MODEL}, tools={ALLOWED_TOOLS})")
    try:
        server.serve_forever()
    except KeyboardInterrupt:
        pass


if __name__ == "__main__":
    sys.exit(main())
setup_mattermost.py Mattermost の初期設定(管理者、チーム、チャンネル、incoming webhook)
"""Mattermost の初期設定(標準ライブラリのみ)。

    python setup_mattermost.py [--url http://localhost:8065] [--admin-password ...]

1. 最初のユーザー(= システム管理者)admin を作る
2. チーム ops とチャンネル alerts を作る
3. alerts への incoming webhook を作る
4. webhook の id を .env(MATTERMOST_HOOK_ID)に書き、alertmanager/alertmanager.tmpl.yml から alertmanager/alertmanager.yml を生成する

再実行しても既存のものを使い回す(冪等)。終わったら `docker compose restart alertmanager` で通知先を反映する。
"""

from __future__ import annotations

import argparse
import json
import re
import sys
import time
import urllib.error
import urllib.request
from pathlib import Path

HERE = Path(__file__).resolve().parent


def api(base: str, method: str, path: str, body: dict | None = None, token: str | None = None):
    req = urllib.request.Request(f"{base}/api/v4{path}", method=method,
                                 data=json.dumps(body).encode("utf-8") if body is not None else None,
                                 headers={"Content-Type": "application/json"})
    if token:
        req.add_header("Authorization", f"Bearer {token}")
    try:
        with urllib.request.urlopen(req, timeout=30) as resp:
            raw = resp.read()
            return resp.status, (json.loads(raw) if raw else None), resp.headers
    except urllib.error.HTTPError as e:
        return e.code, json.loads(e.read() or b"{}"), e.headers


def wait_ready(base: str) -> None:
    for _ in range(60):
        try:
            with urllib.request.urlopen(f"{base}/api/v4/system/ping", timeout=5) as resp:
                if resp.status == 200:
                    return
        except (urllib.error.URLError, TimeoutError, OSError):
            pass
        time.sleep(3)
    raise SystemExit("Mattermost が起動しない")


def main() -> int:
    p = argparse.ArgumentParser()
    p.add_argument("--url", default="http://localhost:8065")
    p.add_argument("--admin-user", default="admin")
    p.add_argument("--admin-password", default="Admin-pass-1234")
    p.add_argument("--admin-email", default="admin@example.invalid")
    p.add_argument("--team", default="ops")
    p.add_argument("--channel", default="alerts")
    a = p.parse_args()
    base = a.url.rstrip("/")
    wait_ready(base)

    # 1. 管理者(最初のユーザーがシステム管理者になる)
    status, body, _ = api(base, "POST", "/users", {"email": a.admin_email, "username": a.admin_user, "password": a.admin_password})
    print(f"user: {status} {body.get('username') if isinstance(body, dict) else body}")
    status, body, headers = api(base, "POST", "/users/login", {"login_id": a.admin_user, "password": a.admin_password})
    if status != 200:
        raise SystemExit(f"login failed: {status} {body}")
    token = headers["Token"]

    # 2. チームとチャンネル
    status, team, _ = api(base, "GET", f"/teams/name/{a.team}", token=token)
    if status != 200:
        status, team, _ = api(base, "POST", "/teams", {"name": a.team, "display_name": a.team.capitalize(), "type": "O"}, token=token)
    print(f"team: {team['name']} ({team['id']})")
    status, ch, _ = api(base, "GET", f"/teams/{team['id']}/channels/name/{a.channel}", token=token)
    if status != 200:
        status, ch, _ = api(base, "POST", "/channels", {"team_id": team["id"], "name": a.channel, "display_name": a.channel, "type": "O"}, token=token)
    print(f"channel: {ch['name']} ({ch['id']})")

    # 3. incoming webhook(display_name で既存を探す)
    status, hooks, _ = api(base, "GET", f"/hooks/incoming?team_id={team['id']}&per_page=200", token=token)
    hook = next((h for h in (hooks or []) if h.get("display_name") == "alert-investigation"), None)
    if not hook:
        status, hook, _ = api(base, "POST", "/hooks/incoming", {
            "channel_id": ch["id"], "display_name": "alert-investigation",
            "description": "Alertmanager の発火通知と AI 調査レポート",
        }, token=token)
    print(f"webhook: {hook['id']}")

    # 4. .env と alertmanager.yml
    env_path = HERE / ".env"
    lines = env_path.read_text(encoding="utf-8").splitlines() if env_path.exists() else []
    lines = [l for l in lines if not l.startswith("MATTERMOST_HOOK_ID=")]
    lines.append(f"MATTERMOST_HOOK_ID={hook['id']}")
    env_path.write_text("\n".join(lines) + "\n", encoding="utf-8")
    tmpl = (HERE / "alertmanager" / "alertmanager.tmpl.yml").read_text(encoding="utf-8")
    (HERE / "alertmanager" / "alertmanager.yml").write_text(tmpl.replace("{{MATTERMOST_HOOK_ID}}", hook["id"]), encoding="utf-8")
    print(f"wrote .env and alertmanager/alertmanager.yml. next: docker compose restart alertmanager")
    print(f"login: {base}  user={a.admin_user}  channel=~{a.channel}")
    return 0


if __name__ == "__main__":
    sys.exit(main())
PROMPT.md 調査プロンプト(v2)
あなたは Grafana に接続された SRE の調査エージェントです。
Alertmanager から次のアラートが届きました。

```json
{alert_json}
```

監視対象は計算ジョブの投入システムで、`job-api` と `slot-manager` の 2 サービスからなります。
`job-api` はジョブ 1 件ごとに `slot-manager` の `POST /reserve` を呼び、`slot-manager` は DB で計算枠を確認して確保します。
3 つのデータソースが使えます。

- Prometheus(メトリクス): `jobs_total{outcome}`、`jobs_duration_seconds_bucket` など。`mcp__grafana__query_prometheus`
- Loki(ログ): `{service_name="job-api"}`、`{service_name="slot-manager"}`。`mcp__grafana__query_loki_logs`
- Tempo(トレース): `mcp__grafana__tempo_traceql-search` と `mcp__grafana__tempo_get-trace`

MCP ツール(`mcp__grafana__*`)だけを使って原因を調べてください。
ツールは読み取り専用です。設定変更や書き込みはできません。

クエリの書き方:

- PromQL の時間範囲は `startsAt` の前後 30 分。range クエリには `stepSeconds=60` を渡す
- LogQL は 1 本のストリームにフィルタを足す(例: `{service_name="slot-manager"} |~ "(?i)error|warn"`)。`or` で 2 本をつながない
- TraceQL の例: `{resource.service.name="job-api" && duration > 500ms}`。検索には開始・終了時刻を必ず渡す
- ツールがエラーを返したら、クエリを 1 度だけ書き直すか別の手段に切り替える

手順:

1. `query_prometheus` でアラート式そのものと、関連するメトリクスをアラートの `startsAt` の前後 30 分で確認する
2. 仮説を 2〜3 個立て、それぞれをクエリで検証する
3. メトリクスで原因が閉じないときは、`startsAt` の前後 5 分のログとトレースを見る
   - ログ: `query_loki_logs` で `startsAt` 前後の ERROR / WARNING 行を引き、原因を示す行をそのまま引用する
   - トレース: `tempo_traceql-search` で遅い、または失敗したトレースを探し、`tempo_get-trace` でどのサービスのどのスパンが時間を使っているかを見る
4. `get_annotations` で直近のデプロイや過去の調査結果があれば参照する

調べ方の原則:

- エラーの原因は、エラーを返した側(下流のサービス)のログから引く。呼び出し側に残る HTTP ステータス(409 や 502)は結果であって原因ではない
- 失敗が一部に偏っていないか(種別・対象・ユーザーなど)を、ログの属性で集計して確かめる。偏りがあれば、その値を結論に書く
- 過去の調査結果(annotation)は参考に留める。今回のアラートの時間帯のデータで確認できたことだけを根拠にし、過去の障害を今回の原因に流用しない

回答は日本語で、次の見出し構成の Markdown にしてください。

## 症状

アラートの内容と、実際に観測した値(数字と時刻を含める)。

## 仮説と検証

| #   | 仮説 | 検証に使ったクエリ | 結果 | 判定 |
| --- | ---- | ------------------ | ---- | ---- |

## 結論

最も可能性の高い原因を 1〜3 行で。確信度(高 / 中 / 低)を添える。
根拠にしたログ行やスパン名があれば、そのまま引用する。

## 次の一手

運用者が取るべき行動を箇条書きで 3 つまで。
mcp.json Claude Code に渡す MCP サーバ定義
{
  "mcpServers": {
    "grafana": {
      "type": "http",
      "url": "http://localhost:8097/mcp"
    }
  }
}
local_agent.py ローカル LLM 用の最小エージェント
"""ローカル LLM 用の最小エージェント。受け口(investigate.py)から `claude -p` の代わりに呼ぶ。

    uv run --no-project --with mcp --with openai python local_agent.py \
        --model qwen3.6-35b-a3b --mcp-config mcp.json [--allowedTools ...]

- stdin にプロンプト、stdout に `claude --output-format json` と同じ形の JSON を返す
- MCP サーバ(mcp.json の先頭)に Streamable HTTP でつなぎ、ツール一覧を OpenAI 互換の function calling に変換する
- LLM は OpenAI 互換 API(既定 http://localhost:8099/v1、llama.cpp)。ツール呼び出しがなくなるまで回す
- コンテキストが小さいので、ツールの結果は LOCAL_MAX_TOOL_CHARS で切り詰め、
  prompt_tokens が LOCAL_COMPACT_AT を超えたら古いツール結果を短くする

Claude Code のシステムプロンプト(約 18k トークン)を持たないので、同じ 32k のモデルでも調査が最後まで回る。
受け口の CLAUDE_BIN に "uv run --no-project --with mcp --with openai python local_agent.py" を入れると差し替わる。
"""

from __future__ import annotations

import argparse
import asyncio
import json
import os
import sys
import time

from mcp import ClientSession
from openai import OpenAI

try:  # mcp SDK の版で名前が違う
    from mcp.client.streamable_http import streamable_http_client as streamablehttp_client
except ImportError:  # pragma: no cover
    from mcp.client.streamable_http import streamablehttp_client

LLM_BASE = os.environ.get("LOCAL_LLM_BASE", "http://localhost:8099/v1")
MAX_TURNS = int(os.environ.get("LOCAL_MAX_TURNS", "30"))
# 既定は llama.cpp 64k 向け(ツール 20 回 / 結果 6,000 文字 / 52k で圧縮)。32k なら 14 / 3000 / 24000 にする
MAX_TOOL_CALLS = int(os.environ.get("LOCAL_MAX_TOOL_CALLS", "20"))
MAX_TOOL_CHARS = int(os.environ.get("LOCAL_MAX_TOOL_CHARS", "6000"))
COMPACT_AT = int(os.environ.get("LOCAL_COMPACT_AT", "52000"))
KEEP_RECENT = int(os.environ.get("LOCAL_KEEP_RECENT", "4"))
TOOL_ALLOW = [t for t in os.environ.get("LOCAL_TOOLS", "").split(",") if t]

# 枠組みの説明だけを書く。調査の手順やクエリの書き方は Claude Code と共通の PROMPT.md に置く
SYSTEM = (
    "ツールは 1 回に 1 つずつ呼び、結果を読んでから次を決めてください。"
    f"ツール呼び出しは合計 {MAX_TOOL_CALLS} 回までです。各結果の末尾に残り回数を示します。"
    "残りが少なくなったら、手元の事実で報告を書いてください。"
    "tempo_get-trace の結果は「service | span | duration_ms | detail」の表で返ります。"
)


def log(msg: str) -> None:
    print(f"[local_agent] {msg}", file=sys.stderr, flush=True)


def to_openai_tool(t) -> dict:
    schema = (getattr(t, "input_schema", None) or getattr(t, "inputSchema", None)
              or {"type": "object", "properties": {}})
    return {"type": "function", "function": {
        "name": t.name,
        "description": (t.description or "")[:800],
        "parameters": schema,
    }}


def shape_result(name: str, text: str) -> str:
    """大きな JSON を小さなモデル向けに要約する。Tempo の検索結果とトレースが対象。他はそのまま返す"""
    if not text.lstrip().startswith("{"):
        return text
    try:
        data = json.loads(text)
    except json.JSONDecodeError:
        return text
    if name.endswith("traceql-search") and isinstance(data.get("traces"), list):
        lines = ["traceID | root service | root span | duration_ms"]
        for t in data["traces"][:15]:
            lines.append(f"{t.get('traceID')} | {t.get('rootServiceName')} | {t.get('rootTraceName')} | {t.get('durationMs')}")
        lines.append(f"({len(data['traces'])} traces。詳細は tempo_get-trace で 1 件ずつ)")
        return "\n".join(lines)
    if name.endswith("get-trace") and isinstance(data.get("trace"), dict):
        rows = []
        for svc in data["trace"].get("services", []):
            for scope in svc.get("scopes", []):
                for sp in scope.get("spans", []):
                    a = sp.get("attributes", {}) or {}
                    detail = a.get("db.operation.name") or a.get("http.route") or a.get("url.full") or ""
                    status = sp.get("status", {}).get("code", "")
                    err = " ERROR" if "ERROR" in status else ""
                    rows.append((int(sp.get("startTimeUnixNano", 0)), svc.get("serviceName"), sp.get("name"),
                                 float(sp.get("durationMs", 0)), detail, err,
                                 a.get("http.response.status_code", "")))
        rows.sort()
        lines = [f"trace {data['trace'].get('traceId')} : service | span | duration_ms | detail"]
        for _, s, n, d, detail, err, code in rows[:40]:
            lines.append(f"{s} | {n} | {d:.1f} | {detail} {code}{err}".rstrip())
        if len(rows) > 40:
            lines.append(f"(全 {len(rows)} スパンのうち先頭 40)")
        return "\n".join(lines)
    return text


def compact(messages: list[dict]) -> int:
    """古いツール結果を短くしてコンテキストを空ける。最近の KEEP_RECENT 件は残す"""
    idx = [i for i, m in enumerate(messages) if m["role"] == "tool" and not m.get("_compacted")]
    n = 0
    for i in idx[:-KEEP_RECENT] if len(idx) > KEEP_RECENT else []:
        m = messages[i]
        if len(m["content"]) > 300:
            m["content"] = m["content"][:300] + "…(古い結果なので省略)"
            m["_compacted"] = True
            n += 1
    return n


def dedupe_paragraphs(text: str) -> str:
    """反復ループで同じ段落が並んだ出力を、初出だけ残して畳む(判定のため元の長さは注記する)"""
    paras = text.split("\n\n")
    out, seen = [], set()
    for p in paras:
        k = p.strip()
        if k and k in seen:
            continue
        seen.add(k)
        out.append(p)
    if len(out) < len(paras):
        out.append(f"\n(同じ段落の繰り返し {len(paras) - len(out)} 個を省略)")
    return "\n\n".join(out)


def strip_private(messages: list[dict]) -> list[dict]:
    return [{k: v for k, v in m.items() if not k.startswith("_")} for m in messages]


async def run(model: str, mcp_url: str, prompt: str) -> dict:
    t0 = time.monotonic()
    usage = {"input_tokens": 0, "output_tokens": 0}
    tool_log: list[str] = []
    async with streamablehttp_client(mcp_url) as streams:
        read, write = streams[0], streams[1]
        async with ClientSession(read, write) as session:
            await session.initialize()
            tools = (await session.list_tools()).tools
            if TOOL_ALLOW:
                tools = [t for t in tools if t.name in TOOL_ALLOW]
            oai_tools = [to_openai_tool(t) for t in tools]
            log(f"tools={len(oai_tools)} model={model} llm={LLM_BASE}")

            client = OpenAI(base_url=LLM_BASE, api_key=os.environ.get("LOCAL_LLM_KEY", "none"), timeout=600)
            messages: list[dict] = [{"role": "system", "content": SYSTEM}, {"role": "user", "content": prompt}]
            turns = 0
            result = ""
            seen: dict[str, int] = {}
            calls_left = MAX_TOOL_CALLS
            for turn in range(MAX_TURNS):
                last = turn == MAX_TURNS - 1 or calls_left <= 0
                if calls_left <= 0 and messages[-1]["role"] != "user":
                    messages.append({"role": "user", "content": "ツール呼び出しの上限に達しました。ここまでに得た事実だけで、"
                                    "指示された見出し構成(症状 / 仮説と検証 / 結論 / 次の一手)の報告を書いてください。"})
                resp = client.chat.completions.create(
                    model=model, messages=strip_private(messages),
                    tools=oai_tools, tool_choice="none" if last else "auto",
                    max_tokens=int(os.environ.get("LOCAL_MAX_TOKENS", "4000")),
                    # Qwen3 系の推奨値。温度を下げすぎると同じ文を繰り返し始める
                    temperature=float(os.environ.get("LOCAL_TEMPERATURE", "0.6")),
                    top_p=0.95,
                    presence_penalty=float(os.environ.get("LOCAL_PRESENCE_PENALTY", "1.0")),
                )
                turns += 1
                if resp.usage:
                    usage["input_tokens"] += resp.usage.prompt_tokens or 0
                    usage["output_tokens"] += resp.usage.completion_tokens or 0
                    log(f"turn {turns}: prompt_tokens={resp.usage.prompt_tokens} completion={resp.usage.completion_tokens}")
                    if (resp.usage.prompt_tokens or 0) > COMPACT_AT:
                        log(f"compacted {compact(messages)} tool results")
                msg = resp.choices[0].message
                assistant: dict = {"role": "assistant", "content": msg.content or ""}
                reasoning = getattr(msg, "reasoning_content", None)
                if reasoning:
                    assistant["reasoning_content"] = reasoning
                if msg.tool_calls:
                    assistant["tool_calls"] = [{
                        "id": tc.id, "type": "function",
                        "function": {"name": tc.function.name, "arguments": tc.function.arguments or "{}"},
                    } for tc in msg.tool_calls]
                messages.append(assistant)
                if not msg.tool_calls:
                    result = msg.content or ""
                    break
                for tc in msg.tool_calls:
                    name = tc.function.name
                    try:
                        args = json.loads(tc.function.arguments or "{}")
                    except json.JSONDecodeError:
                        args = {}
                    key = name + json.dumps(args, sort_keys=True, ensure_ascii=False)
                    tool_log.append(name)
                    if seen.get(key, 0) >= 1:
                        # 同じ呼び出しの繰り返し。結果を返さず、手元の情報で書くよう促す
                        seen[key] += 1
                        log(f"  {name} repeated x{seen[key]} -> skipped")
                        messages.append({"role": "tool", "tool_call_id": tc.id,
                                         "content": "この呼び出しは既に同じ引数で実行済みです。前の結果を使い、"
                                                    "新しい事実が要らなければ報告を書いてください。"})
                        continue
                    seen[key] = 1
                    try:
                        r = await session.call_tool(name, args)
                        text = "\n".join(getattr(c, "text", "") for c in r.content) or "(結果なし)"
                        if getattr(r, "is_error", None) or getattr(r, "isError", None):
                            text = "error: " + text
                        else:
                            text = shape_result(name, text)
                    except Exception as e:  # noqa: BLE001
                        text = f"error: {e}"
                    if len(text) > MAX_TOOL_CHARS:
                        text = text[:MAX_TOOL_CHARS] + f"…(全 {len(text)} 文字のうち先頭のみ)"
                    calls_left -= 1
                    text += f"\n(残りのツール呼び出し: {max(calls_left, 0)} 回)"
                    log(f"  {name}({json.dumps(args, ensure_ascii=False)[:120]}) -> {len(text)} chars")
                    messages.append({"role": "tool", "tool_call_id": tc.id, "content": text})
            else:
                result = messages[-1].get("content") or "(上限ターン数に達したので中断)"
    if "</think>" in result:
        # llama.cpp が思考部分を分離し損ねたとき、本文だけ残す
        result = result.split("</think>", 1)[1].lstrip()
    result = dedupe_paragraphs(result)
    return {
        "type": "result", "subtype": "success", "is_error": False,
        "duration_ms": int((time.monotonic() - t0) * 1000),
        "num_turns": turns, "result": result, "usage": usage,
        "tool_calls": tool_log, "model": model,
    }


def main() -> int:
    if hasattr(sys.stdout, "reconfigure"):
        sys.stdout.reconfigure(encoding="utf-8")
        sys.stdin.reconfigure(encoding="utf-8")
    parser = argparse.ArgumentParser()
    parser.add_argument("--model", default=os.environ.get("CLAUDE_MODEL", "qwen3.6-35b-a3b"))
    parser.add_argument("--mcp-config", default="mcp.json")
    args, _unknown = parser.parse_known_args()   # claude -p 互換の他の引数は無視する
    with open(args.mcp_config, encoding="utf-8") as fh:
        servers = json.load(fh)["mcpServers"]
    mcp_url = next(iter(servers.values()))["url"]
    prompt = sys.stdin.read()
    out = asyncio.run(run(args.model, mcp_url, prompt))
    print(json.dumps(out, ensure_ascii=False))
    return 0


if __name__ == "__main__":
    sys.exit(main())

負荷と故障注入

loadgen.py 平常時のジョブ投入(5 rps)
"""平常時の負荷。一定のレートで計算ジョブを投入し続ける(標準ライブラリのみ)。

使い方: python loadgen.py [--rps 5] [--seconds 0(無期限)] [--url http://localhost:8100/jobs]

応答を待たずに次を送る(open loop)ので、サーバが遅くなっても送信レートは落ちない。
同時実行が --max-inflight を超えた分は送らずに数える。
"""

from __future__ import annotations

import argparse
import json
import random
import threading
import time
import urllib.error
import urllib.request

USERS = ["sato", "suzuki", "takahashi", "tanaka"]
JOB_TYPES = ["mip", "lp", "simulation", "ml"]


def run(url: str, rps: float, seconds: float, max_inflight: int = 500, timeout: float = 10.0,
        label: str = "loadgen") -> dict:
    stats = {"sent": 0, "ok": 0, "failed": 0, "error": 0, "skipped": 0, "latency": []}
    inflight = threading.Semaphore(max_inflight)
    lock = threading.Lock()

    def one() -> None:
        body = json.dumps({
            "user": random.choice(USERS), "job": random.choice(JOB_TYPES), "size": random.randint(1, 3),
        }).encode("utf-8")
        req = urllib.request.Request(url, data=body, headers={"Content-Type": "application/json"})
        t0 = time.perf_counter()
        try:
            with urllib.request.urlopen(req, timeout=timeout) as resp:
                resp.read()
                key = "ok" if resp.status == 200 else "failed"
        except urllib.error.HTTPError:
            key = "failed"
        except (urllib.error.URLError, TimeoutError, OSError):
            key = "error"
        finally:
            inflight.release()
        with lock:
            stats[key] += 1
            stats["latency"].append(time.perf_counter() - t0)

    deadline = time.time() + seconds if seconds > 0 else float("inf")
    interval = 1.0 / rps
    next_at = time.perf_counter()
    last_report = time.time()
    try:
        while time.time() < deadline:
            now = time.perf_counter()
            if now < next_at:
                time.sleep(next_at - now)
            next_at += interval
            if inflight.acquire(blocking=False):
                stats["sent"] += 1
                threading.Thread(target=one, daemon=True).start()
            else:
                stats["skipped"] += 1
            if time.time() - last_report >= 10:
                last_report = time.time()
                with lock:
                    lat = sorted(stats["latency"][-200:]) or [0.0]
                    p95 = lat[int(len(lat) * 0.95) - 1] if len(lat) > 1 else lat[0]
                    print(f"[{label}] sent={stats['sent']} ok={stats['ok']} failed={stats['failed']} "
                          f"error={stats['error']} skipped={stats['skipped']} p95={p95:.3f}s", flush=True)
    except KeyboardInterrupt:
        pass
    time.sleep(min(timeout, 2.0))
    return stats


if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--url", default="http://localhost:8100/jobs")
    parser.add_argument("--rps", type=float, default=5)
    parser.add_argument("--seconds", type=float, default=0, help="0 で無期限(Ctrl+C で止める)")
    parser.add_argument("--max-inflight", type=int, default=500)
    args = parser.parse_args()
    s = run(args.url, args.rps, args.seconds, args.max_inflight)
    print(f"done sent={s['sent']} ok={s['ok']} failed={s['failed']} error={s['error']} skipped={s['skipped']}")
faults/load_spike.py 実験 1: 投入の急増
"""実験 1(レベル 1): リクエストの急増。平常の 5 rps に加えて高レートでジョブを投入する。

使い方: python faults/load_spike.py [--rps 80 --seconds 180 --max-inflight 60]
slot-manager の DB プール(4 本 × 約 80 ms/ジョブ = 50 rps)を超えるので待ちが伸び、
JobApiLatencyHigh(p95 > 0.5 秒が 1 分)が発火する。原因はメトリクス(リクエスト率)だけで分かる。
同時実行を --max-inflight で頭打ちにして、待ち行列が無限に伸びてタイムアウト(failed)に化けるのを防ぐ
(頭打ちなしで 80 rps を送ると job-api の httpx 接続プールが枯渇して 502 が混じり、失敗率のアラートも鳴る)。
"""

from __future__ import annotations

import argparse
import sys
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from loadgen import run  # noqa: E402

parser = argparse.ArgumentParser()
parser.add_argument("--url", default="http://localhost:8100/jobs")
parser.add_argument("--rps", type=float, default=80)
parser.add_argument("--seconds", type=float, default=180)
parser.add_argument("--max-inflight", type=int, default=60)
args = parser.parse_args()

s = run(args.url, args.rps, args.seconds, max_inflight=args.max_inflight, label="spike")
print(f"done sent={s['sent']} ok={s['ok']} failed={s['failed']} error={s['error']} skipped={s['skipped']}")
faults/slot_out.py 実験 2: 特定種別の枠切れ(2b は --slots 45)
"""実験 2(レベル 2): 計算枠の枯渇。1 つのジョブ種別の枠を 0 にして、その種別の投入だけを失敗させる。

使い方: python faults/slot_out.py [--job mip --seconds 180]
失敗率が 25% 前後になり JobErrorRateHigh(失敗率 > 20% が 1 分)が発火する。
どの種別かはメトリクスに出ず、slot-manager のログ「failed to reserve: no free slot: mip」にしか無い。

--slots 45 のように少量を残すと、先に「slots running low」の警告ログが出て
Loki ruler の SlotsRunningLow(実験 2b)が鳴る。
"""

from __future__ import annotations

import argparse
import json
import time
import urllib.request

parser = argparse.ArgumentParser()
parser.add_argument("--url", default="http://localhost:8101/admin/fault")
parser.add_argument("--job", default="mip")
parser.add_argument("--slots", type=int, default=0)
parser.add_argument("--seconds", type=float, default=180)
args = parser.parse_args()


def post(body: dict) -> dict:
    req = urllib.request.Request(args.url, data=json.dumps(body).encode("utf-8"),
                                 headers={"Content-Type": "application/json"})
    with urllib.request.urlopen(req, timeout=10) as resp:
        return json.loads(resp.read())


print(f"slots[{args.job}] = {args.slots} for {args.seconds:.0f}s", flush=True)
post({"slots": {args.job: args.slots}})
try:
    time.sleep(args.seconds)
finally:
    r = post({"slots": {args.job: 10**9}})
    print(f"restored: slots[{args.job}] = {r['slots'][args.job]}")
faults/slow_db.py 実験 3: DB の遅延
"""実験 3(レベル 3): 下流の DB の遅延。slot-manager の DB 問い合わせに待ち時間を足す。

使い方: python faults/slow_db.py [--latency-ms 300 --seconds 180]
投入レートは変えないのに job-api の p95 が伸び、実験 1 と同じ JobApiLatencyHigh が発火する。
トレースを見ると slot-manager の db.query スパンが時間を占めている。
"""

from __future__ import annotations

import argparse
import json
import time
import urllib.request

parser = argparse.ArgumentParser()
parser.add_argument("--url", default="http://localhost:8101/admin/fault")
parser.add_argument("--latency-ms", type=float, default=300)
parser.add_argument("--seconds", type=float, default=180)
args = parser.parse_args()


def post(body: dict) -> dict:
    req = urllib.request.Request(args.url, data=json.dumps(body).encode("utf-8"),
                                 headers={"Content-Type": "application/json"})
    with urllib.request.urlopen(req, timeout=10) as resp:
        return json.loads(resp.read())


print(f"latency_ms = {args.latency_ms:.0f} for {args.seconds:.0f}s", flush=True)
post({"latency_ms": args.latency_ms})
try:
    time.sleep(args.seconds)
finally:
    r = post({"latency_ms": 0})
    print(f"restored: latency_ms = {r['fault']['latency_ms']}")
faults/retry_storm.py 実験 3b: 二重投入のバグ
"""実験 3b(任意): 二重投入のバグ。job-api が 1 ジョブあたり slot-manager を N 回呼ぶ。

使い方: python faults/retry_storm.py [--calls 8 --seconds 180]
5 回だと 1 ジョブ 0.43 秒で p95 がしきい値 0.5 秒に届かないことがある。8 回で約 0.7 秒。
投入レートも DB も変わらないのに p95 が伸びる。原因はトレースの構造(子スパンが N 本)にしか現れない。
"""

from __future__ import annotations

import argparse
import json
import time
import urllib.request

parser = argparse.ArgumentParser()
parser.add_argument("--url", default="http://localhost:8100/admin/fault")
parser.add_argument("--calls", type=int, default=8)
parser.add_argument("--seconds", type=float, default=180)
args = parser.parse_args()


def post(body: dict) -> dict:
    req = urllib.request.Request(args.url, data=json.dumps(body).encode("utf-8"),
                                 headers={"Content-Type": "application/json"})
    with urllib.request.urlopen(req, timeout=10) as resp:
        return json.loads(resp.read())


print(f"duplicate_calls = {args.calls} for {args.seconds:.0f}s", flush=True)
post({"duplicate_calls": args.calls})
try:
    time.sleep(args.seconds)
finally:
    r = post({"duplicate_calls": 1})
    print(f"restored: duplicate_calls = {r['fault']['duplicate_calls']}")

テスト

tests/test_investigate.py claude をモックにして受け口を通す
"""受け口の起動テスト。claude をモックし、偽の Grafana に annotation が届くことを確認する。

実行: uv run --no-project python tests/test_investigate.py
"""

import json
import os
import subprocess
import sys
import tempfile
import time
import urllib.request
from http.server import BaseHTTPRequestHandler, HTTPServer
from pathlib import Path
from threading import Thread

HERE = Path(__file__).resolve().parent
ROOT = HERE.parent
received: list[dict] = []


class FakeGrafana(BaseHTTPRequestHandler):
    def do_POST(self):  # noqa: N802
        n = int(self.headers.get("Content-Length", "0"))
        received.append({
            "path": self.path,
            "auth": self.headers.get("Authorization"),
            "body": json.loads(self.rfile.read(n)),
        })
        self.send_response(200)
        self.send_header("Content-Type", "application/json")
        self.end_headers()
        self.wfile.write(b'{"id":1,"message":"Annotation added"}')

    def log_message(self, *a):  # noqa: D102
        return


def main() -> int:
    grafana = HTTPServer(("127.0.0.1", 0), FakeGrafana)
    Thread(target=grafana.serve_forever, daemon=True).start()
    gport = grafana.server_address[1]

    with tempfile.TemporaryDirectory() as tmp:
        # claude の代わりにモックを起動させる(CLAUDE_BIN は空白区切りのコマンドでもよい)
        env = {
            **os.environ,
            "PORT": "9195",
            "CLAUDE_BIN": f"{sys.executable} {HERE / 'mock_claude.py'}",
            "REPORT_DIR": tmp,
            "GRAFANA_URL": f"http://127.0.0.1:{gport}",
            "GRAFANA_SA_TOKEN": "",
            "GRAFANA_BASIC_AUTH": "admin:admin",
        }
        proc = subprocess.Popen([sys.executable, str(ROOT / "investigate.py")], env=env)
        try:
            time.sleep(1.0)
            payload = (HERE / "sample_alert.json").read_bytes()
            req = urllib.request.Request(
                "http://127.0.0.1:9195/alert", data=payload,
                headers={"Content-Type": "application/json"}, method="POST",
            )
            with urllib.request.urlopen(req, timeout=5) as resp:
                assert resp.status == 200, resp.status
            # 同じ fingerprint を再送しても 2 回目は調査しない
            with urllib.request.urlopen(req, timeout=5):
                pass
            for _ in range(50):
                if received:
                    break
                time.sleep(0.1)
        finally:
            proc.terminate()
            proc.wait(timeout=5)

        reports = sorted(Path(tmp).glob("*_JobApiLatencyHigh.md"))
        assert len(reports) == 1, reports
        text = reports[0].read_text(encoding="utf-8")
        assert "## 結論" in text and "モックの結論" in text, text
        assert (Path(tmp) / "results.jsonl").exists()
        assert len(received) == 1, received
        assert received[0]["path"] == "/api/annotations"
        assert received[0]["auth"].startswith("Basic "), received[0]["auth"]
        assert "JobApiLatencyHigh" in received[0]["body"]["tags"]
        # プロンプトに Loki / Tempo の手順が入っていること
        prompt = (ROOT / "PROMPT.md").read_text(encoding="utf-8")
        assert "query_loki_logs" in prompt and "tempo_traceql-search" in prompt
    print("ok")
    return 0


if __name__ == "__main__":
    sys.exit(main())
tests/mock_claude.py
"""claude CLI のモック。stdin のプロンプトを受け取り、--output-format json と同じ形の結果を返す。"""

import json
import sys

sys.stdin.reconfigure(encoding="utf-8")
sys.stdout.reconfigure(encoding="utf-8")  # 本物の claude と同じく UTF-8 で返す
prompt = sys.stdin.read()
print(json.dumps({
    "type": "result",
    "subtype": "success",
    "is_error": False,
    "duration_ms": 1234,
    "num_turns": 5,
    "result": "## 症状\n\n(モック)プロンプト長 " + str(len(prompt)) + " 文字\n\n## 結論\n\nモックの結論\n",
    "total_cost_usd": 0.01,
    "usage": {"input_tokens": 100, "output_tokens": 50},
}, ensure_ascii=False))
tests/sample_alert.json
{
  "version": "4",
  "groupKey": "{}:{alertname=\"JobApiLatencyHigh\"}",
  "status": "firing",
  "receiver": "ai-investigator",
  "groupLabels": { "alertname": "JobApiLatencyHigh" },
  "commonLabels": { "alertname": "JobApiLatencyHigh", "severity": "warning" },
  "externalURL": "http://alertmanager:9093",
  "alerts": [
    {
      "status": "firing",
      "labels": {
        "alertname": "JobApiLatencyHigh",
        "severity": "warning"
      },
      "annotations": {
        "summary": "ジョブ投入 API の p95 レイテンシが 1 分以上 0.5 秒を超えている",
        "description": "直近 2 分の p95 = 1.2s"
      },
      "startsAt": "2026-09-12T05:23:45.000Z",
      "endsAt": "0001-01-01T00:00:00Z",
      "generatorURL": "http://prometheus:9090/graph?g0.expr=...",
      "fingerprint": "a1b2c3d4e5f60718"
    }
  ]
}