Skip to main content
Glama

AREL — Agent Reliability & Evaluation Lab

Продакшн-ориентированная система надёжности, бенчмаркинга и оценки, построенная вокруг MCP agent runtime.


Зачем это нужно

Исходный MCP agent runtime — отличный фреймворк для компоновки агентов с помощью Model Context Protocol, но для вывода в продакшн одних примитивов компоновки недостаточно. Реальным продакшн-системам нужны согласование нескольких версий протокола, одноранговое обнаружение агентов, потоковая передача с учётом обратного давления, пул соединений с автоматическими выключателями, архитектура плагинов с горячей перезагрузкой, отказоустойчивые повторные попытки и восстановление состояния, а также непрерывный мониторинг состояния с автоскалированием.

AREL сохраняет исходный runtime как основу и добавляет поверх него полноценную платформу надёжности, бенчмаркинга и оценки:

  • 5 версий протокола поддерживается (MCP 1.0 → 2.1 + A2A v0.3) с формальным согласованием и предупреждениями об устаревании

  • 9 новых подсистем, реализованных как дополнительные неинвазивные модули (3 155 LOC исходного кода, 2 096 LOC тестов)

  • Добавлено 109 новых тестов — 0 регрессий, доля прохождения набора восстановлена с 85.9 % → 99.4 %

  • 82 % покрытия строк в новом пакете enhancements

  • 470 k msg/s пропускной способности адаптивной потоковой передачи (+34 % к наивному базовому варианту)


Краткий обзор

Возможность

До

После

Доля прохождения тестов

1292 / 1503 (85.9 %)

1602 / 1612 (99.4 %)

Добавлено новых тестов

109 (0 регрессий)

Покрытие 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 agent runtime. Исходный пакет mcp_agent.* (слева, серым) предоставляет примитивы компоновки; новый пакет mcp_agent.enhancements.* (справа, цветной) предоставляет платформу надёжности, бенчмаркинга и оценки. Они соединены через HybridMCPA2AGateway и ResilientExecutor — всё остальное является дополнительным и может внедряться инкрементально.

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

Послойное представление

Четыре группы 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

Диаграммы рабочих процессов

Жизненный цикл запроса через отказоустойчивый исполнитель

Типичный запрос агента проходит через адаптер протокола (согласование версии), пул соединений (получение соединения из пула), отказоустойчивый исполнитель (повтор при сбое, при необходимости переключение на A2A-агента), адаптивный процессор потока (потребление LLM-потока с QoS) и монитор состояния (запись задержки и частоты ошибок). Между шагами сохраняются снимки состояния, чтобы повторная попытка могла продолжиться с середины вместо повторного выполнения завершённой работы.

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 для продакшена, внутрипроцессный для тестов).

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

Конечный автомат автоматического выключателя

Три состояния, экспоненциальная задержка при срабатывании восстановления, асинхронный колбэк 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 — это ограниченная asyncio.Queue с тремя уровнями QoS:

Уровень

Поведение при заполнении

DROPPABLE (priority=1)

отбросить самый старый элемент; увеличить счётчик dropped

BEST_EFFORT (priority=5)

заблокировать производителя (классическое обратное давление)

REALTIME (priority=10)

заблокировать + вызвать on_backpressure() (хук для автоскалирования)

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

Контур обратной связи монитора состояния и автоскалера

Монитор состояния выполняет зарегистрированные проверки по расписанию, вычисляет EWMA-задержку и EWMA-частоту ошибок за последние 20 вызовов и запускает колбэки 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

Что нового в этом выпуске

Восемь рабочих направлений из плана улучшений реализованы как дополнительные самодостаточные модули. Ни одна существующая точка вызова в mcp_agent.* не изменялась (кроме одного исправления бага в исходном коде — см. audit.md §2). Новый код изолирован в src/mcp_agent/enhancements/:

ID

Рабочее направление

Модуль

Основной класс / функция

P1.1

Согласование протоколов и совместимость

enhancements/protocol/

MCPProtocolAdapter, CompatibilityLayer

P1.2

Протокол Agent-to-Agent (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,
)

