PydanticRPC

Uma biblioteca Python para construir serviços gRPC/ConnectRPC com modelos Pydantic, com geração automática de protobuf e exposição de ferramentas para assistentes de IA.

Documentação

🚀 PydanticRPC

PydanticRPC é uma biblioteca Python que permite expor rapidamente modelos Pydantic por meio de serviços gRPC/Connect RPC sem escrever nenhum arquivo protobuf. Em vez disso, ela gera automaticamente arquivos protobuf em tempo de execução a partir das assinaturas de métodos dos seus objetos Python e das assinaturas de tipos dos seus modelos Pydantic.

Abaixo está um exemplo de um serviço gRPC simples que expõe um 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())

E aqui está um exemplo de um serviço Connect RPC simples que expõe o mesmo agente como uma aplicação 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())

💡 Principais Recursos

  • 🔄 Geração Automática de Protobuf: Cria automaticamente arquivos protobuf correspondentes às assinaturas de métodos dos seus objetos Python.
  • ⚙️ Geração Dinâmica de Código: Gera stubs de servidor e cliente usando grpcio-tools.
  • ✅ Integração com Pydantic: Usa pydantic para validação de tipos e serialização robustas.
  • 📄 Exportação de Arquivos Protobuf: Exporta os arquivos protobuf gerados para uso em outras linguagens.
  • Para gRPC:
    • 💚 Verificação de Saúde: Suporte integrado para verificações de saúde gRPC usando grpc_health.v1.
    • 🔎 Reflexão de Servidor: Suporte integrado para reflexão de servidor gRPC.
    • ⚡ Suporte Assíncrono: Crie facilmente serviços gRPC assíncronos com AsyncIOServer.
  • Para Connect-RPC:
    • 🌐 Suporte Completo ao Protocolo: Suporte nativo a Connect-RPC via connect-python
    • 🔄 Todos os Padrões de Streaming: Unário, streaming de servidor, streaming de cliente e streaming bidirecional
    • 🌐 Aplicações WSGI/ASGI: Execute como aplicações WSGI ou ASGI padrão para implantação fácil
  • 🛠️ Arquivos Protobuf e Código Pré-gerados: Pré-gera arquivos proto e código correspondente via CLI. Ao definir a variável de ambiente (PYDANTIC_RPC_SKIP_GENERATION), você pode pular a geração em tempo de execução.
  • 🤖 Suporte a MCP (Model Context Protocol): Exponha seus serviços como ferramentas para assistentes de IA usando o SDK oficial do MCP, com suporte aos transportes stdio e HTTP/SSE.

⚠️ Notas Importantes para Connect-RPC

Ao usar Connect-RPC com ASGIApp:

  • Formato do Caminho do Endpoint: Endpoints Connect-RPC usam nomes de métodos em CamelCase no caminho: /<package>.<service>/<Method> (por exemplo, /chat.v1.ChatService/SendMessage)
  • Content-Type: Defina Content-Type: application/json ou application/connect+json para requisições
  • Requisito HTTP/2: Streaming bidirecional requer HTTP/2. Use Hypercorn em vez de uvicorn para suporte a HTTP/2
  • Testes: Use buf curl para testar endpoints Connect-RPC com suporte adequado a streaming

Para exemplos detalhados e instruções de teste, consulte o diretório de exemplos.

📦 Instalação

Instale o PydanticRPC via pip:

pip install pydantic-rpc

Para suporte a CLI com executores de servidor integrados:

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

🆕 Recursos Aprimorados (v0.10.0+)

Nota: Todos os novos recursos são totalmente compatíveis com versões anteriores. O código existente continua funcionando sem modificações.

API de Inicialização Aprimorada

Todas as classes de servidor agora suportam inicialização opcional com serviços:

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

Tratamento de Erros com Decoradores

Mapeie automaticamente exceções para códigos de status 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]

🚀 Começando

O PydanticRPC suporta dois protocolos principais:

  • gRPC: Serviços gRPC tradicionais com Server e AsyncIOServer
  • Connect-RPC: RPC moderno baseado em HTTP com ASGIApp e WSGIApp

🔧 Exemplo de Serviço 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())

⚙️ Exemplo de Serviço gRPC Assí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())

O AsyncIOServer lida automaticamente com o desligamento gracioso nos sinais SIGTERM e SIGINT.

🌐 Exemplo de Aplicação 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

🌐 Exemplo de Aplicação 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

🏆 Exemplo de Connect-RPC com Streaming

O PydanticRPC fornece suporte nativo a Connect-RPC via connect-python, incluindo recursos completos de streaming e convenções de nomenclatura PEP 8. Confira nossos exemplos ASGI:

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

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

Isso iniciará uma aplicação ASGI baseada em connect-python que usa os mesmos modelos Pydantic para atender requisições Connect-RPC.

Suporte a Streaming com connect-python

connect-python fornece suporte completo para RPCs de streaming com 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, instale protoc-gen-connect-python para executar o exemplo connect-python.

♻️ Pulando a Geração de Protobuf

Por padrão, o PydanticRPC gera arquivos .proto e código em tempo de execução. Se você deseja pular a etapa de geração de código (por exemplo, em ambiente de produção), defina a variável de ambiente abaixo:

export PYDANTIC_RPC_SKIP_GENERATION=true

Quando esta variável estiver definida como "true", o PydanticRPC carregará módulos pré-gerados existentes em vez de gerá-los em tempo de execução.

🪧 Definindo o diretório de geração de Protobuf e Connect RPC/gRPC

Por padrão, seus arquivos serão gerados no diretório de trabalho atual de onde você executou o código, mas você pode definir um diretório específico personalizado definindo a variável de ambiente abaixo:

