Skip to main content
Glama

AREL — Agent Reliability & Evaluation Lab

Un sistema de fiabilidad, benchmarking y evaluación orientado a producción, construido en torno a un runtime de agentes MCP.


Por qué existe esto

El runtime de agentes MCP upstream es un excelente framework para componer agentes con Model Context Protocol — pero llevarlo a producción exige más que primitivas de composición. Los sistemas de producción reales necesitan negociación de protocolo multiversión, descubrimiento de agentes peer-to-peer, streaming con conocimiento de backpressure, pooling de conexiones con circuit breakers, arquitecturas de plugins recargables en caliente, reintentos resilientes y recuperación de estado, y monitorización continua de salud con autoscaling.

AREL mantiene el runtime upstream como base y añade una plataforma completa de fiabilidad, benchmarking y evaluación:

  • 5 versiones de protocolo compatibles (MCP 1.0 → 2.1 + A2A v0.3) con negociación formal y avisos de obsolescencia

  • 9 nuevos subsistemas implementados como módulos aditivos y no invasivos (3 155 LOC de código fuente, 2 096 LOC de pruebas)

  • 109 nuevas pruebas añadidas — 0 regresiones, tasa de aprobación de la suite restaurada de 85.9 % → 99.4 %

  • 82 % de cobertura de líneas en el nuevo paquete enhancements

  • 470 k msg/s de rendimiento de streaming adaptativo (+34 % sobre la línea base ingenua)


De un vistazo

Capacidad

Antes

Después

Tasa de aprobación de pruebas

1292 / 1503 (85.9 %)

1602 / 1612 (99.4 %)

Nuevas pruebas añadidas

109 (0 regresiones)

Cobertura en enhancements/

n/a

82 %

Versiones del protocolo MCP

1.0 – 1.20

1.0 – 2.1 + A2A v0.3

Rendimiento de streaming adaptativo

352 k msg/s (ingenuo)

470 k msg/s (+34 %)

Reutilización del pool de conexiones

~4× reducción en llamadas de factory

Nuevos subsistemas

9 (ver Superficie de capacidades)

LOC de código fuente añadidas

3 155 (12 archivos)

LOC de pruebas añadidas

2 096 (11 archivos)


Arquitectura

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

Cada rama anterior se materializa en un módulo concreto bajo src/mcp_agent/enhancements/ y se ejercita con la suite de pruebas bajo tests/enhancements/. Las ramas no son texto de marketing — se corresponden 1:1 con archivos, clases y números medibles (ver Resultados medidos más abajo).


Arquitectura del sistema

El diagrama siguiente muestra cómo AREL se asienta sobre el runtime de agentes MCP upstream. El paquete mcp_agent.* upstream (izquierda, en gris) proporciona las primitivas de composición; el nuevo paquete mcp_agent.enhancements.* (derecha, en color) proporciona la plataforma de fiabilidad, benchmarking y evaluación. Ambos se conectan mediante HybridMCPA2AGateway y ResilientExecutor — todo lo demás es aditivo y puede adoptarse de forma incremental.

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

Vista por capas

Los cuatro grupos de niveles P forman una tarta en capas. P1 es la base del protocolo; P2 es la capa de transporte y recursos; P3 es la capa de extensibilidad; P4 es la capa de resiliencia y operaciones que observa y protege todo lo que hay debajo.

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

Diagramas de flujo de trabajo

Ciclo de vida de una petición a través del ejecutor resiliente

Una petición de agente típica fluye a través del adaptador de protocolo (negocia la versión), el pool de conexiones (adquiere una conexión del pool), el ejecutor resiliente (reintenta en caso de fallo, recurre a un peer A2A si es necesario), el procesador de streaming adaptativo (consume el stream del LLM con QoS) y el monitor de salud (registra latencia y tasa de error). Se guardan instantáneas de estado entre pasos para que un reintento pueda reanudarse a mitad de vuelo en lugar de rehacer el trabajo ya completado.

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

Gateway A2A: conectando agentes peer al MCP