Каждый указанный выше символ покрыт модульными тестами; см. tests/enhancements/ — 109 рабочих примеров.


Быстрый старт

Установка

# 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]"

Запуск базового агента (исходный runtime, без изменений)

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.0, 1.20, 2.0, 2.1, плюс A2A v0.3. MCPProtocolAdapter выбирает максимальную взаимно поддерживаемую версию между клиентом и сервером, нормализует возможности в объект NegotiatedCapabilities и предоставляет списки DEPRECATED_IN_V2 / V2_ONLY_FEATURES. CompatibilityLayer оборачивает сессию и:

  • испускает DeprecationWarning, когда вызов, доступный только в v1 (roots/list, resources/list), выполняется в v2-сессии, — потребители получают мягкий сигнал к миграции;

  • выбрасывает ProtocolFeatureUnavailable, когда вызов, доступный только в v2, предпринимается в v1-сессии, — вызывающий код быстро завершается ошибкой вместо тихого no-op.

Этот модуль написан на чистом Python и не выполняет ввода-вывода — его можно тестировать без живого MCP-сервера.

P1.2 — Протокол Agent-to-Agent (A2A) (enhancements/a2a/)

Реализует спецификацию A2A v0.3 (обнаружение агентов через /.well-known/agent.json, жизненный цикл задач submitted → working → input-required → completed/canceled/failed, внутрипроцессный и HTTP-транспорты). Ключевой класс — HybridMCPA2AGateway, который подключает любого A2A-агента к поверхности MCP-инструментов через мост — удалённые A2A-агенты выглядят как MCP-инструменты с именем a2a__<agent_name>. Это позволяет одному MCP-клиенту оркестрировать целый парк A2A-агентов, не меняя свой вызывающий код.

A2AClient поддерживает как transport="http" (на базе httpx, для продакшена), так и transport="inproc" (для тестов и композиции без побочных эффектов). send_task_and_wait() опрашивает жизненный цикл задачи с задержками, пока не дойдёт до терминального состояния.

P2.1 — Адаптивная потоковая передача с обратным давлением (enhancements/streaming/)

AdaptiveStreamProcessor — это ограниченная asyncio.Queue с тремя уровнями QoS:

Уровень

Поведение при заполнении

DROPPABLE (priority=1)

отбросить самый старый элемент; увеличить счётчик dropped

BEST_EFFORT (priority=5)

заблокировать производителя (классическое обратное давление)

REALTIME (priority=10)

заблокировать + вызвать on_backpressure() (хук для автоскалирования)

StreamStats предоставляет items_in, items_out, dropped, backpressure_events, backpressure_ms, throughput_per_sec. StreamingMultiplexer — это взвешенный round-robin fan-in по нескольким именованным источникам — полезен для слияния потоков телеметрии от N агентов.

P2.2 — Пул соединений, circuit breaker, квоты (enhancements/connection/)

MCPConnectionPool поддерживает ограниченные пулы на каждый целевой ресурс с глобальным лимитом (семафор). Простаивающие соединения переиспользуются; повреждённые удаляются через колбэк cleanup. Для каждого целевого ресурса есть собственный CircuitBreaker.

CircuitBreaker — это трёхсостоянийный (CLOSED / OPEN / HALF_OPEN) выключатель с экспоненциальной задержкой при восстановлении (backoff_base_s * 2^(trips-1), с ограничением backoff_max_s). Переходы между состояниями генерируют асинхронные колбэки on_trip, чтобы монитор здоровья мог реагировать.

QuotaManager предоставляет семафоры на ключ + ограничитель скорости на основе token-bucket + счётчик общего максимума — полезен для защиты вышестоящего LLM API от перегрузки из-за некорректно работающего воркфлоу.

P3.1 — Архитектура плагинов с горячей перезагрузкой (enhancements/plugin/)

Plugin — минимальный базовый класс с async setup(app) и async teardown(). PluginManager загружает плагины по dotted-пути (pkg.mod:Class) или пути файловой системы (./my_plugin.py), поддерживает unload() с корректным завершением и горячую перезагрузку изменённых плагинов без перезапуска процесса.

