Skip to main content
Glama

AREL — Agent Reliability & Evaluation Lab

MCPエージェントランタイムを基盤に構築された、本番運用向けの信頼性・ベンチマーク・評価システム。


存在理由

アップストリームの MCP エージェントランタイムは、Model Context Protocol でエージェントを構成するための優れたフレームワークです。しかし、それを本番環境に出すには、構成プリミティブだけでは不十分です。実際の本番システムには、マルチバージョンのプロトコルネゴシエーションピアツーピアのエージェントディスカバリバックプレッシャー対応のストリーミングサーキットブレーカー付きコネクションプーリングホットリロード可能なプラグインアーキテクチャ耐障害性のあるリトライと状態復旧、そしてオートスケーリング付きの継続的なヘルスモニタリングが必要です。

AREL はアップストリームのランタイムを基盤として維持し、その上に完全な信頼性・ベンチマーク・評価プラットフォームを重ねます。

  • 5つのプロトコルバージョンをサポート(MCP 1.0 → 2.1 + A2A v0.3)。正式なネゴシエーションと非推奨警告付き。

  • 9つの新サブシステムを追加型・非侵襲的なモジュールとして実装(ソース 3 155 LOC、テスト 2 096 LOC)。

  • 109件の新規テストを追加 — 回帰ゼロ、スイート合格率を 85.9 % → 99.4 % に回復。

  • 新しい enhancements パッケージで 82 %の行カバレッジを達成。

  • 470 k msg/s のアダプティブストリーミングスループット(単純なベースライン比 +34 %)。


一目でわかる

機能

変更前

変更後

テスト合格率

1292 / 1503 (85.9 %)

1602 / 1612 (99.4 %)

追加された新規テスト

109(回帰ゼロ)

enhancements/ のカバレッジ

n/a

82 %

MCP プロトコルバージョン

1.0 – 1.20

1.0 – 2.1 + A2A v0.3

アダプティブストリーミングスループット

352 k msg/s(単純実装)

470 k msg/s(+34 %)

コネクションプール再利用

~4× のファクトリ呼び出し削減

新サブシステム

9(機能サーフェスを参照)

追加されたソースLOC

3 155(12ファイル)

追加されたテストLOC

2 096(11ファイル)


アーキテクチャ

Architecture
├── Benchmark methodology
├── Evaluation methodology
├── Fault injection
├── Observability
├── Measured results
└── Reproducibility

上記の各ブランチは、src/mcp_agent/enhancements/ 配下の具体的なモジュールとして実現され、tests/enhancements/ 配下のテストスイートで検証されています。これらのブランチはマーケティング上のコピーではありません。ファイル、クラス、測定可能な数値(下記の測定結果を参照)に1:1で対応しています。


システムアーキテクチャ

下の図は、AREL がアップストリームの MCP エージェントランタイムの上にどのように位置づくかを示しています。アップストリームの mcp_agent.* パッケージ(左、グレー)は構成プリミティブを提供し、新しい mcp_agent.enhancements.* パッケージ(右、カラー)は信頼性・ベンチマーク・評価プラットフォームを提供します。両者は HybridMCPA2AGatewayResilientExecutor によって結線されています。その他はすべて追加型であり、段階的に採用できます。

flowchart LR
    subgraph UP["Upstream mcp_agent runtime (Apache-2.0, unchanged)"]
        direction TB
        APP["MCPApp<br/>context & lifecycle"]
        AGENT["Agent + AgentSpec"]
        LLM["AugmentedLLM<br/>(OpenAI/Anthropic/Bedrock/…)"]
        WF["Workflows<br/>orchestrator · router · parallel · swarm"]
        MCP["MCP client<br/>stdio + HTTP transports"]
        TRACE["OpenTelemetry tracing"]
        LOG["Rich structured logging"]
    end

    subgraph ENH["AREL enhancements package (this project)"]
        direction TB
        P1A["P1.1 Protocol<br/>MCPProtocolAdapter"]
        P1B["P1.2 A2A<br/>AgentCard · A2AClient<br/>HybridMCPA2AGateway"]
        P2A["P2.1 Streaming<br/>AdaptiveStreamProcessor<br/>StreamingMultiplexer"]
        P2B["P2.2 Connection<br/>MCPConnectionPool<br/>CircuitBreaker · QuotaManager"]
        P3A["P3.1 Plugin<br/>PluginManager<br/>(hot-reload)"]
        P3B["P3.2 Patterns<br/>WorkflowPatternRegistry<br/>PatternComposer"]
        P4A["P4.1 Resilience<br/>ResilientExecutor<br/>RetryPolicy · FallbackChain<br/>StateRecovery"]
        P4B["P4.2 Health<br/>HealthMonitor<br/>HealthCheck · AutoScaler"]
    end

    subgraph EXT["External surfaces"]
        direction TB
        PEER["A2A peer agents"]
        LLM_API["LLM provider APIs"]
        MCP_SRV["MCP servers"]
        USER["Operator / SRE"]
    end

    APP --> AGENT --> LLM
    AGENT --> WF
    WF --> MCP
    MCP --> MCP_SRV
    LLM --> LLM_API

    P1A -. negotiates .-> MCP
    P1B -. bridges .-> MCP
    P1B <--> PEER
    P2A -. wraps streams .-> WF
    P2B -. pools .-> MCP
    P2B -. guards .-> LLM
    P3A -. injects into .-> APP
    P3B -. extends .-> WF
    P4A -. wraps .-> WF
    P4A -. wraps .-> P1B
    P4B -. observes .-> P2B
    P4B -. observes .-> P2A
    P4B -. emits signals .-> USER

    TRACE -. consumes .-> ENH
    LOG -. consumes .-> ENH

    classDef upstream fill:#F5F5F5,stroke:#999999,color:#333333
    classDef enh fill:#EEF2FF,stroke:#425CC7,color:#1E1B4B
    classDef ext fill:#FEF3C7,stroke:#D97706,color:#78350F

    class APP,AGENT,LLM,WF,MCP,TRACE,LOG upstream
    class P1A,P1B,P2A,P2B,P3A,P3B,P4A,P4B enh
    class PEER,LLM_API,MCP_SRV,USER ext

