Todos los artículos

// Arquitectura que aguanta

Arquitectura Hexagonal con Agentes IA: Patrones de Integración Prácticos

Integra MCP y agentes IA con arquitectura hexagonal. Patrones producción-ready en Python FastAPI pgvector para sistemas escalables y testeables.

6 de julio de 202614 min de lectura

Problem

La llegada de los Model Context Protocol (MCP) y los agentes IA ha desafiado las arquitecturas backend tradicionales. En desarrollo, conectar un agente con un MCP server parece trivial:

python
# Prototype: conexión directa a MCP
from anthropic import Anthropic
from mcp import ClientSession

client = Anthropic()
mcp_session = ClientSession("localhost:3000")

async def process_query(query: str):
    tools = await mcp_session.list_tools()
    response = client.messages.create(
        model="claude-3.5-sonnet",
        messages=[{"role": "user", "content": query}],
        tools=tools,
    )
    return response

Este enfoque funciona perfectamente hasta que intentas escalar a producción. Los problemas que enfrentamos en Orbit Analytics incluyen:

  • Acoplamiento directo: Cada cambio en el MCP server rompe el código del agente
  • Testabilidad imposible: Tests de integración requieren MCP servers reales
  • Dificultad de mocking: Sin interfaces claras, no puedes simular tools en tests unitarios
  • Eventos de dominio perdidos: Cuando un agente modifica datos, no hay un contrato de domain events
  • Integración de vectores: pgvector se mezcla con lógica de negocio sin separación

La realidad es que las arquitecturas tradicionales de capas no manejan bien la naturaleza dinámica de los agentes IA. Los MCP tools cambian frecuentemente, los modelos se actualizan, y las interacciones son inherentemente no determinísticas.

Core Concept

La Arquitectura Hexagonal (Ports & Adapters) ofrece un marco conceptual perfecto para integrar agentes IA de forma limpia. La clave es tratar a los MCP servers y modelos IA como adaptadores secundarios que se conectan a través de puertos bien definidos.

El Patrón Hexagonal para Agentes IA

terminal
                    ┌─────────────────────────────────┐
  FastAPI ──────────►│  Puerto (Entrada)               │
  (REST Endpoint)   │  AgentOrchestratorPort          │
                    │                                 │
                    │   NÚCLEO DE APLICACIÓN          │
                    │   (Python Puro)                 │
                    │   - Use Cases                   │
                    │   - Domain Entities             │
                    │   - Domain Events              │
                    │                                 │
                    │   Puerto (Salida) ────────► pgvector
                    │   VectorStorePort              │
                    │                                 │
                    │   Puerto (Salida) ────────► MCP Server
                    │   MCPServerPort                 │
                    └─────────────────────────────────┘
Tip

La regla de oro: tu núcleo de aplicación nunca debe saber nada sobre MCP, Anthropic, o pgvector. Esos son detalles de infraestructura que viven exclusivamente en los adaptadores.

Domain Events como Contrato

Los agentes IA generan eventos que deben ser parte del dominio de negocio:

python
# core/domain/events.py
from dataclasses import dataclass
from datetime import datetime
from typing import Any
from enum import Enum

class EventType(str, Enum):
    AGENT_QUERY_INITIATED = "agent_query_initiated"
    AGENT_TOOL_INVOKED = "agent_tool_invoked"
    AGENT_RESPONSE_GENERATED = "agent_response_generated"
    VECTOR_STORE_UPDATED = "vector_store_updated"

@dataclass
class DomainEvent:
    event_type: EventType
    aggregate_id: str
    payload: dict[str, Any]
    timestamp: datetime
    correlation_id: str
    user_id: str

Esta abstracción permite que cualquier implementación de storage (PostgreSQL, Redis, Kafka) capture los eventos sin afectar el núcleo.

Implementation

Paso 1: Definir el Puerto MCP

El puerto define qué herramientas necesita el núcleo, no cómo se implementan:

python
# core/ports/mcp_server_port.py
from abc import ABC, abstractmethod
from typing import Any

class MCPTool(ABC):
    @abstractmethod
    def name(self) -> str:
        pass

    @abstractmethod
    def description(self) -> str:
        pass

    @abstractmethod
    async def invoke(self, params: dict[str, Any]) -> Any:
        pass

class MCPServerPort(ABC):
    """Puerto de salida para MCP tools"""

    @abstractmethod
    async def list_tools(self) -> list[MCPTool]:
        pass

    @abstractmethod
    async def invoke_tool(self, tool_name: str, params: dict[str, Any]) -> Any:
        pass

    @abstractmethod
    async def health_check(self) -> bool:
        pass

