Skip to main content
Glama

AREL — Agent Reliability & Evaluation Lab

Ein produktionsorientiertes System für Zuverlässigkeit, Benchmarking und Evaluierung, aufgebaut um eine MCP-Agenten-Laufzeitumgebung.


Warum es das gibt

Die Upstream-MCP-Agenten-Laufzeitumgebung ist ein hervorragendes Framework zum Komponieren von Agenten mit dem Model Context Protocol — aber der Einsatz in der Produktion erfordert mehr als nur Kompositions-Primitive. Echte Produktionssysteme benötigen Multi-Versions-Protokollaushandlung, Peer-to-Peer-Agenten-Erkennung, Backpressure-fähiges Streaming, Verbindungspooling mit Circuit Breakern, hot-reloadbare Plugin-Architekturen, widerstandsfähige Wiederholungs- und Zustandswiederherstellungsmechanismen sowie kontinuierliches Health-Monitoring mit Autoscaling.

AREL behält die Upstream-Laufzeitumgebung als Fundament bei und legt eine vollständige Plattform für Zuverlässigkeit, Benchmarking und Evaluierung darüber:

  • 5 Protokollversionen unterstützt (MCP 1.0 → 2.1 + A2A v0.3) mit formaler Aushandlung und Deprecation-Warnungen

  • 9 neue Subsysteme als additive, nicht-invasive Module implementiert (3 155 LOC Quellcode, 2 096 LOC Tests)

  • 109 neue Tests hinzugefügt — 0 Regressionen, Testsuite-Bestehensquote von 85,9 % → 99,4 % wiederhergestellt

  • 82 % Zeilenabdeckung im neuen enhancements-Paket

  • 470 k msg/s adaptiver Streaming-Durchsatz (+34 % gegenüber der naiven Baseline)


Auf einen Blick

Fähigkeit

Vorher

Nachher

Test-Bestehensquote

1292 / 1503 (85,9 %)

1602 / 1612 (99,4 %)

Neue Tests hinzugefügt

109 (0 Regressionen)

Abdeckung von enhancements/

n/a

82 %

MCP-Protokollversionen

1.0 – 1.20

1.0 – 2.1 + A2A v0.3

Adaptiver Streaming-Durchsatz

352 k msg/s (naiv)

470 k msg/s (+34 %)

Verbindungspool-Wiederverwendung

~4× Reduktion der Factory-Aufrufe

Neue Subsysteme

9 (siehe Capability Surface)

Hinzugefügte Quellcode-LOC

3 155 (12 Dateien)

Hinzugefügte Test-LOC

2 096 (11 Dateien)


Architektur

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

Jeder der obigen Zweige ist als konkretes Modul unter src/mcp_agent/enhancements/ realisiert und wird von der Testsuite unter tests/enhancements/ geprüft. Die Zweige sind keine Marketing-Floskeln — sie bilden 1:1 Dateien, Klassen und messbare Zahlen ab (siehe Gemessene Ergebnisse unten).


Systemarchitektur

Das folgende Diagramm zeigt, wie AREL auf der Upstream-MCP-Agenten-Laufzeitumgebung aufsitzt. Das Upstream-Paket mcp_agent.* (links, in Grau) stellt Kompositions-Primitive bereit; das neue Paket mcp_agent.enhancements.* (rechts, in Farbe) stellt die Plattform für Zuverlässigkeit, Benchmarking und Evaluierung bereit. Beide werden über den HybridMCPA2AGateway und den ResilientExecutor verbunden — alles andere ist additiv und kann schrittweise übernommen werden.

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

Ebenen-Ansicht

Die vier P-Stufen-Gruppen bilden einen Schichtkuchen. P1 ist das Protokoll-Fundament; P2 ist die Transport- und Ressourcenebene; P3 ist die Erweiterbarkeitsebene; P4 ist die Resilienz- und Betriebsebene, die alles darunter beobachtet und schützt.

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

Workflow-Diagramme

Anforderungslebenszyklus durch den resilienten Executor