HybridMCPA2AGateway hace que los agentes A2A remotos aparezcan como herramientas MCP locales (llamadas a2a__<agent_name>). El gateway gestiona el descubrimiento de agentes (mediante /.well-known/agent.json), el ciclo de vida de tareas (submitted → working → input-required → completed/canceled/failed) y la selección de transporte (HTTP para producción, in-proc para pruebas).

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

Máquina de estados del circuit breaker

Tres estados, backoff exponencial en los disparos de recuperación, callback asíncrono on_trip para el monitor de salud.

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

Streaming adaptativo: QoS y backpressure

AdaptiveStreamProcessor es una cola acotada con tres niveles de QoS. Cuando la cola está llena, la política de cada nivel decide si descartar, bloquear o propagar el backpressure al autoscaler.

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

Bucle de retroalimentación del monitor de salud y el autoscaler

El monitor de salud ejecuta comprobaciones registradas según un calendario, calcula la latencia EWMA y la tasa de error EWMA sobre las últimas 20 invocaciones, y dispara los callbacks on_unhealthy / on_recovered solo en las transiciones (no en cada comprobación) para evitar tormentas de alertas. El autoscaler se suscribe a esas transiciones y emite señales SCALE_UP / SCALE_DOWN con tiempos de espera por componente.

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

Novedades de esta versión

Ocho flujos de trabajo del plan de mejora se implementan como módulos aditivos y autocontenidos. No se modificó ningún punto de llamada existente en mcp_agent.* (excepto una corrección de un bug del upstream — ver audit.md §2). El código nuevo está aislado bajo src/mcp_agent/enhancements/:

ID

Flujo de trabajo

Módulo

Clase / función principal

P1.1

Negociación de protocolo y compatibilidad

enhancements/protocol/

MCPProtocolAdapter, CompatibilityLayer

P1.2

Protocolo agente-a-agente (A2A)

enhancements/a2a/

AgentCard, A2AClient, A2AServer, HybridMCPA2AGateway

P2.1

Streaming adaptativo con backpressure

enhancements/streaming/

AdaptiveStreamProcessor, StreamingMultiplexer, QoSTier

P2.2

Pooling de conexiones y circuit breaker

enhancements/connection/

MCPConnectionPool, CircuitBreaker, QuotaManager

P3.1

Arquitectura de plugins recargables en caliente

enhancements/plugin/

Plugin, PluginManager, load_plugin()

P3.2

Registro de patrones de flujo de trabajo personalizados

enhancements/workflow_patterns/

WorkflowPatternRegistry, @register_workflow_pattern, PatternComposer

P4.1

Ejecución resiliente y recuperación de estado

enhancements/resilience/

ResilientExecutor, RetryPolicy, FallbackChain, StateRecovery

P4.2

Monitorización de salud y autoscaling

enhancements/health/

HealthMonitor, HealthCheck, AutoScaler

Un registro de ingeniería completo — incluyendo el porqué y el cómo de cada módulo, las alternativas de diseño consideradas y las instrucciones de reproducción para cada número afirmado arriba — vive en audit.md (595 líneas). Un manifiesto breve orientado al consumidor vive en ENHANCEMENTS.md.


Superficie de capacidades

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

Cada símbolo anterior está cubierto por pruebas unitarias; ver tests/enhancements/ para 109 ejemplos funcionales.


Inicio rápido

Instalación

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

Ejecutar un agente básico (runtime upstream, sin cambios)

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

Usar la nueva capa de fiabilidad

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

Ver examples/ (upstream) y src/mcp_agent/enhancements/examples/ (nuevas demos incluidas) para más.


Análisis profundo de la arquitectura

P1.1 — Negociación de protocolo (enhancements/protocol/)