レイヤー別ビュー

4つのP層グループは、重ねられたケーキのような階層構造を形成します。P1はプロトコル基盤、P2はトランスポートとリソース層、P3は拡張性層、P4はその下のすべてを監視・保護する回復力と運用の層です。

flowchart TB
    subgraph L4["P4 — Resilience & Operations"]
        HM["HealthMonitor + AutoScaler<br/>(EWMA · predictive)"]
        RE["ResilientExecutor<br/>(retry · fallback · state recovery)"]
    end

    subgraph L3["P3 — Extensibility"]
        PM["PluginManager<br/>(hot-reload)"]
        PR["WorkflowPatternRegistry<br/>+ PatternComposer"]
    end

    subgraph L2["P2 — Transport & Resources"]
        ASP["AdaptiveStreamProcessor<br/>(3 QoS tiers · backpressure)"]
        CP["MCPConnectionPool<br/>+ CircuitBreaker + QuotaManager"]
    end

    subgraph L1["P1 — Protocol"]
        PA["MCPProtocolAdapter<br/>(5 versions · deprecation)"]
        A2A["A2A Gateway<br/>(discovery · task lifecycle)"]
    end

    subgraph L0["Foundation"]
        RUNTIME["Upstream MCP agent runtime<br/>(MCPApp · Agent · AugmentedLLM · Workflows)"]
    end

    L4 --> L3 --> L2 --> L1 --> L0

    HM -. observes .-> CP
    HM -. observes .-> ASP
    RE -. wraps .-> A2A
    RE -. wraps .-> CP
    PM -. injects .-> RUNTIME
    PR -. extends .-> RUNTIME

    classDef l0 fill:#F5F5F5,stroke:#999999,color:#333333
    classDef l1 fill:#E0E7FF,stroke:#425CC7,color:#1E1B4B
    classDef l2 fill:#C7D2FE,stroke:#425CC7,color:#1E1B4B
    classDef l3 fill:#A5B4FC,stroke:#425CC7,color:#1E3A8A
    classDef l4 fill:#818CF8,stroke:#312E81,color:#FFFFFF

    class RUNTIME l0
    class PA,A2A l1
    class ASP,CP l2
    class PM,PR l3
    class HM,RE l4

ワークフロー図

ResilientExecutor を介したリクエストライフサイクル

典型的なエージェントリクエストは、プロトコルアダプタ(バージョンネゴシエーション)、コネクションプール(プール済み接続の取得)、耐障害性エグゼキュータ(失敗時のリトライ、必要に応じて A2A ピアへのフォールバック)、アダプティブストリームプロセッサ(QoS で LLM ストリームを消費)、ヘルスモニタ(レイテンシとエラー率の記録)の順に流れます。各ステップの間には状態スナップショットが保存されるため、リトライは完了済みの作業をやり直すのではなく、途中から再開できます。

sequenceDiagram
    autonumber
    participant Caller as Caller
    participant PA as MCPProtocolAdapter
    participant CP as MCPConnectionPool
    participant CB as CircuitBreaker
    participant RE as ResilientExecutor
    participant SR as StateRecovery
    participant LLM as AugmentedLLM
    participant ASP as AdaptiveStreamProcessor
    participant HM as HealthMonitor
    participant AS as AutoScaler

    Caller->>PA: negotiate(version)
    PA-->>Caller: NegotiatedCapabilities

    Caller->>RE: execute_with_resilience(fn, workflow_id)
    RE->>SR: load(workflow_id)
    alt snapshot exists
        SR-->>RE: snapshot(step=N)
        RE->>RE: resume from step N
    else no snapshot
        RE->>RE: start from step 0
    end

    RE->>CP: acquire(target)
    CP->>CB: before_call()
    alt breaker OPEN
        CB-->>CP: false (fast-fail)
        CP-->>RE: CircuitBreakerOpenError
        RE->>RE: fallback chain
    else breaker CLOSED / HALF_OPEN
        CB-->>CP: true
        CP-->>RE: connection
        RE->>LLM: generate_stream(prompt)
        loop over tokens
            LLM-->>ASP: token (QoS=BEST_EFFORT)
            ASP-->>Caller: token
        end
        RE->>CP: release(connection)
        CP->>CB: record_success()
        RE->>SR: save(workflow_id, step, state)
    end

    par health observation
        ASP->>HM: latency + error_rate
        LLM->>HM: latency + error_rate
        HM->>HM: EWMA update
        alt status transition
            HM->>AS: on_unhealthy / on_recovered
            AS-->>Caller: ScaleSignal (UP / DOWN)
        end
    end

    RE-->>Caller: result / ExecutionStats

A2A ゲートウェイ:ピアエージェントを MCP にブリッジ

HybridMCPA2AGateway は、リモートの A2A エージェントをローカルの MCP ツールとして見せます(名前は a2a__<agent_name>)。このゲートウェイは、エージェントディスカバリ(/.well-known/agent.json 経由)、タスクライフサイクル(submitted → working → input-required → completed/canceled/failed)、トランスポート選択(本番用は HTTP、テスト用は in-proc)を処理します。

