PydanticRPC

Una biblioteca de Python para construir servicios gRPC/ConnectRPC con modelos Pydantic, que incluye generación automática de protobuf y exposición de herramientas para asistentes de IA.

Documentación

🚀 PydanticRPC

PydanticRPC es una biblioteca de Python que te permite exponer rápidamente modelos Pydantic a través de servicios gRPC/Connect RPC sin escribir ningún archivo protobuf. En su lugar, genera automáticamente archivos protobuf sobre la marcha a partir de las firmas de métodos de tus objetos Python y las firmas de tipo de tus modelos Pydantic.

A continuación se muestra un ejemplo de un servicio gRPC simple que expone un agente PydanticAI:

import asyncio

from openai import AsyncOpenAI
from pydantic_ai import Agent
from pydantic_ai.models.openai import OpenAIModel
from pydantic_rpc import AsyncIOServer, Message


# `Message` is just an alias for Pydantic's `BaseModel` class.
class CityLocation(Message):
    city: str
    country: str


class Olympics(Message):
    year: int

    def prompt(self):
        return f"Where were the Olympics held in {self.year}?"


class OlympicsLocationAgent:
    def __init__(self):
        client = AsyncOpenAI(
            base_url="http://localhost:11434/v1",
            api_key="ollama_api_key",
        )
        ollama_model = OpenAIModel(
            model_name="llama3.2",
            openai_client=client,
        )
        self._agent = Agent(ollama_model)

    async def ask(self, req: Olympics) -> CityLocation:
        result = await self._agent.run(req.prompt())
        return result.data


if __name__ == "__main__":
    # New enhanced initialization API (optional - backward compatible)
    s = AsyncIOServer(service=OlympicsLocationAgent(), port=50051)
    loop = asyncio.get_event_loop()
    loop.run_until_complete(s.run())

Y aquí hay un ejemplo de un servicio Connect RPC simple que expone el mismo agente como una aplicación ASGI:

import asyncio

from openai import AsyncOpenAI
from pydantic_ai import Agent
from pydantic_ai.models.openai import OpenAIModel
from pydantic_rpc import ASGIApp, Message


class CityLocation(Message):
    city: str
    country: str


class Olympics(Message):
    year: int

    def prompt(self):
        return f"Where were the Olympics held in {self.year}?"


class OlympicsLocationAgent:
    def __init__(self):
        client = AsyncOpenAI(
            base_url="http://localhost:11434/v1",
            api_key="ollama_api_key",
        )
        ollama_model = OpenAIModel(
            model_name="llama3.2",
            openai_client=client,
        )
        self._agent = Agent(ollama_model, result_type=CityLocation)

    async def ask(self, req: Olympics) -> CityLocation:
        result = await self._agent.run(req.prompt())
        return result.data

# New enhanced initialization API (optional - backward compatible)
app = ASGIApp(service=OlympicsLocationAgent())

💡 Características Clave

  • 🔄 Generación Automática de Protobuf: Crea automáticamente archivos protobuf que coinciden con las firmas de métodos de tus objetos Python.
  • ⚙️ Generación Dinámica de Código: Genera stubs de servidor y cliente usando grpcio-tools.
  • ✅ Integración con Pydantic: Usa pydantic para una validación de tipos y serialización robusta.
  • 📄 Exportación de Archivos Protobuf: Exporta los archivos protobuf generados para usarlos en otros lenguajes.
  • Para gRPC:
    • 💚 Verificación de Salud: Soporte integrado para verificaciones de salud gRPC usando grpc_health.v1.
    • 🔎 Reflexión de Servidor: Soporte integrado para reflexión de servidor gRPC.
    • ⚡ Soporte Asíncrono: Crea fácilmente servicios gRPC asíncronos con AsyncIOServer.
  • Para Connect-RPC:
    • 🌐 Soporte Completo de Protocolo: Soporte nativo de Connect-RPC a través de connect-python
    • 🔄 Todos los Patrones de Streaming: Unario, streaming de servidor, streaming de cliente y streaming bidireccional
    • 🌐 Aplicaciones WSGI/ASGI: Se ejecuta como aplicaciones WSGI o ASGI estándar para una implementación fácil
  • 🛠️ Archivos Protobuf y Código Pre-generados: Pre-genera archivos proto y código correspondiente a través de la CLI. Al configurar la variable de entorno (PYDANTIC_RPC_SKIP_GENERATION), puedes omitir la generación en tiempo de ejecución.
  • 🤖 Soporte MCP (Protocolo de Contexto de Modelo): Expón tus servicios como herramientas para asistentes de IA usando el SDK oficial de MCP, con soporte para transportes stdio y HTTP/SSE.