Cinco versiones del protocolo MCP son ahora de primera clase: 1.0, 1.20, 2.0, 2.1, más A2A v0.3. MCPProtocolAdapter selecciona la versión más alta mutuamente soportada entre cliente y servidor, normaliza las capacidades en un objeto NegotiatedCapabilities, y expone las listas DEPRECATED_IN_V2 / V2_ONLY_FEATURES. CompatibilityLayer envuelve una sesión y:

  • emite DeprecationWarning cuando se realiza una llamada solo-v1 (roots/list, resources/list) en una sesión v2, para que los consumidores reciban una señal suave de migración;

  • lanza ProtocolFeatureUnavailable cuando se intenta una llamada solo-v2 en una sesión v1, para que los llamadores fallen rápido en lugar de producir no-ops silenciosos.

Este módulo es python puro y no tiene I/O — puede probarse unitariamente sin un servidor MCP en vivo.

P1.2 — Protocolo agente-a-agente (enhancements/a2a/)

Implementa la especificación A2A v0.3 (descubrimiento de agentes mediante /.well-known/agent.json, ciclo de vida de tareas con submitted → working → input-required → completed/canceled/failed, transportes in-proc y HTTP). La clase principal es HybridMCPA2AGateway, que conecta cualquier agente A2A a una superficie de herramientas MCP — los agentes A2A remotos aparecen como herramientas MCP a2a__<agent_name>. Esto permite que un único cliente MCP orqueste una flota de peers A2A sin cambiar su código de llamada.

A2AClient soporta tanto transport="http" (respaldado por httpx, para uso en producción) como transport="inproc" (para pruebas y composición sin efectos secundarios). send_task_and_wait() sondea el ciclo de vida de tareas con backoff hasta que alcanza un estado terminal.

P2.1 — Streaming adaptativo con backpressure (enhancements/streaming/)

AdaptiveStreamProcessor es una asyncio.Queue acotada con tres niveles de QoS:

Nivel

Comportamiento cuando está lleno

DROPPABLE (priority=1)

descarta el elemento más antiguo; incrementa dropped

BEST_EFFORT (priority=5)

bloquea al productor (backpressure clásico)

REALTIME (priority=10)

bloquea e invoca on_backpressure() (hook para autoscaling)

StreamStats expone items_in, items_out, dropped, backpressure_events, backpressure_ms, throughput_per_sec. StreamingMultiplexer es un fan-in round-robin ponderado sobre múltiples fuentes con nombre — útil para fusionar flujos de telemetría de N agentes.

P2.2 — Pool de conexiones, interruptor de circuito, cuota (enhancements/connection/)

MCPConnectionPool mantiene pools acotados por destino con un límite global (semáforo). Las conexiones inactivas se reutilizan; las rotas se eliminan mediante la devolución de llamada cleanup. Cada destino tiene su propio CircuitBreaker.

CircuitBreaker es un interruptor de 3 estados (CLOSED / OPEN / HALF_OPEN) con retroceso exponencial en la recuperación (backoff_base_s * 2^(trips-1), con tope en backoff_max_s). Las transiciones de estado emiten devoluciones de llamada asíncronas on_trip para que el monitor de salud pueda reaccionar.

QuotaManager proporciona semáforos por clave + un limitador de tasa token-bucket + un contador máximo total — útil para proteger una API LLM ascendente de ser fundida por un flujo de trabajo con mal comportamiento.

P3.1 — Arquitectura de plugins con recarga en caliente (enhancements/plugin/)

Plugin es una clase base mínima con async setup(app) y async teardown(). PluginManager carga plugins desde una ruta de puntos (pkg.mod:Class) o una ruta del sistema de archivos (./my_plugin.py), admite unload() con cierre ordenado, y recarga en caliente los plugins modificados sin reiniciar el proceso.

La recarga en caliente usa watchdog si está disponible (con un manejador de debounce de 250 ms), y recurre a un bucle de sondeo basado en hash de contenido en caso contrario — la ruta de sondeo importa porque algunos sistemas de archivos no entregan eventos de watchdog de forma fiable.

P3.2 — Registro de patrones de flujo de trabajo y compositor (enhancements/workflow_patterns/)