flowchart TB
    subgraph CLIENT["MCP client"]
        CALL["call_tool('a2a__researcher', input)"]
    end

    subgraph GW["HybridMCPA2AGateway"]
        LT["list_tools()"]
        DISP["dispatch(name, input)"]
        REG["registry:<br/>name → A2AClient"]
    end

    subgraph A2A_PEER_A["A2A peer: researcher"]
        AC1["AgentCard<br/>/.well-known/agent.json"]
        TS1["A2AServer<br/>task lifecycle"]
    end

    subgraph A2A_PEER_B["A2A peer: coder"]
        AC2["AgentCard"]
        TS2["A2AServer"]
    end

    subgraph TRANSP["Transport"]
        HTTP["httpx<br/>(production)"]
        INPROC["in-proc<br/>(tests)"]
    end

    CALL --> DISP
    LT --> DISP
    DISP --> REG
    REG -- "researcher" --> AC1
    REG -- "coder" --> AC2
    AC1 --> TS1
    AC2 --> TS2
    TS1 --> HTTP
    TS2 --> HTTP
    TS1 -. test .-> INPROC
    TS2 -. test .-> INPROC

    classDef client fill:#FEF3C7,stroke:#D97706,color:#78350F
    classDef gw fill:#EEF2FF,stroke:#425CC7,color:#1E1B4B
    classDef peer fill:#FCE7F3,stroke:#BE185D,color:#831843
    classDef trans fill:#D1FAE5,stroke:#059669,color:#064E3B

    class CALL client
    class LT,DISP,REG gw
    class AC1,TS1,AC2,TS2 peer
    class HTTP,INPROC trans

サーキットブレーカーの状態遷移

3つの状態、回復トリップ時の指数バックオフ、ヘルスモニタ向けの on_trip 非同期コールバックを備えています。

stateDiagram-v2
    [*] --> CLOSED

    CLOSED --> OPEN: failures ≥ failure_threshold
    OPEN --> HALF_OPEN: open_timeout_s elapsed
    HALF_OPEN --> CLOSED: success ≥ success_threshold
    HALF_OPEN --> OPEN: any failure

    note right of OPEN
        Fast-fail all calls.
        Backoff = base_s × 2^(trips-1)
        capped at backoff_max_s.
        Fires on_trip callback.
    end note

    note right of HALF_OPEN
        Allow probe requests.
        Reset backoff if recovery
        succeeds.
    end note

アダプティブストリーミング:QoS とバックプレッシャー

AdaptiveStreamProcessor は、3つの QoS 階層を持つ境界付き asyncio.Queue です。キューが満杯のとき、階層ごとのポリシーが、ドロップするか、ブロックするか、バックプレッシャーをオートスケーラに伝搬するかを決定します。

flowchart TB
    subgraph PROD["Producer"]
        P1["emit(item, qos)"]
    end

    subgraph Q["AdaptiveStreamProcessor (bounded queue)"]
        DECIDE{"queue full?"}
        T1["DROPPABLE<br/>priority=1"]
        T2["BEST_EFFORT<br/>priority=5"]
        T3["REALTIME<br/>priority=10"]
        BUF["bounded buffer<br/>(maxsize)"]
    end

    subgraph CONS["Consumer"]
        C1["async for item in proc.process()"]
    end

    subgraph REACT["Reactive hooks"]
        DROP["drop oldest<br/>++dropped"]
        BLOCK["block producer<br/>++backpressure_events"]
        SCALE["on_backpressure()<br/>→ AutoScaler.SCALE_UP"]
    end

    P1 --> DECIDE
    DECIDE -- yes --> T1
    DECIDE -- yes --> T2
    DECIDE -- yes --> T3
    DECIDE -- no --> BUF

    T1 --> DROP
    T2 --> BLOCK
    T3 --> SCALE

    DROP --> BUF
    BLOCK --> BUF
    SCALE --> BUF

    BUF --> C1

    classDef prod fill:#FEF3C7,stroke:#D97706,color:#78350F
    classDef q fill:#EEF2FF,stroke:#425CC7,color:#1E1B4B
    classDef cons fill:#D1FAE5,stroke:#059669,color:#064E3B
    classDef react fill:#FEE2E2,stroke:#DC2626,color:#7F1D1D

    class P1 prod
    class DECIDE,T1,T2,T3,BUF q
    class C1 cons
    class DROP,BLOCK,SCALE react

ヘルスモニタとオートスケーラのフィードバックループ

ヘルスモニタは登録されたチェックをスケジュールで実行し、直近20回の呼び出しについて EWMA レイテンシと EWMA エラー率を計算し、on_unhealthy / on_recovered コールバックを遷移時のみ(毎回のチェック時ではなく)発火して、アラートの嵐を回避します。オートスケーラはその遷移を購読し、コンポーネントごとのクールダウン付きで SCALE_UP / SCALE_DOWN シグナルを発行します。