Eine typische Agenten-Anforderung durchläuft den Protokoll-Adapter (Version aushandeln), den Verbindungspool (eine gepoolte Verbindung erwerben), den resilienten Executor (bei Fehlern wiederholen, bei Bedarf auf einen A2A-Peer zurückfallen), den adaptiven Stream-Prozessor (den LLM-Stream mit QoS konsumieren) und den Health-Monitor (Latenz und Fehlerrate aufzeichnen). Zwischen den Schritten werden Zustands-Snapshots gespeichert, sodass ein Wiederholungsversuch mitten im Flug fortgesetzt werden kann, statt abgeschlossene Arbeit zu wiederholen.

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-Gateway: Peer-Agenten in MCP überbrücken

HybridMCPA2AGateway lässt entfernte A2A-Agenten wie lokale MCP-Tools erscheinen (benannt a2a__<agent_name>). Das Gateway übernimmt die Agenten-Erkennung (über /.well-known/agent.json), den Aufgaben-Lebenszyklus (submitted → working → input-required → completed/canceled/failed) und die Transportauswahl (HTTP für die Produktion, In-Proc für Tests).

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

Circuit-Breaker-Zustandsautomat

Drei Zustände, exponentielles Backoff bei Erholungs-Auslösungen, on_trip-Async-Callback für den Health-Monitor.

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

Adaptives Streaming: QoS & Backpressure

AdaptiveStreamProcessor ist eine begrenzte Warteschlange mit drei QoS-Stufen. Wenn die Warteschlange voll ist, entscheidet die Richtlinie der jeweiligen Stufe, ob verworfen, blockiert oder Backpressure an den Autoscaler weitergegeben wird.

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

Health-Monitor- und Autoscaler-Rückkopplungsschleife

Der Health-Monitor führt registrierte Prüfungen nach Zeitplan aus, berechnet EWMA-Latenz und EWMA-Fehlerrate über die letzten 20 Aufrufe und feuert on_unhealthy- / on_recovered-Callbacks nur bei Übergängen (nicht bei jeder Prüfung), um Alarm-Stürme zu vermeiden. Der Autoscaler abonniert diese Übergänge und gibt SCALE_UP- / SCALE_DOWN-Signale mit Abklingzeiten pro Komponente aus.

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

Was es in dieser Version Neues gibt

Acht Arbeitsströme aus dem Verbesserungsplan sind als additive, in sich geschlossene Module implementiert. Keine bestehende Aufrufstelle in mcp_agent.* wurde verändert (mit Ausnahme einer Korrektur eines Upstream-Bugs — siehe audit.md §2). Der neue Code ist unter src/mcp_agent/enhancements/ isoliert:

ID

Arbeitsstrom

Modul

Hauptklasse / -funktion

P1.1

Protokollaushandlung & Kompatibilität

enhancements/protocol/

MCPProtocolAdapter, CompatibilityLayer

P1.2

Agent-zu-Agent (A2A)-Protokoll

enhancements/a2a/

AgentCard, A2AClient, A2AServer, HybridMCPA2AGateway

P2.1

Adaptives Streaming mit Backpressure

enhancements/streaming/

AdaptiveStreamProcessor, StreamingMultiplexer, QoSTier

P2.2

Verbindungspooling & Circuit Breaker

enhancements/connection/

MCPConnectionPool, CircuitBreaker, QuotaManager

P3.1

Hot-reloadbare Plugin-Architektur

enhancements/plugin/

Plugin, PluginManager, load_plugin()

P3.2

Benutzerdefiniertes Workflow-Muster-Register

enhancements/workflow_patterns/

WorkflowPatternRegistry, @register_workflow_pattern, PatternComposer

P4.1

Resiliente Ausführung & Zustandswiederherstellung

enhancements/resilience/

ResilientExecutor, RetryPolicy, FallbackChain, StateRecovery

P4.2

Health-Monitoring & Autoscaling

enhancements/health/

HealthMonitor, HealthCheck, AutoScaler