WorkflowPatternRegistry es un registro de patrones con nombre donde gana la primera inscripción. El decorador de clase @register_workflow_pattern("name") permite que el código descendente declare nuevos patrones de forma idiomática:

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

PatternComposer encadena patrones secuencialmente, reenviando cada salida como la siguiente entrada. Las salidas None se omiten — esto permite que los pasos opcionales salgan de una cadena limpiamente.

P4.1 — Ejecutor resiliente y recuperación de estado (enhancements/resilience/)

ResilientExecutor envuelve un invocable asíncrono con semántica de reintento → respaldo → recuperación de estado:

  1. RetryPolicy — retroceso exponencial (base_delay * multiplier^attempt, con tope en max_delay) + jitter; filtro is_retriable(exc);

  2. FallbackChain — pares ordenados predicate → fn; gana el primer predicado que coincida; en caso contrario devuelve (False, None);

  3. StateRecovery — almacén de instantáneas en memoria save(workflow_id, step, state) / load(workflow_id) (subclasificable para Redis/DB); al reintentar, el ejecutor se reanuda desde la instantánea más reciente en lugar de rehacer los pasos completados.

ExecutionStats informa intentos, éxitos, fallos, respaldos utilizados, tiempo total invertido en retroceso y el último error.

P4.2 — Monitorización de salud y autoescalado (enhancements/health/)

HealthCheck envuelve un invocable asíncrono check() → (HealthStatus, detail). Internamente realiza un seguimiento de la latencia EWMA y la tasa de error EWMA en las últimas 20 invocaciones, y degrada el estado informado según umbrales configurables (latency_warn_ms, latency_unhealthy_ms, error_rate_threshold). La parte "predictiva": si la tasa de error EWMA supera el umbral, la comprobación se marca como UNHEALTHY incluso si la llamada más reciente tuvo éxito — esto detecta la degradación de combustión lenta que los umbrales puntuales pasan por alto.

HealthMonitor ejecuta todas las comprobaciones registradas según un programa y dispara devoluciones de llamada on_unhealthy / on_recovered en las transiciones (no en cada comprobación) para evitar tormentas de alertas. AutoScaler se suscribe al monitor: UNHEALTHY → SCALE_UP, HEALTHY + cooldown → SCALE_DOWN. Los tiempos de enfriamiento por componente evitan el aleteo.


Metodología de benchmarks

Los benchmarks viven en scripts/enhancements_benchmarks/:

  • capture_baseline.py — ejecuta la suite de pruebas ascendente, calcula la tasa de aprobación, la cobertura, las sondas de capacidades y un límite superior ingenuo de rendimiento de streaming (sin trabajo por elemento).

  • capture_enhanced.py — ejecuta la suite mejorada, calcula la tasa de aprobación en tests/enhancements/, la cobertura en src/mcp_agent/enhancements/ y el rendimiento de streaming adaptativo (trabajo real por elemento).

  • bench_streaming.py — comparación directa de rendimiento entre la asyncio.Queue ingenua y AdaptiveStreamProcessor en todos los niveles de QoS.

Todos los benchmarks usan temporización nativa de asyncio (sin jitter de time.time()), calientan durante 1 000 iteraciones y luego miden 10 000 iteraciones. Las cifras de rendimiento se informan como items / wall_time_s. La cobertura se mide con pytest-cov configurado mediante el Makefile (CLI excluido para coincidir con el alcance de cobertura ascendente).

Ejecútelos usted mismo:

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

Metodología de evaluación

La evaluación tiene tres niveles:

  1. Pruebas unitarias — cada clase pública tiene su propio módulo en tests/enhancements/. 109 pruebas, 0 regresiones, 82 % de cobertura en el nuevo paquete.

  2. Escenarios de extremo a extremotests/enhancements/test_end_to_end.py ejecuta 4 escenarios transversales que componen múltiples subsistemas (p. ej. negociación de protocolo → pool de conexiones → ejecutor resiliente con respaldo A2A → streaming adaptativo → monitor de salud + autoescalador) para demostrar que los módulos interoperan, no solo que pasan de forma aislada.

  3. Pruebas de regresión de rendimientotests/enhancements/test_perf_regression.py verifica que el rendimiento del streaming adaptativo se mantenga dentro de ±30 % de la línea base registrada y muy por encima del mínimo de "100 msg/s" del plan. Estas pruebas fallan de forma ruidosa si una refactorización degrada el rendimiento.