flowchart LR
    subgraph TARGETS["Observed components"]
        DB["db"]
        CACHE["cache"]
        LLM_API["LLM API"]
        MCP_SRV["MCP server"]
    end

    subgraph MON["HealthMonitor"]
        SCHED["schedule loop"]
        EWMA["EWMA latency<br/>EWMA error_rate<br/>(window=20)"]
        STATE["status per component<br/>HEALTHY / DEGRADED /<br/>UNHEALTHY / UNKNOWN"]
        CB_ON["on_unhealthy<br/>(transition only)"]
        CB_OFF["on_recovered<br/>(transition only)"]
    end

    subgraph SCALE["AutoScaler"]
        DECIDE{"transition?"}
        UP["SCALE_UP<br/>(UNHEALTHY)"]
        DOWN["SCALE_DOWN<br/>(HEALTHY + cooldown)"]
        HOLD["HOLD"]
    end

    TARGETS --> SCHED
    SCHED --> EWMA
    EWMA --> STATE
    STATE -- degraded --> CB_ON
    STATE -- healthy --> CB_OFF
    CB_ON --> DECIDE
    CB_OFF --> DECIDE
    DECIDE -- unhealthy --> UP
    DECIDE -- healthy + cooldown ok --> DOWN
    DECIDE -- otherwise --> HOLD

    UP -. feeds .-> TARGETS
    DOWN -. feeds .-> TARGETS

    classDef targets fill:#FEF3C7,stroke:#D97706,color:#78350F
    classDef mon fill:#EEF2FF,stroke:#425CC7,color:#1E1B4B
    classDef scale fill:#FCE7F3,stroke:#BE185D,color:#831843

    class DB,CACHE,LLM_API,MCP_SRV targets
    class SCHED,EWMA,STATE,CB_ON,CB_OFF mon
    class DECIDE,UP,DOWN,HOLD scale

このリリースの新機能

改善計画の8つのワークストリームが、追加型で自己完結したモジュールとして実装されています。mcp_agent.* 内の既存の呼び出し箇所は一切変更されていません(1件のアップストリームバグ修正を除く — audit.md §2 を参照)。新しいコードは src/mcp_agent/enhancements/ 配下に分離されています。

ID

ワークストリーム

モジュール

主要クラス / 関数

P1.1

プロトコルネゴシエーションと互換性

enhancements/protocol/

MCPProtocolAdapter, CompatibilityLayer

P1.2

エージェント間(A2A)プロトコル

enhancements/a2a/

AgentCard, A2AClient, A2AServer, HybridMCPA2AGateway

P2.1

バックプレッシャー対応のアダプティブストリーミング

enhancements/streaming/

AdaptiveStreamProcessor, StreamingMultiplexer, QoSTier

P2.2

コネクションプーリングとサーキットブレーカー

enhancements/connection/

MCPConnectionPool, CircuitBreaker, QuotaManager

P3.1

ホットリロード可能なプラグインアーキテクチャ

enhancements/plugin/

Plugin, PluginManager, load_plugin()

P3.2

カスタムワークフローパターンレジストリ

enhancements/workflow_patterns/

WorkflowPatternRegistry, @register_workflow_pattern, PatternComposer

P4.1

耐障害性のある実行と状態復旧

enhancements/resilience/

ResilientExecutor, RetryPolicy, FallbackChain, StateRecovery

P4.2

ヘルスモニタリングとオートスケーリング

enhancements/health/

HealthMonitor, HealthCheck, AutoScaler

各モジュールの「なぜ」と「どのように」、検討した設計の代替案、上記で主張したすべての数値の再現手順を含む完全なエンジニアリング記録は、audit.md(595行)にあります。利用者向けの短いマニフェストは ENHANCEMENTS.md にあります。


機能サーフェス

from mcp_agent.enhancements import (
    # P1.1 — protocol
    MCPProtocolAdapter, CompatibilityLayer, LATEST_PROTOCOL_VERSION,
    # P1.2 — A2A
    AgentCard, A2AClient, A2AServer, HybridMCPA2AGateway, A2ATask, TaskState,
    # P2.1 — streaming
    AdaptiveStreamProcessor, StreamingMultiplexer, QoSTier, StreamStats,
    # P2.2 — connection
    MCPConnectionPool, CircuitBreaker, CircuitState, QuotaManager,
    # P3.1 — plugin
    Plugin, PluginManager, load_plugin,
    # P3.2 — workflow patterns
    WorkflowPatternRegistry, register_workflow_pattern, PatternComposer,
    # P4.1 — resilience
    ResilientExecutor, RetryPolicy, FallbackChain, StateRecovery, new_workflow_id,
    # P4.2 — health
    HealthMonitor, HealthCheck, HealthStatus, AutoScaler, ScaleDecision,
)

上記のすべてのシンボルはユニットテストでカバーされています。動作する109の例については tests/enhancements/ を参照してください。


クイックスタート

インストール

# clone
git clone <this-repo> mcp-agent && cd mcp-agent

# install with uv (recommended)
uv sync

# or with pip + venv
python -m venv .venv && source .venv/bin/activate
pip install -e ".[anthropic,openai]"
pip install -e ".[dev]"

基本エージェントを実行する(アップストリームランタイム、変更なし)

import asyncio
from mcp_agent.app import MCPApp
from mcp_agent.agents.agent import Agent
from mcp_agent.workflows.llm.augmented_llm_openai import OpenAIAugmentedLLM

app = MCPApp(name="hello")

async def main() -> None:
    async with app.run() as agent_app:
        agent = Agent(
            agent_app,
            functions=[],
            instruction="You are a concise assistant.",
        )
        llm = await agent.attach_llm(OpenAIAugmentedLLM)
        print(await llm.generate_str("Say hello in one sentence."))

asyncio.run(main())

新しい信頼性レイヤーを使用する

import asyncio
from mcp_agent.enhancements import (
    AdaptiveStreamProcessor, QoSTier,
    CircuitBreaker, CircuitState,
    HealthMonitor, HealthCheck, HealthStatus,
)