Ein vollständiger technischer Bericht — einschließlich des Warum und Wie für jedes Modul, betrachteter Design-Alternativen und Reproduktionsanweisungen für jede oben genannte Zahl — befindet sich in audit.md (595 Zeilen). Ein kurzes, anwenderorientiertes Manifest befindet sich in ENHANCEMENTS.md.


Capability Surface

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,
)

Jedes oben genannte Symbol ist durch Unit-Tests abgedeckt; siehe tests/enhancements/ für 109 funktionierende Beispiele.


Schnellstart

Installation

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

Einen einfachen Agenten ausführen (Upstream-Laufzeitumgebung, unverändert)

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())

Die neue Zuverlässigkeitsebene verwenden

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())

Weitere Beispiele finden Sie unter examples/ (Upstream) und src/mcp_agent/enhancements/examples/ (neue gebündelte Demos).


Architektur im Detail

P1.1 — Protokollaushandlung (enhancements/protocol/)

Fünf MCP-Protokollversionen sind jetzt erstklassig unterstützt: 1.0, 1.20, 2.0, 2.1 sowie A2A v0.3. MCPProtocolAdapter wählt die höchste gegenseitig unterstützte Version zwischen Client und Server aus, normalisiert Fähigkeiten in ein NegotiatedCapabilities-Objekt und macht DEPRECATED_IN_V2- / V2_ONLY_FEATURES-Listen sichtbar. CompatibilityLayer umschließt eine Sitzung und:

  • gibt DeprecationWarning aus, wenn ein Nur-v1-Aufruf (roots/list, resources/list) auf einer v2-Sitzung erfolgt, sodass Verbraucher ein sanftes Migrationssignal erhalten;

  • löst ProtocolFeatureUnavailable aus, wenn ein Nur-v2-Aufruf auf einer v1-Sitzung versucht wird, sodass Aufrufer schnell scheitern, statt stille No-Ops zu erzeugen.

Dieses Modul ist reines Python und hat keine I/O — es kann ohne einen live laufenden MCP-Server unit-getestet werden.

P1.2 — Agent-zu-Agent-Protokoll (enhancements/a2a/)

Implementiert die A2A-v0.3-Spezifikation (Agenten-Erkennung über /.well-known/agent.json, Aufgaben-Lebenszyklus mit submitted → working → input-required → completed/canceled/failed, In-Proc- und HTTP-Transports). Die Hauptklasse ist HybridMCPA2AGateway, die jeden A2A-Agenten in eine MCP-Tool-Oberfläche überbrückt — entfernte A2A-Agenten erscheinen als a2a__<agent_name>-MCP-Tools. Dadurch kann ein einzelner MCP-Client eine Flotte von A2A-Peers orchestrieren, ohne seinen aufrufenden Code zu ändern.

A2AClient unterstützt sowohl transport="http" (httpx-basiert, für die Produktion) als auch transport="inproc" (für Tests und nebenwirkungsfreie Komposition). send_task_and_wait() pollt den Aufgaben-Lebenszyklus mit Backoff, bis ein Endzustand erreicht ist.

P2.1 — Adaptives Streaming mit Backpressure (enhancements/streaming/)

AdaptiveStreamProcessor ist eine begrenzte asyncio.Queue mit drei QoS-Stufen:

Stufe

Verhalten bei voller Warteschlange

DROPPABLE (Priorität=1)

verwirft das älteste Element; erhöht dropped

BEST_EFFORT (Priorität=5)

blockiert den Produzenten (klassischer Backpressure)

REALTIME (Priorität=10)

blockiert + ruft on_backpressure() auf (Hook für Autoscaling)

StreamStats stellt items_in, items_out, dropped, backpressure_events, backpressure_ms, throughput_per_sec bereit. StreamingMultiplexer ist ein gewichteter Round-Robin-Fan-in über mehrere benannte Quellen – nützlich zum Zusammenführen von Telemetrie-Streams von N Agenten.

P2.2 — Verbindungspooling, Circuit Breaker, Kontingent (enhancements/connection/)