Paso 2: Implementar el Adaptador MCP

El adaptador implementa el puerto usando el protocolo MCP real:

python
# infrastructure/adapters/mcp_adapter.py
from anthropic import Anthropic
from mcp import ClientSession
import httpx
from core.ports.mcp_server_port import MCPServerPort, MCPTool

class MCPToolAdapter(MCPTool):
    def __init__(self, tool_def: dict, session: ClientSession):
        self._tool_def = tool_def
        self._session = session

    def name(self) -> str:
        return self._tool_def["name"]

    def description(self) -> str:
        return self._tool_def["description"]

    async def invoke(self, params: dict[str, Any]) -> Any:
        return await self._session.call_tool(self.name(), params)

class RealMCPServerAdapter(MCPServerPort):
    def __init__(self, mcp_server_url: str):
        self._url = mcp_server_url
        self._session = None
        self._client = httpx.AsyncClient(timeout=30.0)

    async def _ensure_session(self):
        if self._session is None:
            self._session = ClientSession(self._url)
            await self._session.connect()

    async def list_tools(self) -> list[MCPTool]:
        await self._ensure_session()
        tools = await self._session.list_tools()
        return [MCPToolAdapter(tool, self._session) for tool in tools]

    async def invoke_tool(self, tool_name: str, params: dict[str, Any]) -> Any:
        await self._ensure_session()
        return await self._session.call_tool(tool_name, params)

    async def health_check(self) -> bool:
        try:
            await self._ensure_session()
            return await self._session.ping() == "pong"
        except Exception:
            return False

Paso 3: Puerto Vector Store

Para pgvector, definimos un puerto genérico de búsqueda vectorial:

python
# core/ports/vector_store_port.py
from abc import ABC, abstractmethod
from typing import Optional
from dataclasses import dataclass

@dataclass
class VectorDocument:
    id: str
    vector: list[float]
    metadata: dict[str, any]
    content: str

@dataclass
class SearchResult:
    document: VectorDocument
    similarity: float

class VectorStorePort(ABC):
    """Puerto de salida para operaciones de búsqueda vectorial"""

    @abstractmethod
    async def insert(self, doc: VectorDocument) -> None:
        pass

    @abstractmethod
    async def search(
        self,
        vector: list[float],
        limit: int = 5,
        filters: Optional[dict[str, any]] = None
    ) -> list[SearchResult]:
        pass

    @abstractmethod
    async def delete(self, doc_id: str) -> None:
        pass

Paso 4: Adaptador pgvector

El adaptador implementa el puerto usando pgvector:

python
# infrastructure/adapters/pgvector_adapter.py
import asyncpg
from core.ports.vector_store_port import (
    VectorStorePort,
    VectorDocument,
    SearchResult
)

class PgVectorAdapter(VectorStorePort):
    def __init__(self, db_url: str, embedding_dim: int = 1536):
        self._db_url = db_url
        self._embedding_dim = embedding_dim
        self._pool = None

    async def _ensure_pool(self):
        if self._pool is None:
            self._pool = await asyncpg.create_pool(self._db_url)

    async def insert(self, doc: VectorDocument) -> None:
        await self._ensure_pool()
        async with self._pool.acquire() as conn:
            await conn.execute(
                """
                INSERT INTO documents (id, vector, metadata, content)
                VALUES ($1, $2::vector, $3::jsonb, $4)
                ON CONFLICT (id) DO UPDATE
                SET vector = $2::vector, metadata = $3::jsonb, content = $4
                """,
                doc.id,
                doc.vector,
                doc.metadata,
                doc.content
            )

    async def search(
        self,
        vector: list[float],
        limit: int = 5,
        filters: Optional[dict[str, any]] = None
    ) -> list[SearchResult]:
        await self._ensure_pool()
        async with self._pool.acquire() as conn:
            if filters:
                query = """
                    SELECT id, vector, metadata, content,
                           1 - (vector <=> $1::vector) as similarity
                    FROM documents
                    WHERE metadata @> $2::jsonb
                    ORDER BY vector <=> $1::vector
                    LIMIT $3
                """
                rows = await conn.fetch(
                    query,
                    vector,
                    filters,
                    limit
                )
            else:
                query = """
                    SELECT id, vector, metadata, content,
                           1 - (vector <=> $1::vector) as similarity
                    FROM documents
                    ORDER BY vector <=> $1::vector
                    LIMIT $2
                """
                rows = await conn.fetch(query, vector, limit)

            return [
                SearchResult(
                    document=VectorDocument(
                        id=row["id"],
                        vector=list(row["vector"]),
                        metadata=dict(row["metadata"]),
                        content=row["content"]
                    ),
                    similarity=row["similarity"]
                )
                for row in rows
            ]