⚠️ Notas Importantes para Connect-RPC

Al usar Connect-RPC con ASGIApp:

  • Formato de Ruta de Endpoint: Los endpoints de Connect-RPC usan nombres de método en CamelCase en la ruta: /<package>.<service>/<Method> (por ejemplo, /chat.v1.ChatService/SendMessage)
  • Tipo de Contenido: Establece Content-Type: application/json o application/connect+json para las solicitudes
  • Requisito de HTTP/2: El streaming bidireccional requiere HTTP/2. Usa Hypercorn en lugar de uvicorn para soporte HTTP/2
  • Pruebas: Usa buf curl para probar endpoints de Connect-RPC con soporte adecuado de streaming

Para ejemplos detallados e instrucciones de prueba, consulta el directorio de ejemplos.

📦 Instalación

Instala PydanticRPC vía pip:

pip install pydantic-rpc

Para soporte CLI con ejecutores de servidor integrados:

pip install pydantic-rpc-cli  # Includes hypercorn and gunicorn

🆕 Características Mejoradas (v0.10.0+)

Nota: Todas las nuevas características son totalmente compatibles con versiones anteriores. El código existente sigue funcionando sin modificaciones.

API de Inicialización Mejorada

Todas las clases de servidor ahora admiten inicialización opcional con servicios:

# Traditional API (still works)
server = AsyncIOServer()
server.set_port(50051)
await server.run(MyService())

# New enhanced API (optional)
server = AsyncIOServer(
    service=MyService(),
    port=50051,
    package_name="my.package"
)
await server.run()

# Same for ASGI/WSGI apps
app = ASGIApp(service=MyService(), package_name="my.package")

Manejo de Errores con Decoradores

Asigna automáticamente excepciones a códigos de estado gRPC/Connect:

from pydantic_rpc import error_handler
import grpc

class MyService:
    @error_handler(ValidationError, status_code=grpc.StatusCode.INVALID_ARGUMENT)
    @error_handler(KeyError, status_code=grpc.StatusCode.NOT_FOUND)
    async def get_user(self, request: GetUserRequest) -> User:
        # Exceptions are automatically converted to proper status codes
        if request.id not in users_db:
            raise KeyError(f"User {request.id} not found")
        return users_db[request.id]

🚀 Primeros Pasos

PydanticRPC admite dos protocolos principales:

  • gRPC: Servicios gRPC tradicionales con Server y AsyncIOServer
  • Connect-RPC: RPC moderno basado en HTTP con ASGIApp y WSGIApp

🔧 Ejemplo de Servicio gRPC Síncrono

from pydantic_rpc import Server, Message

class HelloRequest(Message):
    name: str

class HelloReply(Message):
    message: str

class Greeter:
    # Define methods that accepts a request and returns a response.
    def say_hello(self, request: HelloRequest) -> HelloReply:
        return HelloReply(message=f"Hello, {request.name}!")

if __name__ == "__main__":
    server = Server()
    server.run(Greeter())

⚙️ Ejemplo de Servicio gRPC Asíncrono

import asyncio

from pydantic_rpc import AsyncIOServer, Message


class HelloRequest(Message):
    name: str


class HelloReply(Message):
    message: str


class Greeter:
    async def say_hello(self, request: HelloRequest) -> HelloReply:
        return HelloReply(message=f"Hello, {request.name}!")


async def main():
    # You can specify a custom port (default is 50051)
    server = AsyncIOServer(port=50052)
    await server.run(Greeter())


if __name__ == "__main__":
    asyncio.run(main())

El AsyncIOServer maneja automáticamente el apagado elegante en señales SIGTERM y SIGINT.

🌐 Ejemplo de Aplicación ASGI Connect-RPC

from pydantic_rpc import ASGIApp, Message

class HelloRequest(Message):
    name: str

class HelloReply(Message):
    message: str

class Greeter:
    async def say_hello(self, request: HelloRequest) -> HelloReply:
        return HelloReply(message=f"Hello, {request.name}!")

app = ASGIApp()
app.mount(Greeter())

# Run with uvicorn:
# uvicorn script:app --host 0.0.0.0 --port 8000

🌐 Ejemplo de Aplicación WSGI Connect-RPC

from pydantic_rpc import WSGIApp, Message

class HelloRequest(Message):
    name: str

class HelloReply(Message):
    message: str

class Greeter:
    def say_hello(self, request: HelloRequest) -> HelloReply:
        return HelloReply(message=f"Hello, {request.name}!")

app = WSGIApp()
app.mount(Greeter())

# Run with gunicorn:
# gunicorn script:app