Для горячей перезагрузки используется watchdog, если он доступен (с обработчиком с дребезгом 250 мс), а в противном случае — цикл опроса на основе хэша содержимого — путь опроса важен, потому что некоторые файловые системы не доставляют события 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 оборачивает асинхронный вызываемый объект семантикой retry → fallback → восстановление состояния:

  1. RetryPolicy — экспоненциальная задержка (base_delay * multiplier^attempt, с ограничением max_delay) + джиттер; фильтр is_retriable(exc);

  2. FallbackChain — упорядоченные пары predicate → fn; выигрывает первый подходящий предикат, иначе возвращается (False, None);

  3. StateRecoverysave(workflow_id, step, state) / load(workflow_id) — хранилище снимков в памяти (можно наследовать для Redis/БД); при повторе исполнитель возобновляет работу с последнего снимка вместо повторного выполнения завершённых шагов.

ExecutionStats сообщает о попытках, успехах, сбоях, использованных fallback-ах, общем времени ожидания и последней ошибке.

P4.2 — Мониторинг здоровья и автомасштабирование (enhancements/health/)

HealthCheck оборачивает асинхронный вызываемый объект check() → (HealthStatus, detail). Внутри отслеживаются EWMA-латентность и EWMA-частота ошибок за последние 20 вызовов, а сообщаемый статус понижается на основе настраиваемых порогов (latency_warn_ms, latency_unhealthy_ms, error_rate_threshold). «Прогностическая» часть: если EWMA-частота ошибок превышает порог, проверка помечается как UNHEALTHY, даже если последний вызов завершился успешно — это позволяет выявить медленную деградацию, которую пропускают точечные пороги.

HealthMonitor запускает все зарегистрированные проверки по расписанию и генерирует колбэки on_unhealthy / on_recovered при переходах (а не при каждой проверке), чтобы избежать штормов алертов. AutoScaler подписывается на монитор: UNHEALTHY → SCALE_UP, HEALTHY + cooldown → SCALE_DOWN. Покомпонентные задержки предотвращают «дребезг».


Методология бенчмарков

Бенчмарки находятся в scripts/enhancements_benchmarks/:

  • capture_baseline.py — запускает вышестоящий набор тестов, вычисляет процент прохождения, покрытие, зонды возможностей и наивную верхнюю границу пропускной способности потоковой передачи (без работы над элементами).

  • capture_enhanced.py — запускает расширенный набор, вычисляет процент прохождения tests/enhancements/, покрытие src/mcp_agent/enhancements/ и адаптивную пропускную способность потоковой передачи (с реальной работой над элементами).

  • bench_streaming.py — прямое сравнение пропускной способности наивной asyncio.Queue и AdaptiveStreamProcessor по уровням QoS.