export PYDANTIC_RPC_PROTO_PATH=/your/path

⚠️ Campos Reservados

Você também pode definir uma variável de ambiente para reservar um número definido de campos para a geração de proto, para compatibilidade retroativa e futura.

export PYDANTIC_RPC_RESERVED_FIELDS=1

💎 Recursos Avançados

🌊 Streaming de Resposta (gRPC)

O PydanticRPC suporta respostas de streaming para serviços gRPC e Connect-RPC. Se o tipo de retorno de um método de classe de serviço for typing.AsyncIterator[T], o método é considerado um método de streaming.

Consulte o código de exemplo abaixo:

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

No exemplo acima, o método ask_stream retorna um objeto AsyncIterator[StreamingResult], que é considerado um método de streaming. A classe StreamingResult é um modelo Pydantic que define o tipo de resposta do método de streaming. Você pode usar qualquer modelo Pydantic como tipo de resposta.

Agora, você pode chamar o método ask_stream do servidor descrito acima usando sua ferramenta de cliente gRPC preferida. O exemplo abaixo 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
%

🪶 Mensagens Vazias

Mensagens de requisição/resposta vazias são automaticamente mapeadas para 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!")

🎨 Serialização Personalizada

Os decoradores de serialização do Pydantic são totalmente suportados:

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
        }

Os serializadores são aplicados automaticamente ao converter entre modelos Pydantic e mensagens protobuf.

⚠️ Limitações e Considerações

1. Serializadores de Mensagens Aninhadas agora são suportados (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

Controle de Estratégia de Serialização: Você pode controlar como os serializadores aninhados são aplicados via variável de ambiente:

# 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 no Desempenho:

  • Estratégia DEEP: ~4% de overhead para estruturas aninhadas simples
  • Estratégia SHALLOW: ~2% de overhead (apenas nível superior)
  • Estratégia NONE: Sem overhead (serializadores desabilitados)

2. Novos campos adicionados por serializadores são ignorados

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: O campo result não existe na definição da Mensagem, então não está no esquema protobuf.

3. O tipo deve 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. Campos Union/Opcionais têm suporte 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. Erros falham silenciosamente com fallback

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. Referências circulares são tratadas com elegância

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

Melhores Práticas:

  • Use serializadores principalmente para tipos primitivos (str, int, float, bool)
  • Mantenha a consistência de tipos (int → int, str → str)
  • Evite transformações complexas ou efeitos colaterais
  • Teste casos de erro minuciosamente
  • Esteja ciente de que erros falham silenciosamente

🔒 Suporte a TLS/mTLS

O PydanticRPC fornece suporte integrado para TLS (Transport Layer Security) e mTLS (TLS mútuo) para comunicação 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 um exemplo completo, consulte examples/tls_server.py e examples/tls_client.py.

🔗 Múltiplos Serviços com Interceptadores Personalizados

O PydanticRPC suporta definir e executar múltiplos serviços gRPC em um único 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] Verificação de Saúde Personalizada

TODO

🤖 Suporte a MCP (Model Context Protocol)

O PydanticRPC pode expor seus serviços como ferramentas MCP para assistentes de IA usando FastMCP. Isso permite integração perfeita com qualquer cliente compatível com MCP.

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

Configurando Clientes MCP

Qualquer cliente compatível com MCP pode se conectar ao seu serviço. Por exemplo, para configurar o Claude Desktop:

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

Exemplo de Modo HTTP/ASGI

O MCP também pode ser montado em aplicações 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)

Os endpoints MCP estarão disponíveis em:

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

🗄️ Ferramenta CLI (pydantic-rpc-cli)

A ferramenta CLI fornece recursos poderosos para gerar arquivos protobuf e executar servidores. Instale-a separadamente:

pip install pydantic-rpc-cli

Gerar Arquivos 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

Executar Servidores Diretamente

A CLI pode executar qualquer 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 os arquivos proto gerados com ferramentas como protoc, buf e BSR, você pode gerar código para qualquer linguagem desejada.

📖 Mapeamento de Tipos de Dados

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
subclass of enum.Enumenum
subclass of pydantic.BaseModelmessage

⚠️ Limitações Conhecidas

Tipos Union com Coleções

Devido às restrições de oneof do protobuf, você não pode usar tipos Union que contenham campos repeated (list/tuple) ou map (dict) diretamente. Esta é uma limitação da própria especificação do protobuf.

❌ Não Suportado:

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]

✅ Solução Alternativa - Use Wrappers de Mensagem:

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]

Essa abordagem funciona porque o protobuf permite tipos de mensagem dentro de campos oneof, e as coleções estão contidas dentro dessas mensagens.

🔧 Desenvolvimento

Este projeto usa just como executor de comandos para tarefas de desenvolvimento.

Instalando just

macOS:

brew install just

Linux:

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

Windows: Baixe dos releases do GitHub

Início 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

Executando Exemplos

# 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 mais comandos e opções de desenvolvimento, consulte o Justfile ou execute just --list.

TODO

  • Suporte a Streaming
    • unary-stream
    • stream-unary
    • stream-stream
  • Suporte a Mensagens Vazias (google.protobuf.Empty automático)
  • Suporte a Serializadores Pydantic (@model_serializer, @field_serializer)
  • Suporte a Verificação de Saúde Personalizada
  • Suporte a MCP (Model Context Protocol) via SDK oficial do MCP
  • Adicionar mais exemplos
  • Adicionar testes

📜 Licença

Este projeto é licenciado sob a Licença MIT. Consulte o arquivo LICENSE para obter detalhes.