Paso 5: Caso de Uso Principal

El caso de uso coordina MCP y vector search sin saber detalles de implementación:

python
# core/use_cases/process_agent_query.py
from anthropic import Anthropic
from core.ports.mcp_server_port import MCPServerPort
from core.ports.vector_store_port import VectorStorePort
from core.domain.events import DomainEvent, EventType
from datetime import datetime

class ProcessAgentQueryUseCase:
    def __init__(
        self,
        mcp_server: MCPServerPort,
        vector_store: VectorStorePort,
        event_publisher: EventPublisherPort
    ):
        self._mcp = mcp_server
        self._vector_store = vector_store
        self._events = event_publisher
        self._llm = Anthropic()

    async def execute(
        self,
        query: str,
        user_id: str,
        correlation_id: str
    ) -> str:
        # Publicar evento de inicio
        await self._events.publish(
            DomainEvent(
                event_type=EventType.AGENT_QUERY_INITIATED,
                aggregate_id=correlation_id,
                payload={"query": query, "user_id": user_id},
                timestamp=datetime.utcnow(),
                correlation_id=correlation_id,
                user_id=user_id
            )
        )

        # Generar embedding para búsqueda vectorial
        embedding = await self._generate_embedding(query)
        results = await self._vector_store.search(embedding, limit=3)

        # Obtener herramientas MCP disponibles
        tools = await self._mcp.list_tools()
        tools_list = [{"name": t.name(), "description": t.description()} for t in tools]

        # Construir contexto
        context = "\n".join([r.document.content for r in results])

        # Invocar LLM con contexto y herramientas
        response = self._llm.messages.create(
            model="claude-3.5-sonnet",
            max_tokens=1024,
            messages=[{
                "role": "user",
                "content": f"""
Contexto relevante:
{context}

Pregunta: {query}
"""
            }],
            tools=tools_list
        )

        # Procesar tool calls si existen
        if hasattr(response.content[0], 'text'):
            final_response = response.content[0].text
        else:
            # Procesar tool calls
            for block in response.content:
                if hasattr(block, 'tool_use'):
                    tool_result = await self._mcp.invoke_tool(
                        block.tool_use.name,
                        block.tool_use.input
                    )
                    final_response = f"Tool {block.tool_use.name} ejecutado: {tool_result}"

        # Publicar evento de respuesta
        await self._events.publish(
            DomainEvent(
                event_type=EventType.AGENT_RESPONSE_GENERATED,
                aggregate_id=correlation_id,
                payload={"response": final_response},
                timestamp=datetime.utcnow(),
                correlation_id=correlation_id,
                user_id=user_id
            )
        )

        return final_response

Paso 6: Handler FastAPI

El adaptador de entrada conecta HTTP con el caso de uso:

python
# infrastructure/web/handlers/agent_handler.py
from fastapi import APIRouter, Depends, HTTPException
from core.use_cases.process_agent_query import ProcessAgentQueryUseCase
from pydantic import BaseModel

router = APIRouter(prefix="/api/v1")

class AgentQueryRequest(BaseModel):
    query: str

class AgentQueryResponse(BaseModel):
    response: str
    query_id: str

@router.post("/agent/query", response_model=AgentQueryResponse)
async def process_agent_query(
    request: AgentQueryRequest,
    use_case: ProcessAgentQueryUseCase = Depends(get_agent_use_case)
):
    try:
        correlation_id = str(uuid.uuid4())
        user_id = "demo-user"  # De JWT en producción

        response = await use_case.execute(
            query=request.query,
            user_id=user_id,
            correlation_id=correlation_id
        )

        return AgentQueryResponse(
            response=response,
            query_id=correlation_id
        )
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))

Paso 7: Inyección de Dependencias

python
# infrastructure/dependencies.py
from core.ports.mcp_server_port import MCPServerPort
from core.ports.vector_store_port import VectorStorePort
from core.use_cases.process_agent_query import ProcessAgentQueryUseCase
from infrastructure.adapters.mcp_adapter import RealMCPServerAdapter
from infrastructure.adapters.pgvector_adapter import PgVectorAdapter
from infrastructure.adapters.event_publisher import PostgresEventPublisher