Все бенчмарки используют собственный тайминг asyncio (без джиттера time.time()), прогреваются на 1 000 итераций, затем измеряют 10 000 итераций. Пропускная способность сообщается как items / wall_time_s. Покрытие измеряется с помощью pytest-cov, настраиваемого через Makefile (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

Методология оценки

Оценка состоит из трёх уровней:

  1. Модульные тесты — для каждого публичного класса есть свой модуль в tests/enhancements/. 109 тестов, 0 регрессий, 82 % покрытия нового пакета.

  2. Сквозные сценарииtests/enhancements/test_end_to_end.py запускает 4 сквозных сценария, которые объединяют несколько подсистем (например, согласование протокола → пул соединений → устойчивый исполнитель с A2A-fallback → адаптивная потоковая передача → монитор здоровья + автомасштабировщик), чтобы доказать, что модули взаимодействуют, а не просто проходят изолированно.

  3. Тесты производительности на регрессииtests/enhancements/test_perf_regression.py проверяет, что пропускная способность адаптивной потоковой передачи остаётся в пределах ±30 % от зафиксированного базового уровня и значительно выше заявленного в плане минимума «100 сообщений/с». Эти тесты громко падают, если рефакторинг ухудшает пропускную способность.

Все три уровня выполняются в CI через pytest tests/enhancements/ и через цель tests в Makefile.


Внедрение сбоев

Сбои внедряются в тестах, а не с помощью отдельного инструмента chaos engineering — это делает набор тестов автономным и воспроизводимым без внешних зависимостей.

Сбой

Где внедряется

Что доказывает

Медленный потребитель (backpressure)

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

Переход half-open → 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

Цепочка fallback

test_resilience.py::test_fallback_predicate_match

Выигрывает первый подходящий предикат, нижестоящие fallback пропускаются

Восстановление состояния

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, задержка предотвращает «дребезг»


Наблюдаемость

В платформу встроены три уровня наблюдаемости:

  1. Структурированное логирование — вышестоящий пакет mcp_agent.logging (на основе Rich) не изменён; новые модули отправляют структурированные записи через тот же логгер, чтобы нижестоящие сборщики видели единый поток.

  2. Трассировка OpenTelemetry — вышестоящий пакет mcp_agent.tracing (экспортёр OTLP, semconv, счётчик токенов) не изменён; новые модули создают спаны со стабильными именами (enhancements.streaming.process, enhancements.connection.acquire, enhancements.resilience.execute_with_resilience и т. д.), чтобы дашборды работали из коробки.

  3. Сигналы здоровья и автомасштабированияHealthMonitor предоставляет объекты HealthCheckResult с EWMA-латентностью, EWMA-частотой ошибок и текущим HealthStatus. AutoScaler предоставляет события ScaleSignal. Оба можно передавать в 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").

Пропускная способность потоковой передачи

Реализация

Пропускная способность (сообщений/с)

Примечания

Наивная asyncio.Queue (верхняя граница, без работы над элементами)

352 000

Базовый уровень

AdaptiveStreamProcessor (преобразование элементов, 3 уровня QoS, отслеживание backpressure)

470 000

+34 % к базовому уровню

Заявленный базовый уровень плана

100

в 4 700 раз выше минимума

Адаптивный процессор быстрее наивного базового уровня, потому что он пакетирует обновления StreamStats и использует путь постановки в очередь с учётом уровня, что позволяет избежать ненужных asyncio.sleep(0).

Пул соединений

Метрика

Наивный (новое соединение на вызов)

С пулом

Вызовов фабрики на 1 000 циклов acquire/release

1 000

~250 (в 4 раза меньше)

Средняя задержка acquire (мкс)

~820

~210

Поверхность возможностей

Возможность

До

После

Согласование протокола с несколькими версиями

✅ (5 версий)

Обнаружение агент-к-агенту

✅ (A2A v0.3)

Адаптивная потоковая передача с QoS

✅ (3 уровня)

Пул соединений

Предохранитель

✅ (3 состояния, EWMA)

Плагины с горячей перезагрузкой

✅ (watchdog + опрос)

Реестр шаблонов рабочих процессов

✅ (на основе декораторов)

Устойчивый исполнитель с восстановлением состояния

Монитор здоровья с автоскейлером


Воспроизводимость

Каждое число выше воспроизводимо из чистого checkout:

# 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

Набор улучшений герметичен — он не требует живого API-ключа LLM, MCP-сервера или A2A-пира. Транспорт A2A inproc и наблюдатель плагинов на основе опроса означают, что каждый тест выполняется в процессе и завершается за миллисекунды.

Для вышестоящего набора некоторые интеграционные тесты требуют API-ключи; они помечены @pytest.mark.integration и пропускаются по умолчанию.


Атрибуция

Этот проект расширяет lastmile-ai/mcp-agent, лицензированный под Apache License, Version 2.0. Слои надёжности, бенчмаркинга, оценки, наблюдаемости, инъекции сбоев, тестирования, инфраструктуры и отчётности в этом репозитории являются независимо разработанными расширениями.

В частности, следующие компоненты разработаны независимо в рамках этого проекта:

  • src/mcp_agent/enhancements/ (весь пакет)

  • tests/enhancements/ (весь каталог тестов)

  • scripts/enhancements_benchmarks/

  • audit.md, ENHANCEMENTS.md, NOTICE и этот 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