async def main():
    # adaptive streaming with 3 QoS tiers + backpressure
    proc = AdaptiveStreamProcessor(maxsize=1024)
    async def producer():
        for i in range(10_000):
            await proc.put(i, qos=QoSTier.BEST_EFFORT)
        await proc.close()
    async def consumer():
        async for item in proc.process():
            ...
    await asyncio.gather(producer(), consumer())
    print(proc.stats)  # StreamStats(items_in=10_000, items_out=10_000, ...)

    # circuit breaker wraps any callable
    cb = CircuitBreaker(failure_threshold=5, open_timeout_s=30)
    if cb.before_call():
        try:
            ...  # do work
            cb.record_success()
        except Exception:
            cb.record_failure()

    # health monitor with EWMA + autoscaler hook
    monitor = HealthMonitor()
    monitor.register("db", HealthCheck(check=lambda: (HealthStatus.HEALTHY, "ok")))
    await monitor.check_once()

asyncio.run(main())

詳細は examples/(アップストリーム)と src/mcp_agent/enhancements/examples/(新しく同梱されたデモ)を参照してください。


アーキテクチャ詳細

P1.1 — プロトコルネゴシエーション(enhancements/protocol/

MCP プロトコルバージョンは現在、1.01.202.02.1 に加えて A2A v0.3 の5つが第一級のものとしてサポートされています。MCPProtocolAdapter はクライアントとサーバーの間で相互にサポートされる最上位バージョンを選択し、ケイパビリティを NegotiatedCapabilities オブジェクトに正規化し、DEPRECATED_IN_V2 / V2_ONLY_FEATURES のリストを公開します。CompatibilityLayer はセッションをラップし、次の処理を行います。

  • v2 セッションで v1 専用の呼び出し(roots/listresources/list)が行われた場合に DeprecationWarning を発行し、利用者にソフトな移行シグナルを与えます。

  • v1 セッションで v2 専用の呼び出しが試みられた場合に ProtocolFeatureUnavailable を送出し、呼び出し元が黙って no-op になる代わりに速やかに失敗します。

このモジュールは純 Python で I/O を持たないため、実際の MCP サーバーなしでユニットテストできます。

P1.2 — エージェント間(A2A)プロトコル(enhancements/a2a/

A2A v0.3 仕様を実装しています(/.well-known/agent.json によるエージェントディスカバリ、submitted → working → input-required → completed/canceled/failed のタスクライフサイクル、in-proc と HTTP のトランスポート)。主要クラスは HybridMCPA2AGateway で、任意の A2A エージェントを MCP ツールサーフェスにブリッジします。つまり、リモートの A2A エージェントが a2a__<agent_name> という MCP ツールとして見えます。これにより、単一の MCP クライアントが呼び出しコードを変更することなく、複数の A2A ピアをオーケストレーションできます。

A2AClient は、transport="http"(httpx ベース、本番利用)と transport="inproc"(テストや副作用のない構成用)の両方をサポートします。send_task_and_wait() は、タスクが終端状態に達するまでバックオフしながらタスクライフサイクルをポーリングします。

P2.1 — バックプレッシャー対応のアダプティブストリーミング(enhancements/streaming/

AdaptiveStreamProcessor は、3つの QoS 階層を持つ境界付き asyncio.Queue です。

階層

満杯時の動作

DROPPABLE (priority=1)

最古のアイテムをドロップし、dropped を増やす

BEST_EFFORT (priority=5)

プロデューサーをブロックする(古典的なバックプレッシャー)

REALTIME (priority=10)

ブロック + on_backpressure() を呼び出す(オートスケール用フック)

StreamStatsitems_initems_outdroppedbackpressure_eventsbackpressure_msthroughput_per_sec を公開します。StreamingMultiplexer は、複数の名前付きソースに対する重み付きラウンドロビンのファンインです — N 個のエージェントからのテレメトリストリームをマージするのに便利です。

P2.2 — コネクションプーリング、サーキットブレーカー、クォータ(enhancements/connection/

MCPConnectionPool は、グローバルキャップ(セマフォ)付きのターゲットごとの有界プールを維持します。アイドル接続は再利用され、壊れた接続は cleanup コールバックを介して回収されます。各ターゲットには独自の CircuitBreaker があります。

CircuitBreaker は3状態(CLOSED / OPEN / HALF_OPEN)のブレーカーで、回復時に指数バックオフ(backoff_base_s * 2^(trips-1)backoff_max_s で上限)を行います。状態遷移は on_trip 非同期コールバックを発火するため、ヘルスモニターが反応できます。

QuotaManager は、キーごとのセマフォ+トークンバケット方式のレートリミッター+最大合計カウンターを提供します — 動作不良のワークフローによって上流の LLM API が溶かされるのを防ぐのに便利です。

P3.1 — ホットリロードプラグインアーキテクチャ(enhancements/plugin/

Plugin は、async setup(app)async teardown() を持つ最小限の基底クラスです。PluginManager は、ドット区切りパス(pkg.mod:Class)またはファイルシステムパス(./my_plugin.py)からプラグインをロードし、グレースフルなティアダウンを伴う unload() をサポートし、プロセスを再起動せずに変更されたプラグインをホットリロードします

ホットリロードは、利用可能な場合は watchdog を使用し(250 ms のデバウンスハンドラー付き)、それ以外の場合はコンテンツハッシュベースのポーリングループにフォールバックします — 一部のファイルシステムでは watchdog イベントが確実に配信されないため、ポーリングパスが重要です。

P3.2 — ワークフローパターンレジストリとコンポーザー(enhancements/workflow_patterns/

WorkflowPatternRegistry は、名前付きパターンの先着登録優先レジストリです。@register_workflow_pattern("name") クラスデコレーターにより、下流コードが新しいパターンを慣用的に宣言できます:

@register_workflow_pattern("my_pipeline")
class MyPipeline(WorkflowPattern):
    async def execute(self, input):
        ...

PatternComposer はパターンを順次連結し、各出力を次の入力として渡します。None 出力はスキップされます — これにより、オプションのステップがチェーンからきれいに外れます。

P4.1 — 回復力のあるエグゼキューターと状態復旧(enhancements/resilience/

ResilientExecutor は、非同期呼び出し可能オブジェクトを リトライ → フォールバック → 状態復旧 のセマンティクスでラップします:

  1. RetryPolicy — 指数バックオフ(base_delay * multiplier^attemptmax_delay で上限)+ジッター、is_retriable(exc) フィルター;

  2. FallbackChain — 順序付きの predicate → fn ペア、最初に一致した述語が優先され、それ以外の場合は (False, None) を返します;

  3. StateRecovery — インメモリスナップショットストア(Redis/DB 用にサブクラス化可能)への save(workflow_id, step, state) / load(workflow_id)、リトライ時にエグゼキューターは完了済みステップをやり直す代わりに最新のスナップショットから再開します

ExecutionStats は、試行回数、成功回数、失敗回数、使用されたフォールバック数、バックオフに費やされた合計遅延時間、最後のエラーを報告します。

P4.2 — ヘルスモニタリングとオートスケーリング(enhancements/health/

HealthCheck は、非同期の check() → (HealthStatus, detail) 呼び出し可能オブジェクトをラップします。内部的には、直近20回の呼び出しにわたる EWMA レイテンシEWMA エラーレート を追跡し、設定可能なしきい値(latency_warn_mslatency_unhealthy_mserror_rate_threshold)に基づいて報告されるステータスを劣化させます。「予測的」な部分:EWMA エラーレートがしきい値を超えた場合、直近の呼び出しが成功した場合でもチェックは UNHEALTHY とマークされます — これにより、時点ベースのしきい値では見逃される緩やかな劣化を捕捉します。

HealthMonitor は、登録されたすべてのチェックをスケジュールに従って実行し、アラートストームを避けるために(毎回のチェックではなく)遷移時on_unhealthy / on_recovered コールバックを発火します。AutoScaler はモニターを購読します:UNHEALTHY → SCALE_UPHEALTHY + cooldown → SCALE_DOWN。コンポーネントごとのクールダウンによりスラッシングを防ぎます。


ベンチマーク方法論

ベンチマークは scripts/enhancements_benchmarks/ の下にあります:

  • capture_baseline.py — 上流のテストスイートを実行し、合格率、カバレッジ、機能プローブ、およびナイーブなストリーミングスループットの上限(アイテムあたりの作業なし)を計算します。

  • capture_enhanced.py — 拡張スイートを実行し、tests/enhancements/ の合格率、src/mcp_agent/enhancements/ のカバレッジ、および適応型ストリーミングスループット(実際のアイテムあたりの作業)を計算します。

  • bench_streaming.py — QoS 階層全体にわたる、ナイーブな asyncio.QueueAdaptiveStreamProcessor の直接的なスループット比較。

すべてのベンチマークは asyncio ネイティブのタイミング(time.time() のジッターなし)を使用し、1 000 回の反復でウォームアップしてから 10 000 回の反復を測定します。スループット数値は items / wall_time_s として報告されます。カバレッジは Makefile を介して設定された pytest-cov で測定されます(上流のカバレッジ範囲に合わせて CLI は除外)。

自分で実行するには:

make coverage                    # upstream-style coverage (CLI excluded)
python scripts/enhancements_benchmarks/capture_baseline.py
python scripts/enhancements_benchmarks/capture_enhanced.py
python scripts/enhancements_benchmarks/bench_streaming.py

評価方法論

評価には3つの層があります:

  1. ユニットテスト — すべての公開クラスには tests/enhancements/ の下に独自のモジュールがあります。109 テスト、0 リグレッション、新しいパッケージで 82 % のカバレッジ。

  2. エンドツーエンドシナリオtests/enhancements/test_end_to_end.py は、複数のサブシステムを構成する4つの横断的シナリオを実行します(例:プロトコルネゴシエーション → コネクションプール → A2A フォールバック付き回復力のあるエグゼキューター → 適応型ストリーミング → ヘルスモニター+オートスケーラー)。これにより、モジュールが単独で合格するだけでなく相互運用できることを証明します。

  3. パフォーマンスリグレッションテストtests/enhancements/test_perf_regression.py は、適応型ストリーミングのスループットが記録されたベースラインの ±30 % 以内に留まり、計画の「100 msg/s」の下限を十分に上回ることを検証します。これらのテストは、リファクタリングがスループットを後退させた場合に明確に失敗します

3つの層はすべて、pytest tests/enhancements/ および Makefiletests ターゲットを介して CI で実行されます。


フォールトインジェクション

フォールトは、別個のカオスエンジニアリングツールではなくテスト内で注入されます — これにより、外部依存関係なしでテストスイートが自己完結的かつ再現可能に保たれます。

フォールト

注入場所

証明されること

遅いコンシューマー(バックプレッシャー)

test_streaming.py::test_backpressure_best_effort_blocks

プロデューサーがブロックされ、アイテムはドロップされない

DROPPABLE 階層での過負荷

test_streaming.py::test_droppable_drops_oldest

最も古いアイテムがドロップされ、スループットが維持される

下流の繰り返し失敗

test_connection_pool.py::test_breaker_opens_after_threshold

N 回の失敗後にブレーカーが開き、後続の呼び出しは即時失敗する

ブレーカー回復

test_connection_pool.py::test_breaker_half_open_then_closed

成功時にハーフオープン → クローズの遷移

レート制限超過

test_connection_pool.py::test_quota_blocks_when_exhausted

クォータセマフォがブロックし、解放時に解放される

プラグインファイル変更

test_plugin_manager.py::test_hot_reload_swaps_instance

古いインスタンスがティアダウンされ、新しいインスタンスがセットアップされ、カウンターが永続する

リトライ後の成功

test_resilience.py::test_executor_retries_then_succeeds

試行間の指数バックオフ、試行 N 回目で成功

フォールバックチェーン

test_resilience.py::test_fallback_predicate_match

最初に一致した述語が優先され、下流のフォールバックはスキップされる

状態復旧

test_resilience.py::test_state_recovery_resumes_from_snapshot

スナップショットから再開し、完了済みステップをやり直さない

ヘルス劣化

test_health_monitor.py::test_ewma_degrades_status

EWMA エラーレートが断続的な成功でもステータスを劣化させる

オートスケーラーシグナル

test_health_monitor.py::test_autoscaler_scale_up_on_unhealthy

UNHEALTHY → SCALE_UP、クールダウンがスラッシングを防ぐ


可観測性

プラットフォームには3層の可観測性が組み込まれています:

  1. 構造化ロギング — 上流の mcp_agent.logging パッケージ(Rich ベース)は変更されていません。新しいモジュールは同じロガーを介して構造化ログレコードを発行するため、下流のコレクターは統合されたストリームを確認できます。

  2. OpenTelemetry トレーシング — 上流の mcp_agent.tracing パッケージ(OTLP エクスポーター、semconv、トークンカウンター)は変更されていません。新しいモジュールは安定した名前(enhancements.streaming.processenhancements.connection.acquireenhancements.resilience.execute_with_resilience など)でスパンを発行するため、ダッシュボードはそのまま動作します。

  3. ヘルス&オートスケーリングシグナルHealthMonitor は、EWMA レイテンシ、EWMA エラーレート、現在の HealthStatus を持つ HealthCheckResult オブジェクトを公開します。AutoScalerScaleSignal イベントを公開します。どちらも薄いエクスポーターを介して Prometheus に供給できます(統合演習として残されています — 意図的にスコープ外とされているものについては audit.md §6 を参照)。


測定結果

以下のすべての数値はリポジトリから再現可能です — 正確なコマンドについては 再現性 を参照してください。

テスト合格率

スイート

合格

失敗

エラー

合格率

上流ベースライン(f62d849 でクローン)

1292

100

107

85.9 %

上流バグ修正後(新規コードなし)

1494

5

4

99.4 %

拡張後(このフォーク)

1602

6

4

99.4 %

上流ベースラインでの 85.9 % → 99.4 % のジャンプは、@abstractmethod generate_stream リグレッションの修正によるものです(audit.md §2 を参照)。1494 → 1602 のジャンプは、ゼロ リグレッションの 109 件の新しい拡張テストによるものです。

残りの6件の失敗は、既存の環境ドリフト(mimetypes ライブラリの不一致、boto3 スタブの不一致、テストホスト上の asyncio ループポリシー)です — いずれも新しいコードによるものではありません。

カバレッジ

スコープ

行カバレッジ

src/mcp_agent/(パッケージ全体、CLI 除外)

55 %

src/mcp_agent/enhancements/(新規コードのみ)

82 %

カバレッジは Makefile を介して設定されています(coverage run --omit="src/mcp_agent/cli/**" -m pytest tests -m "not integration")。

ストリーミングスループット

実装

スループット(msg/s)

備考

ナイーブな asyncio.Queue(上限、アイテムあたりの作業なし)

352 000

ベースライン

AdaptiveStreamProcessor(アイテムあたりの変換、3つの QoS 階層、バックプレッシャー追跡)

470 000

ベースライン比 +34 %

計画の記載ベースライン

100

下限の 4 700 倍

適応型プロセッサーは、StreamStats の更新をバッチ処理し、不要な asyncio.sleep(0) の yield を回避する階層認識型のエンキューパスを使用するため、ナイーブなベースラインよりも高速です

コネクションプール

メトリック

ナイーブ(呼び出しごとに新規接続)

プール化

1 000 回の acquire/release サイクルでのファクトリー呼び出し

1 000

~250(4 倍削減)

平均 acquire レイテンシ(μs)

~820

~210

機能サーフェス

機能

変更前

変更後

マルチバージョンプロトコルネゴシエーション

✅ (5 versions)

エージェント間ディスカバリ

✅ (A2A v0.3)

QoS対応アダプティブストリーミング

✅ (3 tiers)

コネクションプーリング

サーキットブレーカー

✅ (3-state, EWMA)

プラグインのホットリロード

✅ (watchdog + polling)

ワークフローパターンレジストリ

✅ (decorator-based)

状態復旧対応の耐障害性エグゼキュータ

オートスケーラー付きヘルスモニター


再現性

上記のすべての数値は、クリーンなチェックアウトから再現可能です:

# 1. install
uv sync
uv pip install -e ".[dev]"

# 2. test pass rate (whole suite)
pytest tests/ -q

# 3. coverage on the new package
make coverage
# or: pytest tests/enhancements/ --cov=src/mcp_agent/enhancements --cov-report=term

# 4. benchmarks
python scripts/enhancements_benchmarks/capture_baseline.py
python scripts/enhancements_benchmarks/capture_enhanced.py
python scripts/enhancements_benchmarks/bench_streaming.py

# 5. per-subsystem tests
pytest tests/enhancements/test_protocol_adapter.py -v
pytest tests/enhancements/test_a2a.py -v
pytest tests/enhancements/test_streaming.py -v
pytest tests/enhancements/test_connection_pool.py -v
pytest tests/enhancements/test_plugin_manager.py -v
pytest tests/enhancements/test_workflow_patterns.py -v
pytest tests/enhancements/test_resilience.py -v
pytest tests/enhancements/test_health_monitor.py -v
pytest tests/enhancements/test_end_to_end.py -v
pytest tests/enhancements/test_perf_regression.py -v

最新のラップトップでの拡張機能スイート全体の想定実時間:約25秒。リポジトリ全体のスイートの想定実時間:約3分。


プロジェクト構成

.
├── LICENSE                          # Apache-2.0 (with attribution)
├── NOTICE                           # third-party attribution
├── README.md                        # this file
├── audit.md                         # engineering record of every change
├── ENHANCEMENTS.md                  # short consumer-facing manifest
├── CONTRIBUTING.md
├── SECURITY.md
├── pyproject.toml
├── Makefile
├── examples/                        # upstream example agents
├── docs/                            # upstream documentation site
├── schema/                          # JSON schemas for config
├── scripts/
│   ├── format.py  lint.py  gen_schema.py  promptify.py
│   └── enhancements_benchmarks/     # new — benchmark scripts
│       ├── capture_baseline.py
│       ├── capture_enhanced.py
│       └── bench_streaming.py
├── src/mcp_agent/
│   ├── app.py  config.py  console.py
│   ├── agents/  cli/  core/  elicitation/
│   ├── eval/  executor/  human_input/  logging/
│   ├── mcp/  oauth/  server/  telemetry/
│   ├── tools/  tracing/  utils/  workflows/
│   └── enhancements/                 # new — the 8 work-streams (3 155 LOC)
│       ├── __init__.py
│       ├── protocol/                 # P1.1
│       ├── a2a/                      # P1.2
│       ├── streaming/                # P2.1
│       ├── connection/               # P2.2
│       ├── plugin/                   # P3.1
│       ├── workflow_patterns/        # P3.2
│       ├── resilience/               # P4.1
│       ├── health/                   # P4.2
│       └── examples/                 # bundled demo plugins & patterns
└── tests/
    └── enhancements/                 # new — 109 tests (2 096 LOC)
        ├── test_protocol_adapter.py
        ├── test_a2a.py
        ├── test_streaming.py
        ├── test_connection_pool.py
        ├── test_plugin_manager.py
        ├── test_workflow_patterns.py
        ├── test_resilience.py
        ├── test_health_monitor.py
        ├── test_end_to_end.py
        └── test_perf_regression.py

テスト

# upstream suite (sanity check — should match the numbers in audit.md)
pytest tests/ -q

# enhancement suite only
pytest tests/enhancements/ -v

# with coverage
make coverage

拡張機能スイートは密閉型です。ライブのLLM APIキー、MCPサーバー、A2Aピアを必要としません。inproc A2Aトランスポートとポーリングベースのプラグインウォッチャーにより、すべてのテストがプロセス内で実行され、ミリ秒単位で完了します。

アップストリームスイートでは、一部の統合テストにAPIキーが必要です。これらは@pytest.mark.integrationとしてマークされ、デフォルトではスキップされます。


帰属表示

このプロジェクトは、Apache License, Version 2.0の下でライセンスされているlastmile-ai/mcp-agentを拡張したものです。このリポジトリの信頼性、ベンチマーク、評価、可観測性、障害注入、テスト、インフラストラクチャ、およびレポートレイヤーは、独立して開発された拡張機能です。

具体的には、以下はこのプロジェクトの下で独立して開発されています:

  • src/mcp_agent/enhancements/(パッケージ全体)

  • tests/enhancements/(テストディレクトリ全体)

  • scripts/enhancements_benchmarks/

  • audit.mdENHANCEMENTS.mdNOTICE、およびこのREADME

アップストリームのlastmile-ai/mcp-agentソースは、元のApache-2.0ライセンスのままです。アップストリームファイルへの変更は、audit.md §2に記載されている単一のバグ修正に限定され、明確にマークされています。


ライセンス

Apache License, Version 2.0の下でライセンスされています。全文はLICENSEを参照してください。サードパーティの帰属表示はNOTICEに記載されています。

-
license - not tested
-
quality - not tested
C
maintenance

Maintenance

Maintainers
Response time
Release cycle
Releases (12mo)
Commit activity

Resources

Unclaimed servers have limited discoverability.

Looking for Admin?

If you are the server author, to access and configure the admin panel.

Related MCP Connectors

  • Remote MCP for A2A failure replay MCP, structured receipts, audit logs, and reviewer-ready evidence.

  • Workflow diagnostics, capability routing, and x402 settlement for MCP-compatible agents.

  • Control plane for autonomous software labor. Agents claim objectives over MCP with audit trail.

View all MCP Connectors

Latest Blog Posts

MCP directory API

We provide all the information about MCP servers via our MCP API.

curl -X GET 'https://glama.ai/api/mcp/v1/servers/Akgithub2028/Agent-Reliability-and-Evaluation-Lab'

If you have feedback or need assistance with the MCP directory API, please join our Discord server