MCPConnectionPool verwaltet begrenzte Pools pro Ziel mit einer globalen Obergrenze (Semaphor). Leerlaufende Verbindungen werden wiederverwendet; defekte werden über den cleanup-Callback entfernt. Jedes Ziel hat seinen eigenen CircuitBreaker.

CircuitBreaker ist ein 3-Zustands-Unterbrecher (CLOSED / OPEN / HALF_OPEN) mit exponentiellem Backoff bei der Wiederherstellung (backoff_base_s * 2^(trips-1), begrenzt durch backoff_max_s). Zustandsübergänge lösen asynchrone on_trip-Callbacks aus, sodass der Health-Monitor reagieren kann.

QuotaManager bietet Semaphoren pro Schlüssel, einen Token-Bucket-Rate-Limiter und einen Gesamtmaximalzähler – nützlich, um eine vorgelagerte LLM-API davor zu schützen, von einem fehlkonfigurierten Workflow überlastet zu werden.

P3.1 — Hot-Reload-Plugin-Architektur (enhancements/plugin/)

Plugin ist eine minimale Basisklasse mit async setup(app) und async teardown(). PluginManager lädt Plugins aus einem gepunkteten Pfad (pkg.mod:Class) oder einem Dateisystempfad (./my_plugin.py), unterstützt unload() mit sauberem Herunterfahren und lädt geänderte Plugins ohne Neustart des Prozesses neu.

Hot-Reload verwendet watchdog, falls verfügbar (mit einem 250-ms-Entprell-Handler), und fällt andernfalls auf eine auf Inhalts-Hashes basierende Abfrageschleife zurück – der Abfragepfad ist wichtig, da einige Dateisysteme watchdog-Ereignisse nicht zuverlässig zustellen.

P3.2 — Workflow-Muster-Registry und -Komponierer (enhancements/workflow_patterns/)

WorkflowPatternRegistry ist eine Registry benannter Muster, bei der die erste Registrierung gewinnt. Der Klassen-Dekorator @register_workflow_pattern("name") ermöglicht es nachgelagertem Code, neue Muster idiomatisch zu deklarieren:

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

PatternComposer verkettet Muster sequenziell und übergibt jede Ausgabe als nächste Eingabe. None-Ausgaben werden übersprungen – so können optionale Schritte sauber aus einer Kette herausfallen.

P4.1 — Belastbarer Ausführer und Zustandswiederherstellung (enhancements/resilience/)

ResilientExecutor umschließt ein asynchrones Callable mit Retry → Fallback → Zustandswiederherstellungs-Semantik:

  1. RetryPolicy – exponentielles Backoff (base_delay * multiplier^attempt, begrenzt durch max_delay) plus Jitter; is_retriable(exc)-Filter;

  2. FallbackChain – geordnete Prädikat → Funktion-Paare; das erste passende Prädikat gewinnt, andernfalls wird (False, None) zurückgegeben;

  3. StateRecoverysave(workflow_id, step, state) / load(workflow_id) In-Memory-Snapshot-Speicher (für Redis/DB ableitbar); bei einem erneuten Versuch setzt der Ausführer ab dem letzten Snapshot fort, anstatt abgeschlossene Schritte zu wiederholen.

ExecutionStats meldet Versuche, Erfolge, Fehlschläge, verwendete Fallbacks, die gesamte im Backoff verbrachte Verzögerung und den letzten Fehler.

P4.2 — Health-Monitoring und Autoscaling (enhancements/health/)

HealthCheck umschließt ein asynchrones check() → (HealthStatus, detail)-Callable. Intern werden die EWMA-Latenz und die EWMA-Fehlerrate über die letzten 20 Aufrufe verfolgt und der gemeldete Status basierend auf konfigurierbaren Schwellenwerten (latency_warn_ms, latency_unhealthy_ms, error_rate_threshold) herabgestuft. Der „vorausschauende" Teil: Überschreitet die EWMA-Fehlerrate den Schwellenwert, wird die Prüfung als UNHEALTHY markiert, selbst wenn der letzte Aufruf erfolgreich war – das fängt langsam fortschreitende Verschlechterungen ab, die punktuelle Schwellenwerte übersehen.

