Examples¶
Each script below is run by grpclib-transports/tests/test_examples.py, so a
change to the library that breaks one of them fails the test suite.
Common service implementations¶
"""Common protocol implementations shared by the runnable examples."""
from __future__ import annotations
from typing import TYPE_CHECKING, override
import greeter.greeter.common as common_pb2
import greeter.greeter.server as server_grpc
import greeter.greeter.worker as worker_grpc
if TYPE_CHECKING:
from collections.abc import AsyncIterator
class Greeter(server_grpc.GreeterBase):
@override
async def say_hello(self, message: common_pb2.HelloRequest) -> common_pb2.HelloReply:
return common_pb2.HelloReply(message=f"Hello, {message.name}!")
@override
async def upload(self, messages: AsyncIterator[common_pb2.HelloRequest]) -> common_pb2.HelloReply:
total = 0
async for request in messages:
total += len(request.payload)
return common_pb2.HelloReply(message=f"Uploaded {total} bytes")
class WorkerGreeter(worker_grpc.GreeterWorkerBase):
@override
async def say_hello(self, message: common_pb2.HelloRequest) -> common_pb2.HelloReply:
return common_pb2.HelloReply(message=f"Hello, {message.name}!")
@override
async def upload(self, messages: AsyncIterator[common_pb2.HelloRequest]) -> common_pb2.HelloReply:
total = 0
async for request in messages:
total += len(request.payload)
return common_pb2.HelloReply(message=f"Uploaded {total} bytes")
class GreeterManager(worker_grpc.GreeterManagerBase):
@override
async def lookup(self, message: common_pb2.ManagerLookupRequest) -> common_pb2.ManagerLookupReply:
return common_pb2.ManagerLookupReply(value=f"manager:{message.key}")
def worker_services() -> list[WorkerGreeter]:
"""The worker payload, here rather than in the script that spawns it.
The forkserver pickles this by name, and the child resolves that name by
importing the module. A factory defined in the ``__main__`` of the script
has no module the child can import: see ``main_module_not_reexecuted``,
which keeps the child from running that script a second time to make one.
"""
return [WorkerGreeter()]
Unix domain socket¶
"""Unix domain socket transport — in-process server and client.
Run with::
python docs/examples/unix_example.py
"""
from __future__ import annotations
import asyncio
import contextlib
import os
import tempfile
import greeter.greeter.common as common_pb2
import greeter.greeter.server as server_grpc
from anyio import Path
from grpclib_transports import Server, connect_unix
from services import Greeter
async def main() -> None:
fd, sock = tempfile.mkstemp(suffix=".sock")
os.close(fd)
await Path(sock).unlink()
channel = None
try:
async with Server() as server:
await server.endpoint([Greeter()]).listen_unix(sock)
channel = connect_unix(sock)
try:
stub = server_grpc.GreeterStub(channel)
response = await stub.say_hello(common_pb2.HelloRequest(name="World"))
assert response.message == "Hello, World!"
print(f"Greeter replied: {response.message}")
finally:
channel.close()
channel = None
finally:
if channel is not None:
channel.close()
with contextlib.suppress(OSError):
await Path(sock).unlink()
if __name__ == "__main__":
asyncio.run(main())
Stdio subprocess¶
"""Stdio transport — subprocess server + client over stdin/stdout.
Run with::
python docs/examples/stdio_example.py
"""
from __future__ import annotations
import asyncio
import sys
import greeter.greeter.common as common_pb2
import greeter.greeter.worker as worker_grpc
from grpclib_transports import Server
async def main() -> None:
async with Server() as server:
workers = server.endpoint([]).for_workers()
async with workers.stdio_channels(
[
sys.executable,
"-m",
"grpclib_transports",
"server",
"--stdio",
"--max-concurrency",
"1",
],
client_factory=worker_grpc.GreeterWorkerStub,
stderr=asyncio.subprocess.PIPE,
) as pool:
stub = pool[0].client
response = await stub.say_hello(common_pb2.HelloRequest(name="Stdio"))
assert response.message == "Hello, Stdio!"
print(f"Greeter replied: {response.message}")
if __name__ == "__main__":
asyncio.run(main())
Multiprocessing pipe pair¶
"""Multiprocessing transport — forkserver worker speaks gRPC over OS pipes.
Run with::
python docs/examples/multiprocessing_example.py
"""
from __future__ import annotations
import asyncio
import greeter.greeter.common as common_pb2
import greeter.greeter.worker as worker_grpc
from grpclib_transports import Server
from services import worker_services
async def main() -> None:
async with Server() as server:
workers = server.endpoint([]).for_workers()
async with workers.multiprocessing_channels(
worker_services,
client_factory=worker_grpc.GreeterWorkerStub,
preload=["greeter"],
max_concurrency=1,
) as pool:
stub = pool[0].client
response = await stub.say_hello(common_pb2.HelloRequest(name="Multiprocessing"))
assert response.message == "Hello, Multiprocessing!"
print(f"Greeter replied: {response.message}")
if __name__ == "__main__":
asyncio.run(main())
SSH¶
"""SSH transport — in-process server and client over SSH on a Unix socket.
Requires asyncssh.
Run with::
python docs/examples/ssh_example.py
"""
from __future__ import annotations
import asyncio
import contextlib
import os
import socket
import tempfile
from typing import Any
import greeter.greeter.common as common_pb2
import greeter.greeter.server as server_grpc
from anyio import Path
from grpclib_transports import SshChannel, SshTransport
from grpclib_transports.protocol import serve_h2
from services import Greeter
async def main() -> None:
import asyncssh
fd, sock = tempfile.mkstemp(suffix=".sock")
os.close(fd)
await Path(sock).unlink()
key = asyncssh.generate_private_key("ssh-ed25519") # pyright: ignore[reportUnknownMemberType] -- asyncssh type stubs are incomplete
class _Server(asyncssh.SSHServer):
def password_auth_supported(self):
return True
def validate_password(self, username: str, password: str):
return True
async def session_handler(stdin: Any, stdout: Any, _stderr: Any):
transport = SshTransport(stdin, stdout)
await serve_h2([Greeter()], stdin, transport)
ssock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
ssock.bind(sock)
ssock.listen()
try:
acceptor = await asyncssh.listen(
sock=ssock,
server_host_keys=[key],
server_factory=_Server,
session_factory=session_handler,
encoding=None,
line_editor=False,
)
try:
csock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
csock.connect(sock)
conn = await asyncssh.connect(
sock=csock,
known_hosts=None,
username="x",
password="x",
)
channel = None
try:
stdin, stdout, _ = await conn.open_session(encoding=None) # pyright: ignore[reportUnknownVariableType,reportUnknownMemberType] -- asyncssh type stubs are incomplete
channel = SshChannel(stdout, stdin)
stub = server_grpc.GreeterStub(channel)
response = await stub.say_hello(common_pb2.HelloRequest(name="SSH"))
assert response.message == "Hello, SSH!"
print(f"Greeter replied: {response.message}")
finally:
if channel is not None:
await channel.aclose()
conn.close()
finally:
acceptor.close()
await acceptor.wait_closed()
finally:
with contextlib.suppress(OSError):
await Path(sock).unlink()
if __name__ == "__main__":
asyncio.run(main())
Bidirectional logical RPC¶
"""Bidirectional logical RPC — in-process request/response without gRPC.
Demonstrates :class:`~grpclib_transports.LogicalRpcPeer` for lightweight
in-process RPC that doesn't require protobuf or HTTP/2.
Run with::
python docs/examples/bidi_example.py
"""
from __future__ import annotations
import asyncio
from grpclib_transports import LogicalFrame, LogicalRpcPeer, PeerClosedError
async def handler(method: str, payload: str) -> str:
print(f" [server] received {method!r} with {payload!r}")
return payload.upper()
async def main() -> None:
up: asyncio.Queue[LogicalFrame | None] = asyncio.Queue()
down: asyncio.Queue[LogicalFrame | None] = asyncio.Queue()
server = LogicalRpcPeer(
send_frame=down.put,
receive_frame=up.get,
handler=handler,
)
client = LogicalRpcPeer(
send_frame=up.put,
receive_frame=down.get,
)
server.start()
client.start()
try:
response = await client.call("echo", "hello world")
assert response == "HELLO WORLD"
print(f"client got: {response}")
await client.event("notify", "side-effect")
finally:
await client.aclose()
await server.aclose()
try:
await client.call("fail", None)
except PeerClosedError:
print("correctly raised PeerClosedError after close")
if __name__ == "__main__":
asyncio.run(main())