# Instancias singleton (en producción, usa un DI framework)
_mcp_adapter: MCPServerPort = None
_vector_adapter: VectorStorePort = None
_event_publisher: EventPublisherPort = None

def get_mcp_adapter() -> MCPServerPort:
    global _mcp_adapter
    if _mcp_adapter is None:
        _mcp_adapter = RealMCPServerAdapter(
            mcp_server_url=os.getenv("MCP_SERVER_URL", "http://localhost:3000")
        )
    return _mcp_adapter

def get_vector_adapter() -> VectorStorePort:
    global _vector_adapter
    if _vector_adapter is None:
        _vector_adapter = PgVectorAdapter(
            db_url=os.getenv("DATABASE_URL")
        )
    return _vector_adapter

def get_event_publisher() -> EventPublisherPort:
    global _event_publisher
    if _event_publisher is None:
        _event_publisher = PostgresEventPublisher(
            db_url=os.getenv("DATABASE_URL")
        )
    return _event_publisher

def get_agent_use_case() -> ProcessAgentQueryUseCase:
    return ProcessAgentQueryUseCase(
        mcp_server=get_mcp_adapter(),
        vector_store=get_vector_adapter(),
        event_publisher=get_event_publisher()
    )
⚠️Atención

En producción, considera usar un framework de DI como dependency-injector o injector para manejar el ciclo de vida de estas dependencias de forma más robusta.

Lessons Learned

Lección 1: Los Tests Unitarios Deben Ser Sinceros

El test malo requiere MCP server real y base de datos. Es lento, frágil y no es realmente un test unitario. El test bueno verifica la lógica del caso de uso sin dependencias externas.

Lección 2: Circuit Breakers son Obligatorios

Los MCP tools pueden fallar. Sin circuit breakers, un tool que tira timeout puede colgar todo el sistema:

python
# infrastructure/adapters/resilient_mcp_adapter.py
from circuitbreaker import circuit
from core.ports.mcp_server_port import MCPServerPort
import time
import logging

logger = logging.getLogger(__name__)

class ResilientMCPServerAdapter(MCPServerPort):
    def __init__(self, inner_adapter: MCPServerPort, failure_threshold: int = 5, timeout: int = 30):
        self._inner = inner_adapter
        self._threshold = failure_threshold
        self._timeout = timeout

    @circuit(failure_threshold=5, recovery_timeout=60)
    async def list_tools(self):
        try:
            return await asyncio.wait_for(
                self._inner.list_tools(),
                timeout=self._timeout
            )
        except asyncio.TimeoutError:
            logger.error("MCP server timeout en list_tools")
            raise
        except Exception as e:
            logger.error(f"Error en MCP list_tools: {e}")
            raise

    @circuit(failure_threshold=5, recovery_timeout=60)
    async def invoke_tool(self, tool_name: str, params: dict):
        try:
            return await asyncio.wait_for(
                self._inner.invoke_tool(tool_name, params),
                timeout=self._timeout
            )
        except asyncio.TimeoutError:
            logger.error(f"MCP tool {tool_name} timeout")
            raise
        except Exception as e:
            logger.error(f"Error en MCP tool {tool_name}: {e}")
            raise

    async def health_check(self) -> bool:
        return await self._inner.health_check()

Lección 3: Domain Events Permiten Debugging

Sin eventos de dominio, cuando un agente produce resultados incorrectos, no tienes traceability. Con el patrón de events:

python
# infrastructure/repositories/event_repository.py
import asyncpg
from core.domain.events import DomainEvent

class PostgresEventRepository:
    def __init__(self, db_url: str):
        self._pool = None
        self._db_url = db_url

    async def save(self, event: DomainEvent):
        await self._ensure_pool()
        async with self._pool.acquire() as conn:
            await conn.execute(
                """
                INSERT INTO domain_events
                (event_type, aggregate_id, payload, timestamp, correlation_id, user_id)
                VALUES ($1, $2, $3::jsonb, $4, $5, $6)
                """,
                event.event_type.value,
                event.aggregate_id,
                event.payload,
                event.timestamp,
                event.correlation_id,
                event.user_id
            )

    async def get_by_correlation(self, correlation_id: str) -> list[DomainEvent]:
        await self._ensure_pool()
        async with self._pool.acquire() as conn:
            rows = await conn.fetch(
                "SELECT * FROM domain_events WHERE correlation_id = $1 ORDER BY timestamp",
                correlation_id
            )
            return [DomainEvent(**dict(row)) for row in rows]