🏆 Ejemplo de Connect-RPC con Streaming

PydanticRPC proporciona soporte nativo de Connect-RPC a través de connect-python, incluyendo capacidades completas de streaming y convenciones de nomenclatura PEP 8. Consulta nuestros ejemplos ASGI:

# Run with uvicorn
uv run uvicorn greeting_asgi:app --port 3000

# Or run streaming example
uv run python examples/streaming_connect_python.py

Esto lanzará una aplicación ASGI basada en connect-python que usa los mismos modelos Pydantic para servir solicitudes Connect-RPC.

Soporte de Streaming con connect-python

connect-python proporciona soporte completo para RPC de streaming con nomenclatura automática PEP 8 (snake_case):

from typing import AsyncIterator
from pydantic_rpc import ASGIApp, Message

class StreamRequest(Message):
    text: str
    count: int

class StreamResponse(Message):
    text: str
    index: int

class StreamingService:
    # Server streaming
    async def server_stream(self, request: StreamRequest) -> AsyncIterator[StreamResponse]:
        for i in range(request.count):
            yield StreamResponse(text=f"{request.text}_{i}", index=i)
    
    # Client streaming
    async def client_stream(self, requests: AsyncIterator[StreamRequest]) -> StreamResponse:
        texts = []
        async for req in requests:
            texts.append(req.text)
        return StreamResponse(text=" ".join(texts), index=len(texts))
    
    # Bidirectional streaming
    async def bidi_stream(
        self, requests: AsyncIterator[StreamRequest]
    ) -> AsyncIterator[StreamResponse]:
        idx = 0
        async for req in requests:
            yield StreamResponse(text=f"Echo: {req.text}", index=idx)
            idx += 1

app = ASGIApp()
app.mount(StreamingService())

[!NOTE] Por favor, instala protoc-gen-connect-python para ejecutar el ejemplo de connect-python.

♻️ Omitir la Generación de Protobuf

Por defecto, PydanticRPC genera archivos .proto y código en tiempo de ejecución. Si deseas omitir el paso de generación de código (por ejemplo, en entorno de producción), establece la variable de entorno a continuación:

export PYDANTIC_RPC_SKIP_GENERATION=true

Cuando esta variable se establece en "true", PydanticRPC cargará módulos pre-generados existentes en lugar de generarlos sobre la marcha.

🪧 Configuración del Directorio de Generación de Protobuf y Connect RPC/gRPC

Por defecto, tus archivos se generarán en el directorio de trabajo actual desde donde ejecutaste el código, pero puedes establecer un directorio específico personalizado configurando la variable de entorno a continuación:

export PYDANTIC_RPC_PROTO_PATH=/your/path

⚠️ Campos Reservados

También puedes establecer una variable de entorno para reservar un número determinado de campos para la generación de proto, para compatibilidad hacia atrás y hacia adelante.

export PYDANTIC_RPC_RESERVED_FIELDS=1

💎 Características Avanzadas

🌊 Streaming de Respuestas (gRPC)

PydanticRPC admite respuestas de streaming tanto para servicios gRPC como Connect-RPC. Si el tipo de retorno de un método de clase de servicio es typing.AsyncIterator[T], el método se considera un método de streaming.

Consulta el código de ejemplo a continuación:

import asyncio
from typing import Annotated, AsyncIterator

from openai import AsyncOpenAI
from pydantic import Field
from pydantic_ai import Agent
from pydantic_ai.models.openai import OpenAIModel
from pydantic_rpc import AsyncIOServer, Message


# `Message` is just a pydantic BaseModel alias
class CityLocation(Message):
    city: Annotated[str, Field(description="The city where the Olympics were held")]
    country: Annotated[
        str, Field(description="The country where the Olympics were held")
    ]


class OlympicsQuery(Message):
    year: Annotated[int, Field(description="The year of the Olympics", ge=1896)]

    def prompt(self):
        return f"Where were the Olympics held in {self.year}?"


class OlympicsDurationQuery(Message):
    start: Annotated[int, Field(description="The start year of the Olympics", ge=1896)]
    end: Annotated[int, Field(description="The end year of the Olympics", ge=1896)]

    def prompt(self):
        return f"From {self.start} to {self.end}, how many Olympics were held? Please provide the list of countries and cities."


class StreamingResult(Message):
    answer: Annotated[str, Field(description="The answer to the query")]