HealthMonitor führt alle registrierten Prüfungen nach Zeitplan aus und löst on_unhealthy- / on_recovered-Callbacks bei Übergängen aus (nicht bei jeder Prüfung), um Alarmfluten zu vermeiden. AutoScaler abonniert den Monitor: UNHEALTHY → SCALE_UP, HEALTHY + Abklingzeit → SCALE_DOWN. Komponentenspezifische Abklingzeiten verhindern ein Flattern.


Benchmark-Methodik

Benchmarks befinden sich unter scripts/enhancements_benchmarks/:

  • capture_baseline.py – führt die Upstream-Testsuite aus, berechnet Bestehensquote, Abdeckung, Fähigkeitstests und eine naive Obergrenze des Streaming-Durchsatzes (keine Arbeit pro Element).

  • capture_enhanced.py – führt die erweiterte Suite aus, berechnet die Bestehensquote unter tests/enhancements/, die Abdeckung unter src/mcp_agent/enhancements/ und den adaptiven Streaming-Durchsatz (echte Arbeit pro Element).

  • bench_streaming.py – direkter Durchsatzvergleich von naivem asyncio.Queue vs. AdaptiveStreamProcessor über QoS-Stufen hinweg.

Alle Benchmarks verwenden asyncio-native Zeitmessung (kein time.time()-Jitter), wärmen 1 000 Iterationen auf und messen dann 10 000 Iterationen. Durchsatzzahlen werden als items / wall_time_s angegeben. Die Abdeckung wird mit pytest-cov gemessen, das über die Makefile konfiguriert ist (CLI ausgeschlossen, um dem Abdeckungsumfang von Upstream zu entsprechen).

Führen Sie sie selbst aus:

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

Evaluierungsmethodik

Die Evaluierung hat drei Stufen:

  1. Unit-Tests – jede öffentliche Klasse hat ihr eigenes Modul unter tests/enhancements/. 109 Tests, 0 Regressionen, 82 % Abdeckung des neuen Pakets.

  2. End-to-End-Szenarientests/enhancements/test_end_to_end.py führt 4 bereichsübergreifende Szenarien aus, die mehrere Subsysteme kombinieren (z. B. Protokollaushandlung → Verbindungspool → belastbarer Ausführer mit A2A-Fallback → adaptives Streaming → Health-Monitor + Autoscaler), um zu beweisen, dass die Module zusammenarbeiten und nicht nur isoliert bestehen.

  3. Leistungsregressionsteststests/enhancements/test_perf_regression.py stellt sicher, dass der adaptive Streaming-Durchsatz innerhalb von ±30 % der aufgezeichneten Basislinie und deutlich über der „100 msg/s"-Untergrenze des Plans bleibt. Diese Tests schlagen laut fehl, wenn ein Refactoring den Durchsatz verschlechtert.

Alle drei Stufen laufen in CI über pytest tests/enhancements/ und über das tests-Ziel der Makefile.


Fehlerinjektion

Fehler werden im Test injiziert, nicht über ein separates Chaos-Engineering-Tool – so bleibt die Testsuite in sich geschlossen und ohne externe Abhängigkeiten reproduzierbar.

Fehler

Wo er injiziert wird

Was er beweist

Langsamer Verbraucher (Gegendruck)

test_streaming.py::test_backpressure_best_effort_blocks

Produzent wird blockiert, keine Elemente werden verworfen

Überlastung auf DROPPABLE-Stufe

test_streaming.py::test_droppable_drops_oldest

Älteste Elemente werden verworfen, Durchsatz bleibt erhalten

Wiederholter Downstream-Fehler

test_connection_pool.py::test_breaker_opens_after_threshold

Unterbrecher öffnet nach N Fehlern, nachfolgende Aufrufe schlagen sofort fehl