Esto permite debugging completo de cualquier conversación de agente.

Excelente

Este patrón de domain events es idéntico al que discutimos en MCP en Producción: Patrones de Integración Real, donde vimos cómo la observabilidad es crítica para sistemas de agentes.

Lección 4: pgvector como Servicio, No como Infraestructura

El error común es mezclar código SQL directamente en el caso de uso. El puerto VectorStorePort permite cambiar de pgvector a Pinecone o Weaviate sin tocar el núcleo:

Lección 5: Startup Time vs Test Speed

Un sistema hexagonal bien diseñado tiene tests que corren en milisegundos:

python
# tests/test_use_cases.py
import pytest
from unittest.mock import AsyncMock
from core.use_cases.process_agent_query import ProcessAgentQueryUseCase

@pytest.fixture
def mock_mcp():
    return AsyncMock()

@pytest.fixture
def mock_vector():
    return AsyncMock()

@pytest.fixture
def mock_events():
    return AsyncMock()

def test_use_case_starts_fast(mock_mcp, mock_vector, mock_events):
    """Este test debería correr en < 100ms"""
    use_case = ProcessAgentQueryUseCase(mock_mcp, mock_vector, mock_events)
    assert use_case is not None

# Resultado real: ~45ms

Comparen esto con un test que inicia PostgreSQL y un MCP server: ~5-10 segundos. La diferencia es exponencial en un suite de cientos de tests.

1

Implementa los puertos primero

Define las interfaces abstractas (MCPServerPort, VectorStorePort) antes de escribir cualquier implementación concreta. Esto forzarte a pensar en qué necesitas, no en cómo lo implementarás.

2

Escribe tests con mocks

Crea implementaciones mock de los puertos y escribe tests del caso de uso. Esto asegura que la lógica es correcta antes de conectar MCP real o pgvector.

3

Implementa adaptadores con resiliencia

Los adaptadores MCP y pgvector deben tener circuit breakers, retries con backoff, y timeouts desde el día uno. No agregues resiliencia después—siempre es más costoso.

4

Integra observabilidad desde el inicio

Cada llamada a MCP y cada query a pgvector debe ser logged y traced. Usa OpenTelemetry como vimos en Observabilidad en Sistemas IA.

5

Documenta los domain events

Crea un catálogo de eventos (agent_query_initiated, agent_tool_invoked, etc.) y manténlo versionado. Esto permite evolución sin romper consumers.

Conclusion

La Arquitectura Hexagonal no es un patrón académico—es una necesidad práctica cuando trabajas con agentes IA en producción. La combinación de MCP + pgvector + agentes IA crea un sistema inherentemente complejo, y la arquitectura hexagonal provee la disciplina para manejar esa complejidad.

Los beneficios clave que hemos observado en producción:

  1. Testability: Tests unitarios que corren en milisegundos, no segundos
  2. Flexibilidad: Cambiar de un MCP server a otro sin tocar la lógica de negocio
  3. Observabilidad: Domain events permiten debugging completo de conversaciones de agentes
  4. Escalabilidad: Adaptadores separados permiten optimizar MCP y pgvector independientemente
  5. Mantenibilidad: Regla de dependencia clara hace que el código sea comprensible
⚠️Atención

La arquitectura hexagonal tiene un costo: mayor ceremonia inicial y más archivos. Pero para sistemas de agentes IA que valen la pena, ese costo se paga a sí mismo rápidamente en reducción de bugs y velocidad de iteration.

La pregunta que debes hacerte no es "¿Debo usar arquitectura hexagonal?" sino "¿Este sistema de agentes IA va a sobrevivir más de 6 meses en producción?" Si la respuesta es sí, los puertos y adaptadores son una inversión, no un costo.

Para profundizar en el stack de vectores, revisa RAG con pgvector: Arquitectura a Escala, donde cubrimos patrones de sharding y partitioning para vector stores a escala.


Anterior en esta serie: Arquitectura Hexagonal en Python: Cuándo Usarla y Cuándo Es Excesiva

Referencias rápidas

Vista general

Recursos externos

Incluye recursos adicionales en el frontmatter para que aparezcan aquí.

Más en esta serie

Serie: Arquitectura de Software Avanzada

// newsletter

¿Te sirvió este artículo?

Recibe los siguientes en tu inbox. Sin spam, cancela cuando quieras.

Discusión

Escrito por Jorge Ochoa. ¿Encontraste un error?

Abrir en GitHub