class OlympicsAgent:
    def __init__(self):
        client = AsyncOpenAI(
            base_url='http://localhost:11434/v1',
            api_key='ollama_api_key',
        )
        ollama_model = OpenAIModel(
            model_name='llama3.2',
            openai_client=client,
        )
        self._agent = Agent(ollama_model)

    async def ask(self, req: OlympicsQuery) -> CityLocation:
        result = await self._agent.run(req.prompt(), result_type=CityLocation)
        return result.data

    async def ask_stream(
        self, req: OlympicsDurationQuery
    ) -> AsyncIterator[StreamingResult]:
        async with self._agent.run_stream(req.prompt(), result_type=str) as result:
            async for data in result.stream_text(delta=True):
                yield StreamingResult(answer=data)


if __name__ == "__main__":
    s = AsyncIOServer()
    loop = asyncio.get_event_loop()
    loop.run_until_complete(s.run(OlympicsAgent()))

En el ejemplo anterior, el método ask_stream devuelve un objeto AsyncIterator[StreamingResult], que se considera un método de streaming. La clase StreamingResult es un modelo Pydantic que define el tipo de respuesta del método de streaming. Puedes usar cualquier modelo Pydantic como tipo de respuesta.

Ahora, puedes llamar al método ask_stream del servidor descrito anteriormente usando tu herramienta de cliente gRPC preferida. El ejemplo a continuación usa buf curl.

% buf curl --data '{"start": 1980, "end": 2024}' -v http://localhost:50051/olympicsagent.v1.OlympicsAgent/AskStream --protocol grpc --http2-prior-knowledge 