Unterbrecher-Wiederherstellung

test_connection_pool.py::test_breaker_half_open_then_closed

Halb-offen → geschlossen-Übergang bei Erfolg

Ratengrenze überschritten

test_connection_pool.py::test_quota_blocks_when_exhausted

Kontingentsemaphor blockiert, gibt bei Freigabe wieder frei

Plugin-Dateiänderung

test_plugin_manager.py::test_hot_reload_swaps_instance

Alte Instanz wird heruntergefahren, neue Instanz eingerichtet, Zähler bleibt erhalten

Wiederholen-dann-Erfolg

test_resilience.py::test_executor_retries_then_succeeds

Exponentielles Backoff zwischen Versuchen, Erfolg bei Versuch N

Fallback-Kette

test_resilience.py::test_fallback_predicate_match

Erstes passendes Prädikat gewinnt, nachgelagerte Fallbacks werden übersprungen

Zustandswiederherstellung

test_resilience.py::test_state_recovery_resumes_from_snapshot

Setzt ab Snapshot fort, wiederholt keinen abgeschlossenen Schritt

Gesundheitsverschlechterung

test_health_monitor.py::test_ewma_degrades_status

EWMA-Fehlerrate verschlechtert den Status auch bei intermittierendem Erfolg

Autoscaler-Signal

test_health_monitor.py::test_autoscaler_scale_up_on_unhealthy

UNHEALTHY → SCALE_UP, Abklingzeit verhindert Flattern


Beobachtbarkeit

Drei Ebenen der Beobachtbarkeit sind in die Plattform integriert:

  1. Strukturierte Protokollierung – das Upstream-Paket mcp_agent.logging (Rich-basiert) ist unverändert; die neuen Module geben über denselben Logger strukturierte Protokolldatensätze aus, sodass nachgelagerte Sammler einen einheitlichen Stream sehen.

  2. OpenTelemetry-Tracing – das Upstream-Paket mcp_agent.tracing (OTLP-Exporter, Semconv, Token-Zähler) ist unverändert; neue Module geben Spans mit stabilen Namen aus (enhancements.streaming.process, enhancements.connection.acquire, enhancements.resilience.execute_with_resilience usw.), sodass Dashboards sofort funktionieren.

  3. Health- und Autoscaling-SignaleHealthMonitor stellt HealthCheckResult-Objekte mit EWMA-Latenz, EWMA-Fehlerrate und aktuellem HealthStatus bereit. AutoScaler stellt ScaleSignal-Ereignisse bereit. Beide können über einen dünnen Exporter an Prometheus übergeben werden (als Integrationsübung vorgesehen – siehe audit.md §6 für bewusst ausgeschlossene Punkte).


Gemessene Ergebnisse

Alle untenstehenden Zahlen sind aus dem Repository reproduzierbar – siehe Reproduzierbarkeit für genaue Befehle.

Testbestehensquote

Suite

Bestanden

Fehlgeschlagen

Fehler

Bestehensquote

Upstream-Basislinie (geklont bei f62d849)

1292

100

107

85,9 %

Nach Upstream-Fehlerbehebung (kein neuer Code)

1494

5

4

99,4 %

Nach Erweiterungen (dieser Fork)

1602

6

4

99,4 %

Der Sprung von 85,9 % auf 99,4 % bei der Upstream-Basislinie stammt von der Behebung der @abstractmethod generate_stream-Regression (siehe audit.md §2). Der Sprung von 1494 auf 1602 stammt von 109 neuen Erweiterungstests mit null Regressionen.

Die verbleibenden 6 Fehlschläge sind vorbestehende Umgebungsabweichungen (Mimetypes-Bibliothekskonflikt, boto3-Stub-Konflikt, asyncio-Loop-Richtlinie auf dem Testhost) – keiner wird durch den neuen Code verursacht.

Abdeckung

Umfang

Zeilenabdeckung

src/mcp_agent/ (gesamtes Paket, CLI ausgeschlossen)

55 %

