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
pydanticpara 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.
- 💚 Verificación de Salud: Soporte integrado para verificaciones de salud gRPC usando
- 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
- 🌐 Soporte Completo de Protocolo: Soporte nativo de Connect-RPC a través de
- 🛠️ 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/jsonoapplication/connect+jsonpara 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
ServeryAsyncIOServer - Connect-RPC: RPC moderno basado en HTTP con
ASGIAppyWSGIApp
🔧 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-pythonpara 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 Python | Tipo Protobuf |
|---|---|
| str | string |
| bytes | bytes |
| bool | bool |
| int | int32 |
| float | float, double |
| list[T], tuple[T] | repeated T |
| dict[K, V] | map<K, V> |
| datetime.datetime | google.protobuf.Timestamp |
| datetime.timedelta | google.protobuf.Duration |
| typing.Union[A, B] | oneof A, B |
| subclase de enum.Enum | enum |
| subclase de pydantic.BaseModel | message |
⚠️ 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.