buf: * Using server reflection to resolve "olympicsagent.v1.OlympicsAgent"
buf: * Dialing (tcp) localhost:50051...
buf: * Connected to [::1]:50051
buf: > (#1) POST /grpc.reflection.v1.ServerReflection/ServerReflectionInfo
buf: > (#1) Accept-Encoding: identity
buf: > (#1) Content-Type: application/grpc+proto
buf: > (#1) Grpc-Accept-Encoding: gzip
buf: > (#1) Grpc-Timeout: 119997m
buf: > (#1) Te: trailers
buf: > (#1) User-Agent: grpc-go-connect/1.12.0 (go1.21.4) buf/1.28.1
buf: > (#1)
buf: } (#1) [5 bytes data]
buf: } (#1) [32 bytes data]
buf: < (#1) HTTP/2.0 200 OK
buf: < (#1) Content-Type: application/grpc
buf: < (#1) Grpc-Message: Method not found!
buf: < (#1) Grpc-Status: 12
buf: < (#1)
buf: * (#1) Call complete
buf: > (#2) POST /grpc.reflection.v1alpha.ServerReflection/ServerReflectionInfo
buf: > (#2) Accept-Encoding: identity
buf: > (#2) Content-Type: application/grpc+proto
buf: > (#2) Grpc-Accept-Encoding: gzip
buf: > (#2) Grpc-Timeout: 119967m
buf: > (#2) Te: trailers
buf: > (#2) User-Agent: grpc-go-connect/1.12.0 (go1.21.4) buf/1.28.1
buf: > (#2)
buf: } (#2) [5 bytes data]
buf: } (#2) [32 bytes data]
buf: < (#2) HTTP/2.0 200 OK
buf: < (#2) Content-Type: application/grpc
buf: < (#2) Grpc-Accept-Encoding: identity, deflate, gzip
buf: < (#2)
buf: { (#2) [5 bytes data]
buf: { (#2) [434 bytes data]
buf: * Server reflection has resolved file "olympicsagent.proto"
buf: * Invoking RPC olympicsagent.v1.OlympicsAgent.AskStream
buf: > (#3) POST /olympicsagent.v1.OlympicsAgent/AskStream
buf: > (#3) Accept-Encoding: identity
buf: > (#3) Content-Type: application/grpc+proto
buf: > (#3) Grpc-Accept-Encoding: gzip
buf: > (#3) Grpc-Timeout: 119947m
buf: > (#3) Te: trailers
buf: > (#3) User-Agent: grpc-go-connect/1.12.0 (go1.21.4) buf/1.28.1
buf: > (#3)
buf: } (#3) [5 bytes data]
buf: } (#3) [6 bytes data]
buf: * (#3) Finished upload
buf: < (#3) HTTP/2.0 200 OK
buf: < (#3) Content-Type: application/grpc
buf: < (#3) Grpc-Accept-Encoding: identity, deflate, gzip
buf: < (#3)
buf: { (#3) [5 bytes data]
buf: { (#3) [25 bytes data]
{
 "answer": "Here's a list of Summer"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [31 bytes data]
{
  "answer": " and Winter Olympics from 198"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [29 bytes data]
{
  "answer": "0 to 2024:\n\nSummer Olympics"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [20 bytes data]
{
  "answer": ":\n1. 1980 - Moscow"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [20 bytes data]
{
  "answer": ", Soviet Union\n2. "
}
buf: { (#3) [5 bytes data]
buf: { (#3) [32 bytes data]
{
  "answer": "1984 - Los Angeles, California"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [15 bytes data]
{
  "answer": ", USA\n3. 1988"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [26 bytes data]
{
  "answer": " - Seoul, South Korea\n4."
}
buf: { (#3) [5 bytes data]
buf: { (#3) [27 bytes data]
{
  "answer": " 1992 - Barcelona, Spain\n"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [20 bytes data]
{
  "answer": "5. 1996 - Atlanta,"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [22 bytes data]
{
  "answer": " Georgia, USA\n6. 200"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [26 bytes data]
{
  "answer": "0 - Sydney, Australia\n7."
}
buf: { (#3) [5 bytes data]
buf: { (#3) [25 bytes data]
{
  "answer": " 2004 - Athens, Greece\n"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [20 bytes data]
{
  "answer": "8. 2008 - Beijing,"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [18 bytes data]
{
  "answer": " China\n9. 2012 -"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [29 bytes data]
{
  "answer": " London, United Kingdom\n10."
}
buf: { (#3) [5 bytes data]
buf: { (#3) [24 bytes data]
{
  "answer": " 2016 - Rio de Janeiro"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [18 bytes data]
{
  "answer": ", Brazil\n11. 202"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [24 bytes data]
{
  "answer": "0 - Tokyo, Japan (held"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [21 bytes data]
{
  "answer": " in 2021 due to the"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [26 bytes data]
{
  "answer": " COVID-19 pandemic)\n12. "
}
buf: { (#3) [5 bytes data]
buf: { (#3) [28 bytes data]
{
  "answer": "2024 - Paris, France\n\nNote"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [41 bytes data]
{
  "answer": ": The Olympics were held without a host"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [26 bytes data]
{
  "answer": " city for one year (2022"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [42 bytes data]
{
  "answer": ", due to the Russian invasion of Ukraine"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [29 bytes data]
{
  "answer": ").\n\nWinter Olympics:\n1. 198"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [27 bytes data]
{
  "answer": "0 - Lake Placid, New York"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [15 bytes data]
{
  "answer": ", USA\n2. 1984"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [27 bytes data]
{
  "answer": " - Sarajevo, Yugoslavia ("
}
buf: { (#3) [5 bytes data]
buf: { (#3) [30 bytes data]
{
  "answer": "now Bosnia and Herzegovina)\n"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [20 bytes data]
{
  "answer": "3. 1988 - Calgary,"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [25 bytes data]
{
  "answer": " Alberta, Canada\n4. 199"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [26 bytes data]
{
  "answer": "2 - Albertville, France\n"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [13 bytes data]
{
  "answer": "5. 1994 - L"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [24 bytes data]
{
  "answer": "illehammer, Norway\n6. "
}
buf: { (#3) [5 bytes data]
buf: { (#3) [23 bytes data]
{
  "answer": "1998 - Nagano, Japan\n"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [16 bytes data]
{
  "answer": "7. 2002 - Salt"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [24 bytes data]
{
  "answer": " Lake City, Utah, USA\n"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [18 bytes data]
{
  "answer": "8. 2006 - Torino"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [17 bytes data]
{
  "answer": ", Italy\n9. 2010"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [40 bytes data]
{
  "answer": " - Vancouver, British Columbia, Canada"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [13 bytes data]
{
  "answer": "\n10. 2014 -"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [20 bytes data]
{
  "answer": " Sochi, Russia\n11."
}
buf: { (#3) [5 bytes data]
buf: { (#3) [16 bytes data]
{
  "answer": " 2018 - Pyeong"
}
buf: { (#3) [5 bytes data]
buf: { (#3) [24 bytes data]
{
  "answer": "chang, South Korea\n12."
}
buf: < (#3)
buf: < (#3) Grpc-Message:
buf: < (#3) Grpc-Status: 0
buf: * (#3) Call complete
buf: < (#2)
buf: < (#2) Grpc-Message:
buf: < (#2) Grpc-Status: 0
buf: * (#2) Call complete
%

🪶 Mensajes Vacíos

Los mensajes de solicitud/respuesta vacíos se asignan automáticamente a google.protobuf.Empty:

from pydantic_rpc import AsyncIOServer, Message


class EmptyRequest(Message):
    pass  # Automatically uses google.protobuf.Empty


class GreetingResponse(Message):
    message: str


class GreetingService:
    async def say_hello(self, request: EmptyRequest) -> GreetingResponse:
        return GreetingResponse(message="Hello!")
    
    async def get_default_greeting(self) -> GreetingResponse:
        # Method with no request parameter (implicitly empty)
        return GreetingResponse(message="Hello, World!")

🎨 Serialización Personalizada

Los decoradores de serialización de Pydantic son totalmente compatibles:

from typing import Any
from pydantic import field_serializer, model_serializer
from pydantic_rpc import Message


class UserMessage(Message):
    name: str
    age: int
    
    @field_serializer('name')
    def serialize_name(self, name: str) -> str:
        """Always uppercase the name when serializing."""
        return name.upper()


class ComplexMessage(Message):
    value: int
    multiplier: int
    
    @model_serializer
    def serialize_model(self) -> dict[str, Any]:
        """Custom serialization with computed fields."""
        return {
            'value': self.value,
            'multiplier': self.multiplier,
            'result': self.value * self.multiplier  # Computed field
        }

Los serializadores se aplican automáticamente al convertir entre modelos Pydantic y mensajes protobuf.

⚠️ Limitaciones y Consideraciones

1. Serializadores de Mensajes Anidados ahora compatibles (v0.8.0+)

class Address(Message):
    city: str
    
    @field_serializer("city")
    def serialize_city(self, city: str) -> str:
        return city.upper()

class User(Message):
    name: str
    address: Address  # ← Address's serializers ARE applied with DEEP strategy
    
    @field_serializer("name")
    def serialize_name(self, name: str) -> str:
        return name.upper()  # ← This IS applied

Control de Estrategia de Serialización: Puedes controlar cómo se aplican los serializadores anidados mediante la variable de entorno:

# Apply serializers at all nesting levels (default)
export PYDANTIC_RPC_SERIALIZER_STRATEGY=deep

# Apply only top-level serializers
export PYDANTIC_RPC_SERIALIZER_STRATEGY=shallow

# Disable all serializers
export PYDANTIC_RPC_SERIALIZER_STRATEGY=none

Impacto en el Rendimiento:

  • Estrategia DEEP: ~4% de sobrecarga para estructuras anidadas simples
  • Estrategia SHALLOW: ~2% de sobrecarga (solo nivel superior)
  • Estrategia NONE: Sin sobrecarga (serializadores deshabilitados)

2. Los nuevos campos agregados por serializadores se ignoran

class ComplexMessage(Message):
    value: int
    multiplier: int
    
    @model_serializer
    def serialize_model(self) -> dict[str, Any]:
        return {
            "value": self.value,
            "multiplier": self.multiplier,
            "result": self.value * self.multiplier  # ← Won't appear in protobuf
        }

Problema: El campo result no existe en la definición de Message, por lo que no está en el esquema protobuf.

3. El tipo debe permanecer consistente

class BadExample(Message):
    number: int
    
    @field_serializer("number")
    def serialize_number(self, number: int) -> str:  # ❌ int → str
        return str(number)  # This will cause issues

4. Los campos Union/Optional tienen soporte limitado

class UnionExample(Message):
    data: str | int | None  # Union type
    
    @field_serializer("data")
    def serialize_data(self, data: str | int | None) -> str | int | None:
        # Serializer may not be applied to Union types
        return data

5. Los errores fallan silenciosamente con respaldo

class RiskyMessage(Message):
    value: int
    
    @field_serializer("value")
    def serialize_value(self, value: int) -> int:
        if value == 0:
            raise ValueError("Cannot serialize zero")
        return value * 2

# If error occurs, original value is used (silent fallback)

6. Las referencias circulares se manejan con elegancia

class Node(Message):
    value: str
    child: "Node | None" = None
    
    @field_serializer("value")
    def serialize_value(self, v: str) -> str:
        return v.upper()

# Circular references are detected and prevented
node1 = Node(value="first")
node2 = Node(value="second")
node1.child = node2
node2.child = node1  # Circular reference

# When converting to protobuf:
# - Circular references are detected
# - Empty proto is returned for repeated objects
# - No infinite recursion occurs
# Note: Pydantic's model_dump() will fail on circular refs,
#       so serializers won't be applied in this case

✅ Uso Recomendado:

class GoodMessage(Message):
    # Use with primitive types
    name: str
    age: int
    
    @field_serializer("name")
    def normalize_name(self, name: str) -> str:
        return name.strip().title()  # Normalization
    
    @field_serializer("age")
    def clamp_age(self, age: int) -> int:
        return max(0, min(age, 150))  # Range limiting

Mejores Prácticas:

  • Usa serializadores principalmente para tipos primitivos (str, int, float, bool)
  • Mantén la consistencia de tipos (int → int, str → str)
  • Evita transformaciones complejas o efectos secundarios
  • Prueba los casos de error a fondo
  • Ten en cuenta que los errores fallan silenciosamente

🔒 Soporte TLS/mTLS

PydanticRPC proporciona soporte integrado para TLS (Seguridad de Capa de Transporte) y mTLS (TLS mutuo) para comunicación gRPC segura.

from pydantic_rpc import AsyncIOServer, GrpcTLSConfig, extract_peer_identity
import grpc

# Basic TLS (server authentication only)
tls_config = GrpcTLSConfig(
    cert_chain=server_cert_bytes,
    private_key=server_key_bytes,
    require_client_cert=False
)

# mTLS (mutual authentication)
tls_config = GrpcTLSConfig(
    cert_chain=server_cert_bytes,
    private_key=server_key_bytes,
    root_certs=ca_cert_bytes,  # CA to verify client certificates
    require_client_cert=True
)

# Create server with TLS
server = AsyncIOServer(tls=tls_config)

# Extract client identity in service methods
class SecureService:
    async def secure_method(self, request, context: grpc.ServicerContext):
        client_identity = extract_peer_identity(context)
        if client_identity:
            print(f"Authenticated client: {client_identity}")

Para un ejemplo completo, consulta examples/tls_server.py y examples/tls_client.py.

🔗 Múltiples Servicios con Interceptores Personalizados

PydanticRPC admite definir y ejecutar múltiples servicios gRPC en un solo servidor:

from datetime import datetime
import grpc
from grpc import ServicerContext

from pydantic_rpc import Server, Message


class FooRequest(Message):
    name: str
    age: int
    d: dict[str, str]


class FooResponse(Message):
    name: str
    age: int
    d: dict[str, str]


class BarRequest(Message):
    names: list[str]


class BarResponse(Message):
    names: list[str]


class FooService:
    def foo(self, request: FooRequest) -> FooResponse:
        return FooResponse(name=request.name, age=request.age, d=request.d)


class MyMessage(Message):
    name: str
    age: int
    o: int | datetime


class Request(Message):
    name: str
    age: int
    d: dict[str, str]
    m: MyMessage


class Response(Message):
    name: str
    age: int
    d: dict[str, str]
    m: MyMessage | str


class BarService:
    def bar(self, req: BarRequest, ctx: ServicerContext) -> BarResponse:
        return BarResponse(names=req.names)


class CustomInterceptor(grpc.ServerInterceptor):
    def intercept_service(self, continuation, handler_call_details):
        # do something
        print(handler_call_details.method)
        return continuation(handler_call_details)


async def app(scope, receive, send):
    pass


if __name__ == "__main__":
    s = Server(10, CustomInterceptor())
    s.run(
        FooService(),
        BarService(),
    )

🩺 [TODO] Verificación de Salud Personalizada

TODO

🤖 Soporte MCP (Protocolo de Contexto de Modelo)

PydanticRPC puede exponer tus servicios como herramientas MCP para asistentes de IA usando FastMCP. Esto permite una integración perfecta con cualquier cliente compatible con MCP.

Ejemplo de Modo Stdio

from pydantic_rpc import Message
from pydantic_rpc.mcp import MCPExporter

class CalculateRequest(Message):
    expression: str

class CalculateResponse(Message):
    result: float

class MathService:
    def calculate(self, req: CalculateRequest) -> CalculateResponse:
        result = eval(req.expression, {"__builtins__": {}}, {})
        return CalculateResponse(result=float(result))

# Run as MCP stdio server
if __name__ == "__main__":
    service = MathService()
    mcp = MCPExporter(service)
    mcp.run_stdio()

Configuración de Clientes MCP

Cualquier cliente compatible con MCP puede conectarse a tu servicio. Por ejemplo, para configurar Claude Desktop:

{
  "mcpServers": {
    "my-math-service": {
      "command": "python",
      "args": ["/path/to/math_mcp_server.py"]
    }
  }
}

Ejemplo de Modo HTTP/ASGI

MCP también se puede montar en aplicaciones ASGI existentes:

from pydantic_rpc import ASGIApp
from pydantic_rpc.mcp import MCPExporter

# Create Connect-RPC ASGI app
app = ASGIApp()
app.mount(MathService())

# Add MCP support via HTTP/SSE
mcp = MCPExporter(MathService())
mcp.mount_to_asgi(app, path="/mcp")

# Run with uvicorn
import uvicorn
uvicorn.run(app, host="127.0.0.1", port=8000)

Los endpoints de MCP estarán disponibles en:

  • SSE: GET http://localhost:8000/mcp/sse
  • Mensajes: POST http://localhost:8000/mcp/messages/

🗄️ Herramienta CLI (pydantic-rpc-cli)

La herramienta CLI proporciona potentes funciones para generar archivos protobuf y ejecutar servidores. Instálala por separado:

pip install pydantic-rpc-cli

Generar Archivos Protobuf

# Generate .proto file from a service class
pydantic-rpc generate myapp.services.UserService --output ./proto/

# Also compile to Python code
pydantic-rpc generate myapp.services.UserService --compile

Ejecutar Servidores Directamente

La CLI puede ejecutar cualquier tipo de servidor:

# Run as gRPC server (auto-detects async/sync)
pydantic-rpc serve myapp.services.UserService --port 50051

# Run as Connect-RPC with ASGI (HTTP/2, uses Hypercorn)
pydantic-rpc serve myapp.services.UserService --asgi --port 8000

# Run as Connect-RPC with WSGI (HTTP/1.1, uses Gunicorn)
pydantic-rpc serve myapp.services.UserService --wsgi --port 8000 --workers 4

Usando los archivos proto generados con herramientas como protoc, buf y BSR, puedes generar código para cualquier lenguaje deseado.

📖 Mapeo de Tipos de Datos

Tipo PythonTipo Protobuf
strstring
bytesbytes
boolbool
intint32
floatfloat, double
list[T], tuple[T]repeated T
dict[K, V]map<K, V>
datetime.datetimegoogle.protobuf.Timestamp
datetime.timedeltagoogle.protobuf.Duration
typing.Union[A, B]oneof A, B
subclase de enum.Enumenum
subclase de pydantic.BaseModelmessage

⚠️ Limitaciones Conocidas

Tipos Union con Colecciones

Debido a las restricciones de oneof de protobuf, no puedes usar tipos Union que contengan campos repeated (lista/tupla) o map (dict) directamente. Esta es una limitación de la propia especificación protobuf.

❌ No Compatible:

from typing import Union, List, Dict
from pydantic_rpc import Message

# These will fail during proto compilation
class MyMessage(Message):
    # Union with list - NOT SUPPORTED
    field1: Union[List[int], str]

    # Union with dict - NOT SUPPORTED
    field2: Union[Dict[str, int], int]

    # Union with nested collections - NOT SUPPORTED
    field3: Union[List[Dict[str, int]], str]

✅ Solución Alternativa - Usar Envoltorios de Mensaje:

from typing import Union, List, Dict
from pydantic_rpc import Message

# Wrap collections in Message types
class IntList(Message):
    values: List[int]

class StringIntMap(Message):
    values: Dict[str, int]

class MyMessage(Message):
    # Now these work!
    field1: Union[IntList, str]
    field2: Union[StringIntMap, int]

Este enfoque funciona porque protobuf permite tipos de mensaje dentro de campos oneof, y las colecciones están contenidas dentro de esos mensajes.

🔧 Desarrollo

Este proyecto usa just como ejecutor de comandos para tareas de desarrollo.

Instalando just

macOS:

brew install just

Linux:

curl --proto '=https' --tlsv1.2 -sSf https://just.systems/install.sh | bash -s -- --to ~/bin

Windows: Descarga desde GitHub releases

Inicio Rápido

# Install dependencies
just install

# Run tests
just test  # or just t

# Format and lint code
just format  # or just f
just lint    # or just l

# Run all checks (lint + tests)
just check   # or just c

# See all available commands
just --list

Ejecutando Ejemplos

# Start servers
just greeting-server  # gRPC server on port 50051
just greeting-asgi    # Connect RPC ASGI on port 8000
just greeting-wsgi    # Connect RPC WSGI on port 3000

# Test with buf curl (in another terminal)
just greet            # gRPC request
just connect-greet    # Connect RPC request
just wsgi-greet       # WSGI request

# Custom names
just greet-name Alice
just connect-greet-name Bob

Para más comandos y opciones de desarrollo, consulta el Justfile o ejecuta just --list.

TODO

  • Soporte de Streaming
    • unario-stream
    • stream-unario
    • stream-stream
  • Soporte de Mensajes Vacíos (google.protobuf.Empty automático)
  • Soporte de Serializadores Pydantic (@model_serializer, @field_serializer)
  • Soporte de Verificación de Salud Personalizada
  • Soporte MCP (Protocolo de Contexto de Modelo) vía SDK oficial de MCP
  • Agregar más ejemplos
  • Agregar pruebas

📜 Licencia

Este proyecto está licenciado bajo la Licencia MIT. Consulta el archivo LICENSE para más detalles.