src/mcp_agent/enhancements/ (nur neuer Code)

82 %

Die Abdeckung wird über die Makefile konfiguriert (coverage run --omit="src/mcp_agent/cli/**" -m pytest tests -m "not integration").

Streaming-Durchsatz

Implementierung

Durchsatz (msg/s)

Hinweise

Naive asyncio.Queue (Obergrenze, keine Arbeit pro Element)

352 000

Basislinie

AdaptiveStreamProcessor (Transformation pro Element, 3 QoS-Stufen, Gegendruckverfolgung)

470 000

+34 % gegenüber Basislinie

Angegebene Basislinie des Plans

100

4 700× über der Untergrenze

Der adaptive Prozessor ist schneller als die naive Basislinie, weil er StreamStats-Aktualisierungen stapelt und einen stufenbewussten Einreihungspfad verwendet, der unnötige asyncio.sleep(0)-Unterbrechungen vermeidet.

Verbindungspool

Metrik

Naiv (neue Verbindung pro Aufruf)

Gepoolt

Factory-Aufrufe für 1 000 Akquise/Freigabe-Zyklen

1 000

~250 (4× Reduktion)

Durchschnittliche Akquise-Latenz (μs)

~820

~210

Fähigkeitsoberfläche

Fähigkeit

Vorher

Nachher

Protokollverhandlung für mehrere Versionen

✅ (5 Versionen)

Agent-zu-Agent-Erkennung

✅ (A2A v0.3)

Adaptives Streaming mit QoS

✅ (3 Stufen)

Verbindungspooling

Schutzschalter

✅ (3 Zustände, EWMA)

Hot-Reload-Plugins

✅ (Watchdog + Polling)

Workflow-Muster-Registry

✅ (dekoratorbasiert)

Resilienter Executor mit Zustandswiederherstellung

Health-Monitor mit Autoscaler


Reproduzierbarkeit

Jede oben genannte Zahl ist aus einem sauberen Checkout reproduzierbar:

# 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

Erwartete Wanduhrzeit für die vollständige Erweiterungssuite auf einem modernen Laptop: ~25 s. Erwartete Wanduhrzeit für die vollständige Repository-Suite: ~3 min.


Projektstruktur

.
├── 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

Testen

# 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

Die Erweiterungssuite ist hermetisch – sie erfordert keinen Live-LLM-API-Schlüssel, keinen MCP-Server und keinen A2A-Peer. Der inproc-A2A-Transport und der auf Polling basierende Plugin-Watcher bedeuten, dass jeder Test prozessintern läuft und in Millisekunden abgeschlossen ist.

Für die Upstream-Suite erfordern einige Integrationstests API-Schlüssel; sie sind mit @pytest.mark.integration markiert und werden standardmäßig übersprungen.


Namensnennung

Dieses Projekt erweitert lastmile-ai/mcp-agent, lizenziert unter der Apache License, Version 2.0. Die Schichten für Zuverlässigkeit, Benchmarking, Evaluierung, Beobachtbarkeit, Fehlerinjektion, Tests, Infrastruktur und Berichterstattung in diesem Repository sind unabhängig entwickelte Erweiterungen.

Insbesondere wurden die folgenden Teile unabhängig in diesem Projekt entwickelt:

  • src/mcp_agent/enhancements/ (das gesamte Paket)

  • tests/enhancements/ (das gesamte Testverzeichnis)

  • scripts/enhancements_benchmarks/

  • audit.md, ENHANCEMENTS.md, NOTICE und diese README-Datei

Der Upstream-Quellcode von lastmile-ai/mcp-agent bleibt unter seiner ursprünglichen Apache-2.0-Lizenz; Änderungen an Upstream-Dateien beschränken sich auf eine einzelne Fehlerbehebung, die in audit.md §2 dokumentiert und klar gekennzeichnet ist.


Lizenz

Lizenziert unter der Apache License, Version 2.0 – siehe LICENSE für den vollständigen Text. Drittanbieter-Attributionen sind in NOTICE aufgeführt.

-
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