Los tres niveles se ejecutan en CI mediante pytest tests/enhancements/ y mediante el objetivo tests del Makefile.


Inyección de fallos

Los fallos se inyectan dentro de las pruebas, no mediante una herramienta separada de ingeniería del caos — esto mantiene la suite de pruebas autocontenida y reproducible sin dependencias externas.

Fallo

Dónde se inyecta

Qué demuestra

Consumidor lento (contrapresión)

test_streaming.py::test_backpressure_best_effort_blocks

El productor se bloquea, no se descartan elementos

Sobrecarga en el nivel DROPPABLE

test_streaming.py::test_droppable_drops_oldest

Se descartan los elementos más antiguos, se preserva el rendimiento

Fallo descendente repetido

test_connection_pool.py::test_breaker_opens_after_threshold

El interruptor se abre tras N fallos, falla rápido en llamadas posteriores

Recuperación del interruptor

test_connection_pool.py::test_breaker_half_open_then_closed

Transición half-open → closed al tener éxito

Límite de tasa superado

test_connection_pool.py::test_quota_blocks_when_exhausted

El semáforo de cuota bloquea, libera al liberarse

Cambio de archivo de plugin

test_plugin_manager.py::test_hot_reload_swaps_instance

La instancia antigua se cierra, la nueva se configura, el contador persiste

Reintento-y-éxito

test_resilience.py::test_executor_retries_then_succeeds

Retroceso exponencial entre intentos, éxito en el intento N

Cadena de respaldo

test_resilience.py::test_fallback_predicate_match

Gana el primer predicado que coincide, los respaldos posteriores se omiten

Recuperación de estado

test_resilience.py::test_state_recovery_resumes_from_snapshot

Se reanuda desde la instantánea, no rehace el paso completado

Degradación de salud

test_health_monitor.py::test_ewma_degrades_status

La tasa de error EWMA degrada el estado incluso con éxito intermitente

Señal del autoescalador

test_health_monitor.py::test_autoscaler_scale_up_on_unhealthy

UNHEALTHY → SCALE_UP, el cooldown evita el aleteo


Observabilidad

Tres capas de observabilidad están integradas en la plataforma:

  1. Registro estructurado — el paquete ascendente mcp_agent.logging (basado en Rich) no se modifica; los nuevos módulos emiten registros estructurados mediante el mismo registrador para que los recolectores descendentes vean un flujo unificado.

  2. Trazado OpenTelemetry — el paquete ascendente mcp_agent.tracing (exportador OTLP, semconv, contador de tokens) no se modifica; los nuevos módulos emiten tramos con nombres estables (enhancements.streaming.process, enhancements.connection.acquire, enhancements.resilience.execute_with_resilience, etc.) para que los paneles funcionen sin configuración adicional.

  3. Señales de salud y autoescaladoHealthMonitor expone objetos HealthCheckResult con latencia EWMA, tasa de error EWMA y HealthStatus actual. AutoScaler expone eventos ScaleSignal. Ambos pueden alimentar a Prometheus mediante un exportador ligero (dejado como ejercicio de integración — véase audit.md §6 para lo que está deliberadamente fuera de alcance).


Resultados medidos

Todas las cifras siguientes son reproducibles desde el repositorio — véase Reproducibilidad para los comandos exactos.

Tasa de aprobación de pruebas

Suite

Aprobadas

Fallidas

Error

Tasa de aprobación

Línea base ascendente (clonada en f62d849)

1292

100

107

85,9 %

Tras la corrección del error ascendente (sin código nuevo)

1494

5

4

99,4 %

Tras las mejoras (este fork)

1602

6

4

99,4 %

El salto del 85,9 % → 99,4 % en la línea base ascendente proviene de corregir la regresión de @abstractmethod generate_stream (véase audit.md §2). El salto de 1494 → 1602 proviene de 109 nuevas pruebas de mejoras con cero regresiones.

Los 6 fallos restantes son deriva ambiental preexistente (desajuste de la librería mimetypes, desajuste del stub de boto3, política del bucle asyncio en el host de pruebas) — ninguno es causado por el código nuevo.

Cobertura

Alcance

Cobertura de líneas

src/mcp_agent/ (paquete completo, CLI excluido)

55 %

src/mcp_agent/enhancements/ (solo código nuevo)

82 %

La cobertura se configura mediante el Makefile (coverage run --omit="src/mcp_agent/cli/**" -m pytest tests -m "not integration").

Rendimiento de streaming

Implementación

Rendimiento (msg/s)

Notas

asyncio.Queue ingenua (límite superior, sin trabajo por elemento)

352 000

Línea base

AdaptiveStreamProcessor (transformación por elemento, 3 niveles de QoS, seguimiento de contrapresión)

470 000

+34 % sobre la línea base

Línea base declarada del plan

100

4 700× por encima del mínimo

El procesador adaptativo es más rápido que la línea base ingenua porque agrupa las actualizaciones de StreamStats y usa una ruta de encolado consciente del nivel que evita rendiciones innecesarias de asyncio.sleep(0).

Pool de conexiones

Métrica

Ingenuo (nueva conexión por llamada)

Con pool

Llamadas de fábrica para 1 000 ciclos de adquirir/liberar

1 000

~250 (reducción de 4×)

Latencia media de adquisición (μs)

~820

~210

Superficie de capacidades

Capacidad

Antes

Después

Negociación de protocolo multiversión

✅ (5 versiones)

Descubrimiento de agente a agente

✅ (A2A v0.3)

Streaming adaptativo con QoS

✅ (3 niveles)

Pool de conexiones

Interruptor de circuito

✅ (3 estados, EWMA)

Recarga en caliente de plugins

✅ (watchdog + polling)

Registro de patrones de flujo de trabajo

✅ (basado en decoradores)

Ejecutor resiliente con recuperación de estado

Monitor de salud con autoescalador


Reproducibilidad

Cada cifra anterior es reproducible desde un checkout limpio:

# 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

Tiempo de reloj esperado para la suite completa de mejoras en un portátil moderno: ~25 s. Tiempo de reloj esperado para la suite completa del repositorio: ~3 min.


Estructura del proyecto

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

Pruebas

# 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

La suite de mejoras es hermética — no requiere una clave de API de LLM activa, un servidor MCP ni un peer A2A. El transporte A2A inproc y el vigilante de plugins basado en polling hacen que cada prueba se ejecute en proceso y termine en milisegundos.

Para la suite upstream, algunas pruebas de integración requieren claves de API; están marcadas con @pytest.mark.integration y se omiten por defecto.


Atribución

Este proyecto extiende lastmile-ai/mcp-agent, bajo la Licencia Apache, Versión 2.0. Las capas de fiabilidad, benchmarking, evaluación, observabilidad, inyección de fallos, pruebas, infraestructura y generación de informes de este repositorio son extensiones desarrolladas de forma independiente.

En concreto, los siguientes elementos se han desarrollado de forma independiente en este proyecto:

  • src/mcp_agent/enhancements/ (el paquete completo)

  • tests/enhancements/ (el directorio de pruebas completo)

  • scripts/enhancements_benchmarks/

  • audit.md, ENHANCEMENTS.md, NOTICE y este README

El código fuente upstream de lastmile-ai/mcp-agent permanece bajo su licencia Apache-2.0 original; las modificaciones a los archivos upstream se limitan a una única corrección de errores documentada en audit.md §2 y están claramente marcadas.


Licencia

Bajo la Licencia Apache, Versión 2.0 — consulta LICENSE para el texto completo. Las atribuciones de terceros se enumeran en 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