Skip to content

Core RPC

The core module provides the server, connection, transport interface, error types, and convenience functions for defining and running RPC services.

Typical Usage

Most users only need serve_pipe (testing) or connect (subprocess):

from vgi_rpc import serve_pipe, connect

# In-process (tests)
with serve_pipe(MyService, MyServiceImpl()) as proxy:
    proxy.my_method(arg=42)

# Subprocess
with connect(MyService, ["python", "worker.py"]) as proxy:
    proxy.my_method(arg=42)

For more control, use RpcServer and RpcConnection directly.

API Reference

RpcServer

RpcServer

RpcServer(
    protocol: type,
    implementation: object,
    *,
    extra_protocols: Sequence[tuple[type, object]] = (),
    identity: object | None = None,
    external_location: ExternalLocationConfig | None = None,
    server_id: str | None = None,
    server_version: str = "",
    enable_describe: bool = False,
    ipc_validation: IpcValidation | None = None,
    include_tracebacks: bool = True,
    grant_keys: GrantKeys | Literal["env"] | None = "env"
)

Dispatches RPC requests to an implementation over IO-stream transports.

Initialize with a protocol type and its implementation.

PARAMETER DESCRIPTION
protocol

The Protocol class defining the RPC interface. If the class declares a protocol_version: ClassVar[str] attribute (canonical semver MAJOR.MINOR.PATCH), the server enforces an exact major+minor match on every dispatched request. See _check_protocol_version for the comparison rule.

TYPE: type

implementation

Object implementing all methods from protocol.

TYPE: object

extra_protocols

Additional (protocol, implementation) pairs to host alongside the primary one. Each is versioned and dispatched independently, so a worker protocol and a shared identity protocol can live in one process. Method names may repeat across protocols — calls resolve on (protocol, method). Protocol names may not repeat; a duplicate raises. One object may implement several protocols.

protocol stays the primary: it is what protocol_name, protocol_hash and implementation report, what the landing page titles, and what framework endpoints with no owning protocol log against. It is never a routing fallback.

TYPE: Sequence[tuple[type, object]] DEFAULT: ()

identity

An :class:~vgi_rpc.rpc._identity.IdentityImpl to host vgi_rpc.Identity.v1. Absent by default, and absent rather than routed-and-refusing when omitted: that is what keeps a dependency upgrade from growing a credential-to-identity oracle on every existing worker.

TYPE: object | None DEFAULT: None

external_location

Optional ExternalLocation configuration.

TYPE: ExternalLocationConfig | None DEFAULT: None

server_id

Optional server identifier; auto-generated if None.

TYPE: str | None DEFAULT: None

server_version

Build version string included in access log entries.

TYPE: str DEFAULT: ''

enable_describe

When True, the server also hosts the vgi_rpc.Reflection.v1 protocol, through which a client discovers what is hosted here and asks for one protocol's machine-readable method metadata. Named for the retired __describe__ method it replaced.

TYPE: bool DEFAULT: False

ipc_validation

Validation level for incoming IPC batches. None (the default) resolves from the VGI_RPC_IPC_VALIDATION environment variable (none / standard / full), falling back to FULL for maximum safety. An explicit value always wins over the environment — the variable is a deployment knob, not an override of stated intent.

TYPE: IpcValidation | None DEFAULT: None

include_tracebacks

Whether EXCEPTION batches carry the remote traceback (log_extra.traceback, frames, cause, context). Included by default on every transport; set False to omit them everywhere this server answers. The exception type, message, code, kind and details are sent either way. See WIRE_PROTOCOL.md §8.

TYPE: bool DEFAULT: True

grant_keys

Sealed-grant configuration (WIRE_PROTOCOL.md §16). "env" (the default) reads VGI_RPC_GRANT_KEYS and friends -- unset means grants are off and nothing changes. With keys, the framework mints sealed grants through issue_grant (unless the worker's IdentityImpl supplies mint_grant) and an HTTP app accepts them back as bearer credentials. None turns grants off regardless of the environment. A malformed key raises here: a worker refuses to start rather than run with a key it misread.

TYPE: GrantKeys | Literal['env'] | None DEFAULT: 'env'

Source code in vgi_rpc/rpc/_server.py
def __init__(
    self,
    protocol: type,
    implementation: object,
    *,
    extra_protocols: Sequence[tuple[type, object]] = (),
    identity: object | None = None,
    external_location: ExternalLocationConfig | None = None,
    server_id: str | None = None,
    server_version: str = "",
    enable_describe: bool = False,
    ipc_validation: IpcValidation | None = None,
    include_tracebacks: bool = True,
    grant_keys: GrantKeys | Literal["env"] | None = "env",
) -> None:
    """Initialize with a protocol type and its implementation.

    Args:
        protocol: The Protocol class defining the RPC interface. If the
            class declares a ``protocol_version: ClassVar[str]`` attribute
            (canonical semver MAJOR.MINOR.PATCH), the server enforces an
            exact major+minor match on every dispatched request. See
            ``_check_protocol_version`` for the comparison rule.
        implementation: Object implementing all methods from *protocol*.
        extra_protocols: Additional ``(protocol, implementation)`` pairs to
            host alongside the primary one. Each is versioned and dispatched
            independently, so a worker protocol and a shared identity
            protocol can live in one process. Method names may repeat across
            protocols — calls resolve on ``(protocol, method)``. Protocol
            *names* may not repeat; a duplicate raises. One object may
            implement several protocols.

            *protocol* stays the primary: it is what ``protocol_name``,
            ``protocol_hash`` and ``implementation`` report, what the
            landing page titles, and what framework endpoints with no owning
            protocol log against. It is never a routing fallback.
        identity: An :class:`~vgi_rpc.rpc._identity.IdentityImpl` to host
            ``vgi_rpc.Identity.v1``.  Absent by default, and absent rather
            than routed-and-refusing when omitted: that is what keeps a
            dependency upgrade from growing a credential-to-identity oracle
            on every existing worker.
        external_location: Optional ExternalLocation configuration.
        server_id: Optional server identifier; auto-generated if ``None``.
        server_version: Build version string included in access log entries.
        enable_describe: When ``True``, the server also hosts the
            ``vgi_rpc.Reflection.v1`` protocol, through which a client
            discovers what is hosted here and asks for one protocol's
            machine-readable method metadata.  Named for the retired
            ``__describe__`` method it replaced.
        ipc_validation: Validation level for incoming IPC batches.
            ``None`` (the default) resolves from the
            ``VGI_RPC_IPC_VALIDATION`` environment variable
            (``none`` / ``standard`` / ``full``), falling back to
            ``FULL`` for maximum safety.  An explicit value always
            wins over the environment — the variable is a deployment
            knob, not an override of stated intent.
        include_tracebacks: Whether EXCEPTION batches carry the remote
            traceback (``log_extra.traceback``, ``frames``, ``cause``,
            ``context``).  Included by default on every transport; set
            ``False`` to omit them everywhere this server answers.  The
            exception type, message, code, kind and details are sent
            either way.  See WIRE_PROTOCOL.md §8.
        grant_keys: Sealed-grant configuration (WIRE_PROTOCOL.md §16).
            ``"env"`` (the default) reads ``VGI_RPC_GRANT_KEYS`` and
            friends -- unset means grants are off and nothing changes.
            With keys, the framework mints sealed grants through
            ``issue_grant`` (unless the worker's ``IdentityImpl`` supplies
            ``mint_grant``) and an HTTP app accepts them back as bearer
            credentials.  ``None`` turns grants off regardless of the
            environment.  A malformed key raises here: a worker refuses to
            start rather than run with a key it misread.

    """
    self._protocol = protocol
    self._include_tracebacks = include_tracebacks
    if grant_keys == "env":
        from vgi_rpc.grants import GrantKeys as _GrantKeys

        # Read at construction, so a malformed key refuses to start the
        # worker rather than failing the first mint.
        grant_keys = _GrantKeys.from_env()
    resolved_grant_keys: GrantKeys | None = grant_keys if not isinstance(grant_keys, str) else None
    self._impl = implementation
    self._server_version = server_version
    self._ipc_validation = IpcValidation.from_env() if ipc_validation is None else ipc_validation
    self._external_config = external_location
    self._server_id = server_id if server_id is not None else uuid.uuid4().hex[:12]
    self._dispatch_hook: _DispatchHook | None = None
    self._transport_kind: TransportKind | None = None
    self._transport_capabilities: frozenset[str] = frozenset()
    self._transport_lock = threading.Lock()

    # Build one binding per hosted protocol. The primary comes first so its
    # values remain what the single-protocol attributes report.
    self._bindings: dict[str, _ProtocolBinding] = {}
    for proto, impl in ((protocol, implementation), *extra_protocols):
        binding = self._build_binding(proto, impl)
        existing = self._bindings.get(binding.name)
        if existing is not None:
            raise ValueError(
                f"Two protocols are hosted under the same name {binding.name!r}: "
                f"{existing.protocol.__name__} and {proto.__name__}. The name is the "
                f"routing key, so it must be unique. Declare a distinct "
                f"`protocol_name: ClassVar[str]` on one of them."
            )
        self._bindings[binding.name] = binding

    # Reflection is registered after the application protocols and before
    # anything reads `_bindings`, so it appears in its own output without
    # being special-cased -- and so `primary` below is still the first
    # *application* protocol, which is what the single-protocol attributes
    # are expected to report.
    if enable_describe:
        from ._reflection import Reflection, ReflectionImpl

        reflection = self._build_binding(Reflection, ReflectionImpl(self), allow_reserved=True)
        self._bindings[reflection.name] = dataclasses.replace(reflection, version_exempt=True)

    # Identity, when the deployment configured it.  Registered after
    # reflection so it appears in reflection's output, and with only the
    # methods whose hooks exist -- see `_build_binding(only=...)`.
    if resolved_grant_keys is not None:
        from ._token_identity import IdentityImpl as _IdentityImpl

        if identity is None:
            # Grants on, no other identity hooks: the framework mints and
            # accepts its own, and hosts issue_grant alone.
            identity = _IdentityImpl(grant_keys=resolved_grant_keys)
        elif isinstance(identity, _IdentityImpl) and identity.grant_keys is None:
            raise ValueError(
                "grant keys were configured (grant_keys= or VGI_RPC_GRANT_KEYS) and an IdentityImpl "
                "was passed without them. Pass IdentityImpl(grant_keys=...) so the minter and the "
                "verifier use the same keys."
            )
    self._identity = identity
    if identity is not None:
        from ._token_identity import Identity, IdentityImpl

        if not isinstance(identity, IdentityImpl):
            raise TypeError(f"identity must be an IdentityImpl, got {type(identity).__name__}")
        offered = identity.offered_methods()
        if offered:
            self._bindings[Identity.protocol_name] = self._build_binding(
                Identity, identity, allow_reserved=True, only=offered
            )

    primary = next(iter(self._bindings.values()))
    self._protocol_version: str | None = primary.version
    self._protocol_version_parts: tuple[int, int, int] | None = primary.version_parts
    self._methods = primary.methods
    self._protocol_hash: str = primary.protocol_hash

    # Legal, and the whole point of namespacing — but an operator should
    # learn about the ambiguity at startup rather than from a dashboard that
    # merges two protocols' traffic under one method name.
    if len(self._bindings) > 1:
        seen: dict[str, str] = {}
        for b in self._bindings.values():
            for method_name in b.methods:
                owner = seen.setdefault(method_name, b.name)
                if owner != b.name:
                    _logger.warning(
                        "Method name %r is defined by both %s and %s. Calls resolve "
                        "correctly on (protocol, method), but anything keyed on the "
                        "bare method name — dashboards, alerts, proxy policy — will "
                        "merge them.",
                        method_name,
                        owner,
                        b.name,
                        extra={"server_id": self._server_id, "method": method_name},
                    )

    # Primary's, for the single-protocol attribute. Dispatch uses the
    # binding's own set so a method taking `ctx` on one protocol does not
    # grant it to a same-named method on another.
    self._ctx_methods: frozenset[str] = primary.ctx_methods

    _logger.info(
        "RpcServer created for %s (server_id=%s, methods=%d)",
        protocol.__name__,
        self._server_id,
        len(self._methods),
        extra={"server_id": self._server_id, "protocol": protocol.__name__, "method_count": len(self._methods)},
    )

    # Auto-attach Sentry instrumentation when the SDK is initialised in
    # this process.  We only consult sentry_sdk when it is already
    # imported, so this never forces the optional dependency on users
    # who have not opted into Sentry.
    if "sentry_sdk" in sys.modules:
        try:
            from vgi_rpc.sentry import _maybe_auto_instrument

            _maybe_auto_instrument(self)
        except ImportError:
            _logger.debug("sentry_sdk imported but vgi_rpc.sentry unavailable", exc_info=True)

methods property

methods: Mapping[str, RpcMethodInfo]

Return method metadata for this server's protocol.

bindings property

bindings: Mapping[str, _ProtocolBinding]

The protocols this server hosts, keyed by wire name, primary first.

Framework-internal: dispatch, state-type resolution and telemetry read it. Ordinary callers want methods or implementation_for.

Read-only. The hosted set is sealed at construction -- earlier than the "before serving starts" WIRE_PROTOCOL.md §3.1 requires -- and there is no registration API afterwards, so a protocol cannot appear on one transport and not another, or change what reflection already reported.

implementation property

implementation: object

The implementation object.

external_config property

external_config: ExternalLocationConfig | None

The ExternalLocation configuration, if any.

server_id property

server_id: str

Short random identifier for this server instance.

protocol_name property

protocol_name: str

Wire name of the primary protocol.

The declared protocol_name ClassVar when there is one, otherwise the class name — so this is the routing key, not a Python identifier that merely resembles it. A server hosting several protocols reports the primary here; per-call labelling reads RpcMethodInfo.protocol_name.

server_version property

server_version: str

Version string passed at construction (empty if not set).

protocol_version property

protocol_version: str | None

Application protocol surface version declared by the Protocol class.

Read from vars(protocol).get("protocol_version") at construction (a ClassVar[str] in canonical semver MAJOR.MINOR.PATCH form, or None when the Protocol opts out). When set, the server enforces an exact major+minor match on every dispatched request via _check_protocol_version.

protocol_hash property

protocol_hash: str

SHA-256 hex digest of the canonical describe payload.

identity property

identity: IdentityImpl | None

The hosted vgi_rpc.Identity.v1 implementation, when there is one.

grant_keys property

grant_keys: GrantKeys | None

The sealed-grant configuration, when grants are on.

include_tracebacks property

include_tracebacks: bool

Whether EXCEPTION batches carry the remote traceback, on every transport.

On by default everywhere. An earlier draft omitted it on HTTP and TCP; that hid chained causes from the DuckDB extension, which puts the remote traceback into the user-visible error. The setting stays so an operator who does not want stack traces leaving the process can turn them off for the whole server.

ctx_methods property

ctx_methods: frozenset[str]

Method names whose implementations accept a ctx parameter.

describe_enabled property

describe_enabled: bool

Whether this server hosts the reflection protocol.

Named for the enable_describe constructor argument it reports, and kept because the HTTP factory gates its human-readable describe page on it.

ipc_validation property

ipc_validation: IpcValidation

Validation level for incoming IPC batches.

transport_kind property

transport_kind: TransportKind | None

Coarse identifier of the bound transport, or None before serving begins.

Set by the framework right before the first request is dispatched (lazy on HTTP for fork-safety). Workers may read this directly, or rely on the on_serve_start lifecycle hook for one-shot startup work.

transport_capabilities property

transport_capabilities: frozenset[str]

Capabilities advertised by the bound transport.

Currently includes "shm" when a :class:ShmPipeTransport is bound. Empty before a transport is bound and for kinds without special capabilities.

implementation_for

implementation_for(info: RpcMethodInfo) -> object

Return the implementation that owns info's method.

implementation keeps returning the primary, because ~9 call sites read it and silently changing its meaning is worse than either leaving it or replacing it. Dispatch and stream rehydration use this instead, so a second protocol's state is never rehydrated against the first protocol's object.

Source code in vgi_rpc/rpc/_server.py
def implementation_for(self, info: RpcMethodInfo) -> object:
    """Return the implementation that owns *info*'s method.

    ``implementation`` keeps returning the primary, because ~9 call sites
    read it and silently changing its meaning is worse than either leaving
    it or replacing it. Dispatch and stream rehydration use this instead, so
    a second protocol's state is never rehydrated against the first
    protocol's object.
    """
    binding = self._bindings.get(info.protocol_name)
    return binding.impl if binding is not None else self._impl

wants_ctx

wants_ctx(info: RpcMethodInfo) -> bool

Whether info's owning implementation declares a ctx parameter.

Resolved against the binding that owns the method rather than against the primary's set. ctx_methods is per binding precisely because two protocols may define the same method name and only one may want a :class:CallContext -- but every dispatch path read the primary's set, so a secondary protocol whose methods all take ctx got none of them. vgi_rpc.Identity.v1 is exactly that shape: both its methods need the caller's :class:AuthContext to apply their guards, and neither could be called at all until this looked at the right binding.

Source code in vgi_rpc/rpc/_server.py
def wants_ctx(self, info: RpcMethodInfo) -> bool:
    """Whether *info*'s owning implementation declares a ``ctx`` parameter.

    Resolved against the binding that owns the method rather than against
    the primary's set.  ``ctx_methods`` is per binding precisely because
    two protocols may define the same method name and only one may want a
    :class:`CallContext` -- but every dispatch path read the *primary's*
    set, so a secondary protocol whose methods all take ``ctx`` got none of
    them.  ``vgi_rpc.Identity.v1`` is exactly that shape: both its methods
    need the caller's :class:`AuthContext` to apply their guards, and
    neither could be called at all until this looked at the right binding.
    """
    binding = self._bindings.get(info.protocol_name)
    return info.name in (binding.ctx_methods if binding is not None else self._ctx_methods)

protocol_hash_for

protocol_hash_for(info: RpcMethodInfo | None) -> str

Return the canonical hash of the protocol that owns info's method.

Pairs with the protocol field, which is already per-binding. The hash was not, and the two disagreeing is worse than either being wrong alone: access-log-spec.md makes protocol_hash the registry key for decoding archived records, so a record naming one protocol and carrying another's digest is decoded against the wrong description -- and nothing about it looks wrong.

Falls back to the primary for a framework endpoint that belongs to no protocol, which is what the spec prescribes for those -- and for an unresolved method, where there is no owner to name.

Source code in vgi_rpc/rpc/_server.py
def protocol_hash_for(self, info: RpcMethodInfo | None) -> str:
    """Return the canonical hash of the protocol that owns *info*'s method.

    Pairs with the ``protocol`` field, which is already per-binding.  The
    hash was not, and the two disagreeing is worse than either being wrong
    alone: ``access-log-spec.md`` makes ``protocol_hash`` the registry key
    for decoding archived records, so a record naming one protocol and
    carrying another's digest is decoded against the wrong description --
    and nothing about it looks wrong.

    Falls back to the primary for a framework endpoint that belongs to no
    protocol, which is what the spec prescribes for those -- and for an
    unresolved method, where there is no owner to name.
    """
    if info is None:
        return self._protocol_hash
    binding = self._bindings.get(info.protocol_name)
    return binding.protocol_hash if binding is not None else self._protocol_hash

check_protocol_agreement

check_protocol_agreement(info: RpcMethodInfo) -> None

Require the request's routing metadata to agree with info.

On HTTP the protocol rides twice: in vgi_rpc.protocol and as a path segment. The metadata field is canonical -- it is the only carrier on the stdio, unix and named-pipe transports -- and the path segment is a required faithful projection, present so an edge device can act on the protocol without an Arrow parser.

Left unchecked, the two may disagree, and then edge policy is applied to one protocol while the worker runs another: the Content-Length/Transfer-Encoding shape. Mirrors the vgi_rpc.method check the HTTP dispatchers already make.

Absent is an error, not an exemption. The rationale that once stood here -- "the single-carrier case: the caller resolved from that key to begin with" -- is true on stdio, unix and named pipes and false here, because on HTTP the caller resolved from the path. Accepting a request with no routing key therefore means routing on the projection alone, which is exactly the case the canonical/projection split exists to catch: an intermediary that rewrites the path cannot touch the metadata, so a rewrite is detectable only while both carriers are required to be present and to agree.

PARAMETER DESCRIPTION
info

The method resolved from the path segment. Its protocol_name is the projection side of the comparison; the request's vgi_rpc.protocol metadata is the canonical side.

TYPE: RpcMethodInfo

RAISES DESCRIPTION
ProtocolNotSpecifiedError

The request carried no routing key.

ProtocolNotSupportedError

The two carriers name different protocols.

Source code in vgi_rpc/rpc/_server.py
def check_protocol_agreement(self, info: RpcMethodInfo) -> None:
    """Require the request's routing metadata to agree with *info*.

    On HTTP the protocol rides twice: in ``vgi_rpc.protocol`` and as a path
    segment.  The metadata field is canonical -- it is the only carrier on
    the stdio, unix and named-pipe transports -- and the path segment is a
    required faithful projection, present so an edge device can act on the
    protocol without an Arrow parser.

    Left unchecked, the two may disagree, and then edge policy is applied
    to one protocol while the worker runs another: the
    Content-Length/Transfer-Encoding shape.  Mirrors the ``vgi_rpc.method``
    check the HTTP dispatchers already make.

    Absent is an error, not an exemption.  The rationale that once stood
    here -- "the single-carrier case: the caller resolved from that key to
    begin with" -- is true on stdio, unix and named pipes and false here,
    because on HTTP the caller resolved from the *path*.  Accepting a
    request with no routing key therefore means routing on the projection
    alone, which is exactly the case the canonical/projection split exists
    to catch: an intermediary that rewrites the path cannot touch the
    metadata, so a rewrite is detectable only while both carriers are
    required to be present and to agree.

    Args:
        info: The method resolved from the path segment.  Its
            ``protocol_name`` is the projection side of the comparison;
            the request's ``vgi_rpc.protocol`` metadata is the canonical
            side.

    Raises:
        ProtocolNotSpecifiedError: The request carried no routing key.
        ProtocolNotSupportedError: The two carriers name different
            protocols.

    """
    md = _current_request_metadata.get()
    raw = md.get(PROTOCOL_KEY) if md is not None else None
    if not raw:
        # NOT YET ENFORCED.  The plan requires this to raise
        # ProtocolNotSpecifiedError -- see the note above -- and C# already
        # implements it that way.  Turning it on here fails 69 tests whose
        # hand-built requests carry no routing key, most of them in the
        # shared cross-language conformance harness, so flipping it is its
        # own change that has to move the harness and the six ports
        # together.  Left as a no-op rather than half-enforced.
        return
    declared = raw.decode(errors="replace")
    if declared != info.protocol_name:
        raise ProtocolNotSupportedError(
            f"Protocol mismatch: the request path resolved to {info.protocol_name!r} but "
            f"the Arrow IPC custom_metadata 'vgi_rpc.protocol' says {declared!r}. "
            f"These must agree."
        )

gate_version

gate_version(info: RpcMethodInfo) -> None

Enforce the declared protocol_version of the binding info belongs to.

A server hosting several protocols has a version per binding, so the gate has to read the resolved method's, not the primary's. Silently gating a secondary protocol against the primary's version is the shape of bug that produces a confusing mismatch message pointing at the wrong protocol.

No-op when the binding's Protocol declares no protocol_version, and for a binding marked version-exempt (reflection): that is the diagnostic path a mismatched client uses to find out what mismatched, so gating it would deny the client its own diagnosis.

PARAMETER DESCRIPTION
info

The resolved method. Its protocol_name selects which binding's declared version the request is gated against.

TYPE: RpcMethodInfo

RAISES DESCRIPTION
ProtocolVersionError

On a major or minor mismatch, with a directional message naming which side to upgrade.

Source code in vgi_rpc/rpc/_server.py
def gate_version(self, info: RpcMethodInfo) -> None:
    """Enforce the declared protocol_version of the binding ``info`` belongs to.

    A server hosting several protocols has a version per binding, so the
    gate has to read the resolved method's, not the primary's.  Silently
    gating a secondary protocol against the primary's version is the shape
    of bug that produces a confusing mismatch message pointing at the wrong
    protocol.

    No-op when the binding's Protocol declares no ``protocol_version``, and
    for a binding marked version-exempt (reflection): that is the
    diagnostic path a mismatched client uses to find out *what* mismatched,
    so gating it would deny the client its own diagnosis.

    Args:
        info: The resolved method.  Its ``protocol_name`` selects which
            binding's declared version the request is gated against.

    Raises:
        ProtocolVersionError: On a major or minor mismatch, with a
            directional message naming which side to upgrade.

    """
    binding = self._bindings.get(info.protocol_name)
    if binding is None or binding.version_exempt:
        return
    server_parts = binding.version_parts
    server_version = binding.version
    if server_parts is None or server_version is None:
        return
    md = _current_request_metadata.get()
    client_version_bytes = md.get(PROTOCOL_VERSION_KEY) if md is not None else None
    protocol_name = binding.name
    if client_version_bytes is None:
        raise ProtocolVersionError(
            f"VGI client/worker protocol_version mismatch for protocol {protocol_name!r}.\n"
            f"  Client: <not declared>\n"
            f"  Server: {server_version}\n"
            f"  Direction: the client did not send a vgi_rpc.protocol_version "
            f"metadata key. This is either a vgi-rpc framework bug or a "
            f"non-VGI client connecting to a VGI worker.",
            protocol=protocol_name,
            server_version=server_version,
        )
    try:
        client_version = client_version_bytes.decode()
    except UnicodeDecodeError as exc:
        raise ProtocolVersionError(
            f"VGI client/worker protocol_version mismatch for protocol {protocol_name!r}.\n"
            f"  Client: <undecodable bytes>\n"
            f"  Server: {server_version}\n"
            f"  Direction: client sent non-UTF-8 protocol_version metadata.",
            protocol=protocol_name,
            client_version="<undecodable>",
            server_version=server_version,
        ) from exc
    try:
        client_parts = parse_version(client_version)
    except ValueError as exc:
        raise ProtocolVersionError(
            f"VGI client/worker protocol_version mismatch for protocol {protocol_name!r}.\n"
            f"  Client: {client_version}\n"
            f"  Server: {server_version}\n"
            f"  Direction: client sent a malformed protocol_version. "
            f"Expected canonical semver MAJOR.MINOR.PATCH.",
            protocol=protocol_name,
            client_version=client_version,
            server_version=server_version,
        ) from exc
    # Exact major+minor match; patch is ignored.
    if client_parts[:2] == server_parts[:2]:
        return
    if (client_parts[0], client_parts[1]) < (server_parts[0], server_parts[1]):
        direction = (
            f"client is too old; upgrade the VGI extension/client to a "
            f"version supporting protocol_version {server_version}."
        )
    else:
        direction = (
            f"server is too old; upgrade the VGI worker to a version supporting protocol_version {client_version}."
        )
    raise ProtocolVersionError(
        f"VGI client/worker protocol_version mismatch for protocol {protocol_name!r}.\n"
        f"  Client: {client_version}\n"
        f"  Server: {server_version}\n"
        f"  Direction: {direction}",
        protocol=protocol_name,
        client_version=client_version,
        server_version=server_version,
    )

serve

serve(
    transport: RpcTransport,
    *,
    auth: AuthContext | None = None,
    peer_evidence: PeerEvidenceSet | None = None,
    transport_metadata: Mapping[str, Any] | None = None
) -> None

Serve requests until close, optionally under one connection identity.

HTTP installs request-scoped identity in middleware and therefore uses the defaults. Stateful raw transports resolve once after accept and pass an immutable connection snapshot here; every unary call, stream turn, and cancellation hook on that connection then sees the same AuthContext and PeerEvidenceSet.

Source code in vgi_rpc/rpc/_server.py
def serve(
    self,
    transport: RpcTransport,
    *,
    auth: AuthContext | None = None,
    peer_evidence: PeerEvidenceSet | None = None,
    transport_metadata: Mapping[str, Any] | None = None,
) -> None:
    """Serve requests until close, optionally under one connection identity.

    HTTP installs request-scoped identity in middleware and therefore uses
    the defaults. Stateful raw transports resolve once after accept and
    pass an immutable connection snapshot here; every unary call, stream
    turn, and cancellation hook on that connection then sees the same
    ``AuthContext`` and ``PeerEvidenceSet``.
    """
    capabilities: frozenset[str] = frozenset()
    if isinstance(transport, ShmPipeTransport):
        kind = TransportKind.PIPE
        capabilities = frozenset({"shm"})
    elif isinstance(transport, UnixTransport):
        kind = TransportKind.UNIX
    elif isinstance(transport, TcpTransport):
        kind = TransportKind.TCP
    elif isinstance(transport, PipeTransport):
        kind = TransportKind.PIPE
    else:
        kind = TransportKind.PIPE
    self._notify_transport(kind, capabilities)
    transport_token = None
    if auth is not None or peer_evidence is not None or transport_metadata is not None:
        source_auth = auth or _ANONYMOUS
        auth_snapshot = AuthContext(
            domain=source_auth.domain,
            authenticated=source_auth.authenticated,
            principal=source_auth.principal,
            claims=MappingProxyType(dict(source_auth.claims)),
        )
        metadata_snapshot = MappingProxyType(dict(transport_metadata or {}))
        transport_token = _current_transport.set(
            _TransportContext(
                auth=auth_snapshot,
                transport_metadata=metadata_snapshot,
                peer_evidence=peer_evidence,
            )
        )

    # Cache the client's dynamically-advertised SHM segment for the life of
    # the connection so later offset-only request/data batches resolve
    # against it (see _ConnectionShm).
    conn_shm = _ConnectionShm()
    try:
        while True:
            try:
                self.serve_one(transport, shm_cache=conn_shm)
            except (EOFError, StopIteration):
                break
            except (BrokenPipeError, ConnectionResetError, ConnectionAbortedError):
                _logger.debug(
                    "serve loop ending due to broken pipe",
                    exc_info=True,
                    extra={"server_id": self._server_id},
                )
                break
            except pa.ArrowInvalid as exc:
                # PyArrow reports a clean socket EOF while waiting for the
                # next concatenated request as ArrowInvalid. This happens
                # when any client closes a connection after its final RPC;
                # some clients use one RPC per connection, while others
                # reuse it. Logging a traceback for each disconnect can
                # serialize thousands of server threads behind the logging
                # lock. Preserve real malformed-IPC warnings and only
                # downgrade the exact empty-stream diagnostic at this
                # connection-loop boundary.
                message = str(exc)
                if "schema message" in message and "null or length 0" in message:
                    _logger.debug(
                        "serve loop ending at client EOF",
                        extra={"server_id": self._server_id},
                    )
                else:
                    _logger.warning(
                        "serve loop ending due to ArrowInvalid",
                        exc_info=True,
                        extra={"server_id": self._server_id},
                    )
                break
    finally:
        conn_shm.close()
        if transport_token is not None:
            _current_transport.reset(transport_token)

serve_one

serve_one(
    transport: RpcTransport,
    *,
    shm_cache: _ConnectionShm | None = None
) -> None

Handle a single RPC call (any method type) over the given transport.

Protocol-level errors (VersionError, RpcError from missing metadata) are caught, written back as error responses, and the method returns normally so the serve loop can continue.

PARAMETER DESCRIPTION
transport

The transport to read the request from and write the response to.

TYPE: RpcTransport

shm_cache

Per-connection SHM segment cache supplied by :meth:serve. When None (e.g. a direct serve_one call), a client-advertised segment is attached and detached per call instead of being cached across requests.

TYPE: _ConnectionShm | None DEFAULT: None

RAISES DESCRIPTION
ArrowInvalid

If the incoming data is not valid Arrow IPC. An error response is written to transport before raising so the client can read a structured RpcError.

Source code in vgi_rpc/rpc/_server.py
def serve_one(self, transport: RpcTransport, *, shm_cache: _ConnectionShm | None = None) -> None:
    """Handle a single RPC call (any method type) over the given transport.

    Protocol-level errors (``VersionError``, ``RpcError`` from missing
    metadata) are caught, written back as error responses, and the
    method returns normally so the serve loop can continue.

    Args:
        transport: The transport to read the request from and write the
            response to.
        shm_cache: Per-connection SHM segment cache supplied by
            :meth:`serve`. When ``None`` (e.g. a direct ``serve_one``
            call), a client-advertised segment is attached and detached
            per call instead of being cached across requests.

    Raises:
        pa.ArrowInvalid: If the incoming data is not valid Arrow IPC.
            An error response is written to *transport* before raising so
            the client can read a structured ``RpcError``.

    """
    token = _current_request_id.set(_generate_request_id())
    stats = CallStatistics()
    stats_token = _current_call_stats.set(stats)
    md_token = _current_request_metadata.set(None)
    rb_token = _current_request_batch.set(None)
    sid_token = _current_stream_id.set("")
    dynamic_shm: ShmSegment | None = None
    try:
        try:
            # A client may route the (single-row) request batch through the
            # shm side channel. Resolve it with the static ShmPipeTransport
            # segment if present, else the connection-cached segment a prior
            # request advertised, else (cache empty) let _read_request attach
            # the segment named in this request's own metadata.
            static_shm = transport.shm if isinstance(transport, ShmPipeTransport) else None
            cached_shm = shm_cache.segment if shm_cache is not None else None
            method_name, kwargs = _read_request(
                transport.reader,
                self._ipc_validation,
                self._external_config,
                shm=static_shm or cached_shm,
                attach_shm=lambda md: _maybe_attach_shm(md, self._transport_kind),
            )
        except pa.ArrowInvalid as exc:
            with contextlib.suppress(BrokenPipeError, OSError):
                _write_error_stream(
                    transport.writer,
                    _EMPTY_SCHEMA,
                    exc,
                    server_id=self._server_id,
                    include_traceback=self._tracebacks,
                )
            raise
        except (VersionError, RpcError) as exc:
            with contextlib.suppress(BrokenPipeError, OSError):
                _write_error_stream(
                    transport.writer,
                    _EMPTY_SCHEMA,
                    exc,
                    server_id=self._server_id,
                    include_traceback=self._tracebacks,
                )
            return

        # __transport_options__ — framework transport-capability handshake,
        # handled before method dispatch (not a registered method, so it never
        # appears in `methods` / `__describe__`). Capabilities ride as response
        # metadata; the response batch is empty. Always available.
        if method_name == TRANSPORT_OPTIONS_METHOD_NAME:
            caps_md = {
                **worker_transport_metadata(),
                REQUEST_VERSION_KEY: REQUEST_VERSION,
                SERVER_ID_KEY: self._server_id.encode(),
            }
            empty = pa.RecordBatch.from_arrays([], schema=_EMPTY_SCHEMA)
            with new_ipc_stream(transport.writer, _EMPTY_SCHEMA) as writer:
                writer.write_batch(empty, custom_metadata=caps_md)
            auth, transport_md = _get_auth_and_metadata()
            _emit_access_log(
                self.protocol_name,
                TRANSPORT_OPTIONS_METHOD_NAME,
                MethodType.UNARY.value,
                self._server_id,
                auth,
                transport_md,
                0.0,
                "ok",
                stats=stats,
                server_version=self._server_version,
                protocol_hash=self._protocol_hash,
                error_code="",
            )
            return

        try:
            info = self._resolve(method_name)
        except (ProtocolNotSpecifiedError, ProtocolNotSupportedError, MethodNotImplementedError) as exc:
            _write_error_stream(
                transport.writer, _EMPTY_SCHEMA, exc, server_id=self._server_id, include_traceback=self._tracebacks
            )
            return

        # Application-protocol-version gate, against the binding that owns
        # the resolved method. Failure writes a typed error stream and
        # returns, so the serve loop continues.
        try:
            self.gate_version(info)
        except ProtocolVersionError as exc:
            err_schema = info.result_schema if info.method_type == MethodType.UNARY else _EMPTY_SCHEMA
            _write_error_stream(
                transport.writer, err_schema, exc, server_id=self._server_id, include_traceback=self._tracebacks
            )
            return

        # Request validation. Both steps are answered with a typed error
        # stream rather than allowed to propagate: an exception escaping
        # here leaves the client waiting on a reply that is never written,
        # so a bad request reads as a dead connection instead of as the
        # error it is. ``_deserialize_params`` was previously unguarded,
        # and it is the step that raises on caller-supplied data — an
        # unrecognised member of a dictionary-encoded enum reaches
        # ``base[value]`` and raises KeyError.
        try:
            _deserialize_params(kwargs, info.param_types, self._ipc_validation)
            _validate_call_signature(
                info.name,
                kwargs,
                info.param_types,
                info.param_defaults,
                info.params_schema,
            )
            _validate_params(info.name, kwargs, info.param_types)
        except Exception as exc:
            err_schema = info.result_schema if info.method_type == MethodType.UNARY else _EMPTY_SCHEMA
            _write_error_stream(
                transport.writer, err_schema, exc, server_id=self._server_id, include_traceback=self._tracebacks
            )
            return

        # Determine the SHM segment for this call's data plane (resolving
        # input batches, routing output/response batches). Prefer the
        # transport-level segment; else the connection cache, refreshed from
        # any segment name this request advertised; else — when no cache is
        # wired — a per-call attach. ``_maybe_attach_shm`` rejects the
        # dynamic path for HTTP (defence-in-depth; HTTP doesn't invoke
        # ``serve_one``).
        #
        # Crucially, only route the *response* through shm when the client
        # signalled shm for THIS exchange — it sent the request via a shm
        # pointer (SHM_OFFSET_KEY) or advertised the segment
        # (SHM_SEGMENT_NAME_KEY). The C++ client resolves shm responses only
        # for those methods; for plain inline requests (bind, catalog_*,
        # transaction_*) it expects an inline response and reports a
        # shm-routed one as empty.
        req_md = _current_request_metadata.get()
        request_used_shm = req_md is not None and (
            req_md.get(SHM_OFFSET_KEY) is not None or req_md.get(SHM_SEGMENT_NAME_KEY) is not None
        )
        shm = transport.shm if isinstance(transport, ShmPipeTransport) else None
        if shm is None and shm_cache is not None:
            shm_cache.refresh(req_md, self._transport_kind)
            shm = shm_cache.segment if request_used_shm else None
        if shm is None and shm_cache is None:
            dynamic_shm = _maybe_attach_shm(req_md, self._transport_kind)
            shm = dynamic_shm

        if info.method_type == MethodType.UNARY:
            self._serve_unary(transport, info, kwargs, stats=stats, shm=shm)
        elif info.method_type == MethodType.STREAM:
            self._serve_stream(transport, info, kwargs, stats=stats, shm=shm)
    finally:
        if dynamic_shm is not None:
            with contextlib.suppress(BufferError):
                dynamic_shm.close()
        _current_stream_id.reset(sid_token)
        _current_request_batch.reset(rb_token)
        _current_request_metadata.reset(md_token)
        _current_call_stats.reset(stats_token)
        _current_request_id.reset(token)

RpcConnection

RpcConnection

RpcConnection(
    protocol: type[P],
    transport: RpcTransport,
    on_log: Callable[[Message], None] | None = None,
    *,
    external_location: ExternalLocationConfig | None = None,
    ipc_validation: IpcValidation = FULL
)

Context manager that provides a typed RPC proxy over a transport.

The type parameter P is the Protocol class, enabling IDE autocompletion for all methods defined on the protocol::

with RpcConnection(MyProtocol, transport) as svc:
    result = svc.add(a=1, b=2)   # IDE sees MyProtocol methods

Initialize with a protocol type and transport.

Source code in vgi_rpc/rpc/_client.py
def __init__(
    self,
    protocol: type[P],
    transport: RpcTransport,
    on_log: Callable[[Message], None] | None = None,
    *,
    external_location: ExternalLocationConfig | None = None,
    ipc_validation: IpcValidation = IpcValidation.FULL,
) -> None:
    """Initialize with a protocol type and transport."""
    self._protocol = protocol
    self._transport = transport
    self._on_log = on_log
    self._external_config = external_location
    self._ipc_validation = ipc_validation

__enter__

__enter__() -> P

Enter the context and return a typed proxy.

Source code in vgi_rpc/rpc/_client.py
def __enter__(self) -> P:
    """Enter the context and return a typed proxy."""
    if wire_transport_logger.isEnabledFor(logging.DEBUG):
        wire_transport_logger.debug("RpcConnection open: protocol=%s", self._protocol.__name__)
    return cast(
        "P",
        _RpcProxy(
            self._protocol,
            self._transport,
            self._on_log,
            external_config=self._external_config,
            ipc_validation=self._ipc_validation,
        ),
    )

__exit__

__exit__(
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: TracebackType | None,
) -> None

Close the transport.

Source code in vgi_rpc/rpc/_client.py
def __exit__(
    self,
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: TracebackType | None,
) -> None:
    """Close the transport."""
    if wire_transport_logger.isEnabledFor(logging.DEBUG):
        wire_transport_logger.debug("RpcConnection close: protocol=%s", self._protocol.__name__)
    self._transport.close()

RpcTransport

RpcTransport

Bases: Protocol

Bidirectional byte stream transport.

reader property

reader: IOBase

Readable binary stream.

writer property

writer: IOBase

Writable binary stream.

close

close() -> None

Close the transport.

Source code in vgi_rpc/rpc/_transport.py
def close(self) -> None:
    """Close the transport."""
    ...

RpcMethodInfo

RpcMethodInfo dataclass

RpcMethodInfo(
    name: str,
    params_schema: Schema,
    result_schema: Schema,
    result_type: object,
    method_type: MethodType,
    has_return: bool,
    doc: str | None,
    param_defaults: dict[str, object] = dict(),
    param_types: dict[str, object] = dict(),
    param_docs: dict[str, str] = dict(),
    header_type: (
        type[ArrowSerializableDataclass] | None
    ) = None,
    is_exchange: bool | None = None,
    protocol_name: str = "",
)

Metadata for a single RPC method, derived from Protocol type hints.

Produced by :func:rpc_methods when introspecting a Protocol class. Each instance describes one method's wire-protocol details: its Arrow schemas, parameter types and defaults, and the original docstring.

ATTRIBUTE DESCRIPTION
name

Method name as it appears on the Protocol.

TYPE: str

params_schema

Arrow schema for the serialized request parameters.

TYPE: Schema

result_schema

Arrow schema for the serialized response (unary only; empty schema for stream methods).

TYPE: Schema

result_type

The raw Python return-type annotation (e.g. float, Stream[MyState]).

TYPE: object

method_type

Whether this is a UNARY or STREAM call.

TYPE: MethodType

has_return

True when the unary method returns a value (False for -> None or stream methods).

TYPE: bool

doc

The method's docstring from the Protocol class, or None if no docstring was provided.

TYPE: str | None

param_defaults

Mapping of parameter name to default value for parameters that have defaults in the Protocol signature.

TYPE: dict[str, object]

param_types

Mapping of parameter name to its Python type annotation (excludes self and return).

TYPE: dict[str, object]

param_docs

Mapping of parameter name to its documented description, parsed from the Protocol method's docstring Args: section.

TYPE: dict[str, str]

header_type

For stream methods with a header, the concrete ArrowSerializableDataclass subclass for the header. None when the method has no header.

TYPE: type[ArrowSerializableDataclass] | None

is_exchange

For stream methods, True if the state class extends ExchangeState (client must send input), False if it extends ProducerState (server-initiated), None if the state class uses raw StreamState (unknown). Always None for unary methods.

TYPE: bool | None

protocol_name

Name of the Protocol this method belongs to. Makes an RpcMethodInfo self-describing, so a server hosting several protocols can label access-log records, spans and errors with the protocol that owns the resolved method rather than a server-wide default. Empty string when the info was built outside rpc_methods (e.g. a synthetic method).

TYPE: str

MethodType

MethodType

Bases: Enum

Classification of RPC method patterns.

Errors

RpcError

RpcError(
    error_type: str,
    error_message: str,
    remote_traceback: str,
    *,
    request_id: str = "",
    error_code: str = "",
    error_kind: str = "",
    error_details: list[dict[str, Any]] | None = None
)

Bases: Exception

Raised on the client side when the server reports an error.

Carries the three layers of the error model (WIRE_PROTOCOL.md §8):

  • :attr:error_code -- the canonical code's name ("UNAVAILABLE"), or "" when the server sent none (a server older than the model). :attr:code reads it as a :class:~vgi_rpc.errors.Code.
  • :attr:error_kind -- the reason a client branches on, or "".
  • :attr:error_details -- the detail objects as received, unknown types included. The typed accessors (:meth:retry_info, :meth:bad_request, ...) return the catalog entry of that type, ignoring the rest.

:meth:is_retryable classifies; nothing here retries. Automatic retry is opt-in because a method may not be idempotent.

Initialize with error details from the remote side.

PARAMETER DESCRIPTION
error_type

The remote exception's class name.

TYPE: str

error_message

The remote message.

TYPE: str

remote_traceback

The remote traceback, or "" when the server has tracebacks turned off.

TYPE: str

request_id

The request correlation ID, when known.

TYPE: str DEFAULT: ''

error_code

The vgi_rpc.error_code value, or "".

TYPE: str DEFAULT: ''

error_kind

The vgi_rpc.error_kind value, or "".

TYPE: str DEFAULT: ''

error_details

The decoded vgi_rpc.error_details objects.

TYPE: list[dict[str, Any]] | None DEFAULT: None

Source code in vgi_rpc/rpc/_common.py
def __init__(
    self,
    error_type: str,
    error_message: str,
    remote_traceback: str,
    *,
    request_id: str = "",
    error_code: str = "",
    error_kind: str = "",
    error_details: list[dict[str, Any]] | None = None,
) -> None:
    """Initialize with error details from the remote side.

    Args:
        error_type: The remote exception's class name.
        error_message: The remote message.
        remote_traceback: The remote traceback, or ``""`` when the server
            has tracebacks turned off.
        request_id: The request correlation ID, when known.
        error_code: The ``vgi_rpc.error_code`` value, or ``""``.
        error_kind: The ``vgi_rpc.error_kind`` value, or ``""``.
        error_details: The decoded ``vgi_rpc.error_details`` objects.

    """
    self.error_type = error_type
    self.error_message = error_message
    self.remote_traceback = remote_traceback
    self.request_id = request_id
    self.error_code = error_code
    self.error_kind = error_kind
    self.error_details: list[dict[str, Any]] = list(error_details) if error_details else []
    super().__init__(f"{error_type}: {error_message}")

code property

code: Code

The canonical code; UNKNOWN when absent or unrecognised.

is_retryable

is_retryable() -> bool

Whether retrying this call is warranted, by the rule in WIRE_PROTOCOL.md §8.

UNAVAILABLE always; RESOURCE_EXHAUSTED only with RetryInfo. When :meth:retry_info is present a retry waits at least that long.

Source code in vgi_rpc/rpc/_common.py
def is_retryable(self) -> bool:
    """Whether retrying this call is warranted, by the rule in WIRE_PROTOCOL.md §8.

    ``UNAVAILABLE`` always; ``RESOURCE_EXHAUSTED`` only with ``RetryInfo``.
    When :meth:`retry_info` is present a retry waits at least that long.
    """
    return is_retryable(self.code, self.error_details)

details

details() -> list[ErrorDetail]

Return the details this client understands, in wire order; unknown types skipped.

Source code in vgi_rpc/rpc/_common.py
def details(self) -> list[ErrorDetail]:
    """Return the details this client understands, in wire order; unknown types skipped."""
    return [d for d in (parse_error_detail(obj) for obj in self.error_details) if d is not None]

error_info

error_info() -> ErrorInfo | None

Return the vgi_rpc.ErrorInfo detail, if present.

Source code in vgi_rpc/rpc/_common.py
def error_info(self) -> ErrorInfo | None:
    """Return the ``vgi_rpc.ErrorInfo`` detail, if present."""
    return self._detail(ErrorInfo)

retry_info

retry_info() -> RetryInfo | None

Return the vgi_rpc.RetryInfo detail, if present.

Source code in vgi_rpc/rpc/_common.py
def retry_info(self) -> RetryInfo | None:
    """Return the ``vgi_rpc.RetryInfo`` detail, if present."""
    return self._detail(RetryInfo)

bad_request

bad_request() -> BadRequest | None

Return the vgi_rpc.BadRequest detail, if present.

Source code in vgi_rpc/rpc/_common.py
def bad_request(self) -> BadRequest | None:
    """Return the ``vgi_rpc.BadRequest`` detail, if present."""
    return self._detail(BadRequest)

precondition_failure

precondition_failure() -> PreconditionFailure | None

Return the vgi_rpc.PreconditionFailure detail, if present.

Source code in vgi_rpc/rpc/_common.py
def precondition_failure(self) -> PreconditionFailure | None:
    """Return the ``vgi_rpc.PreconditionFailure`` detail, if present."""
    return self._detail(PreconditionFailure)

quota_failure

quota_failure() -> QuotaFailure | None

Return the vgi_rpc.QuotaFailure detail, if present.

Source code in vgi_rpc/rpc/_common.py
def quota_failure(self) -> QuotaFailure | None:
    """Return the ``vgi_rpc.QuotaFailure`` detail, if present."""
    return self._detail(QuotaFailure)

resource_info

resource_info() -> ResourceInfo | None

Return the vgi_rpc.ResourceInfo detail, if present.

Source code in vgi_rpc/rpc/_common.py
def resource_info(self) -> ResourceInfo | None:
    """Return the ``vgi_rpc.ResourceInfo`` detail, if present."""
    return self._detail(ResourceInfo)

help

help() -> Help | None

Return the vgi_rpc.Help detail, if present.

Source code in vgi_rpc/rpc/_common.py
def help(self) -> Help | None:
    """Return the ``vgi_rpc.Help`` detail, if present."""
    return self._detail(Help)

localized_message

localized_message() -> LocalizedMessage | None

Return the vgi_rpc.LocalizedMessage detail, if present.

Source code in vgi_rpc/rpc/_common.py
def localized_message(self) -> LocalizedMessage | None:
    """Return the ``vgi_rpc.LocalizedMessage`` detail, if present."""
    return self._detail(LocalizedMessage)

VersionError

Bases: Exception

Raised when a request has a missing or incompatible protocol version.

ProtocolVersionError

ProtocolVersionError(
    message: str = "",
    *,
    protocol: str = "",
    client_version: str = "",
    server_version: str = ""
)

Bases: VersionError

Raised when the client's declared protocol_version is incompatible with the server's.

Subclass of VersionError so existing catch sites in serve_one (and the HTTP unary/stream wrappers) write a typed error stream and continue. Carries a directional message that tells whoever reads it which side to upgrade.

Build the error, keeping the two versions for the precondition detail.

PARAMETER DESCRIPTION
message

The directional, human-readable explanation.

TYPE: str DEFAULT: ''

protocol

The protocol whose version gate refused the call.

TYPE: str DEFAULT: ''

client_version

What the client declared, or "" when it declared none.

TYPE: str DEFAULT: ''

server_version

What the server's binding declares.

TYPE: str DEFAULT: ''

Source code in vgi_rpc/rpc/_common.py
def __init__(
    self,
    message: str = "",
    *,
    protocol: str = "",
    client_version: str = "",
    server_version: str = "",
) -> None:
    """Build the error, keeping the two versions for the precondition detail.

    Args:
        message: The directional, human-readable explanation.
        protocol: The protocol whose version gate refused the call.
        client_version: What the client declared, or ``""`` when it declared none.
        server_version: What the server's binding declares.

    """
    super().__init__(message)
    self.protocol = protocol
    self.client_version = client_version
    self.server_version = server_version

error_details property

error_details: list[PreconditionFailure]

One protocol_version violation naming the protocol, when known.

ProtocolVersionError (a VersionError subclass) is raised at the dispatch boundary when the client's vgi_rpc.protocol_version differs from the server's on the major+minor components (patch is ignored). It is only active when the Protocol class declares a protocol_version; otherwise the check is disabled. __describe__ is exempt so a mismatched client can still introspect to discover the server's version.

Typed marker errors

These exception classes carry a stable error_kind class attribute that the wire serializer surfaces as the vgi_rpc.error_kind metadata key on the EXCEPTION-level batch. Clients can pattern-match the kind instead of substring-searching the error message.

MethodNotImplementedError

Bases: AttributeError

Raised server-side when no handler is registered for the requested RPC method.

Subclass of AttributeError so existing except AttributeError callers keep working. Carries a stable error_kind class attribute that the wire serializer surfaces as a top-level vgi_rpc.error_kind metadata key on the EXCEPTION-level error batch, so clients can pattern-match on the kind rather than substring-searching the message text.

Used by callers that want to detect "old server doesn't know this method" cleanly (e.g. capability detection + fallback to a legacy RPC method).

SessionLostError

Bases: Exception

Raised server-side when a sticky session token cannot be honoured.

Surfaced over the wire with error_kind="session_lost" so clients can pattern-match the kind without substring-searching the message. Causes include: token presented to a different worker than the one that minted it (server_id mismatch), the registry entry aged out via TTL eviction, AAD mismatch (cross-principal replay), or any other validation failure.

Sticky session machinery is HTTP-only; this error never originates from pipe/unix transports.

ServerDrainingError

ServerDrainingError(
    message: str = "", *, retry_after: float = 1.0
)

Bases: Exception

Raised server-side when a sticky-enabled worker is draining and refuses new sessions.

Surfaced over the wire with error_kind="server_draining". Existing sessions continue to serve through TTL or explicit close; only new ctx.open_session calls are rejected. Operators trigger drain via RpcServer.drain() (typically from a SIGTERM handler) ahead of deploy-time worker rotation.

Build the error.

PARAMETER DESCRIPTION
message

Human-readable explanation.

TYPE: str DEFAULT: ''

retry_after

Seconds before a retry -- which a load balancer will usually route to a worker that is not draining.

TYPE: float DEFAULT: 1.0

Source code in vgi_rpc/rpc/_common.py
def __init__(self, message: str = "", *, retry_after: float = 1.0) -> None:
    """Build the error.

    Args:
        message: Human-readable explanation.
        retry_after: Seconds before a retry -- which a load balancer will
            usually route to a worker that is not draining.

    """
    super().__init__(message)
    self.retry_after = retry_after

error_details property

error_details: list[RetryInfo]

The retry hint.

See also IPCError in the Serialization module.

CallStatistics

CallStatistics dataclass

CallStatistics(
    input_batches: int = 0,
    output_batches: int = 0,
    input_rows: int = 0,
    output_rows: int = 0,
    input_bytes: int = 0,
    output_bytes: int = 0,
)

Mutable accumulator of per-call I/O counters for usage accounting.

Created at dispatch start and populated as batches flow through the server. Surfaced through the access log and OTel dispatch hook.

Byte measurement: uses pa.RecordBatch.get_total_buffer_size() which reports logical Arrow buffer sizes (O(columns), negligible cost). This is an approximation — it does not include IPC framing overhead (padding, schema messages, EOS markers).

ATTRIBUTE DESCRIPTION
input_batches

Number of input batches read by the server.

TYPE: int

output_batches

Number of output batches written by the server.

TYPE: int

input_rows

Total rows across all input batches.

TYPE: int

output_rows

Total rows across all output batches.

TYPE: int

input_bytes

Approximate logical bytes across all input batches.

TYPE: int

output_bytes

Approximate logical bytes across all output batches.

TYPE: int

record_input

record_input(batch: RecordBatch) -> None

Record an input batch's row count and buffer size.

Source code in vgi_rpc/rpc/_common.py
def record_input(self, batch: pa.RecordBatch) -> None:
    """Record an input batch's row count and buffer size."""
    self.input_batches += 1
    self.input_rows += batch.num_rows
    self.input_bytes += batch.get_total_buffer_size()

record_output

record_output(batch: RecordBatch) -> None

Record an output batch's row count and buffer size.

Source code in vgi_rpc/rpc/_common.py
def record_output(self, batch: pa.RecordBatch) -> None:
    """Record an output batch's row count and buffer size."""
    self.output_batches += 1
    self.output_rows += batch.num_rows
    self.output_bytes += batch.get_total_buffer_size()

Convenience Functions

run_server

run_server(
    protocol_or_server: type | RpcServer,
    implementation: object | None = None,
) -> None

Serve RPC requests, defaulting to stdin/stdout pipe transport.

This is the recommended entry point for subprocess workers. Accepts either a (protocol, implementation) pair or a pre-built RpcServer.

The function parses sys.argv and supports the following CLI flags:

  • --http — Serve over HTTP instead of stdin/stdout (requires vgi-rpc[http]).
  • --unix PATH — Serve raw framing over a Unix domain socket.
  • --tcp [HOST:]PORT — Serve raw framing over a TCP socket. Host defaults to 127.0.0.1 (loopback only); PORT may be 0 to auto-select. No auth/TLS — use --http for untrusted networks.
  • --host HOST — HTTP bind address (default 127.0.0.1).
  • --port PORT — HTTP port (default 0, auto-select).
  • --describe — Enable the __describe__ introspection method.
  • --access-log PATH — Append JSONL access log records to PATH. The cross-language conformance contract requires every worker to accept this flag; see docs/access-log-spec.md.
  • --max-response-bytes N — HTTP-only. Cap the outgoing HTTP body of every method response at N bytes (including IPC framing). For producer streams, controls when the framework mints a continuation token to split the response across multiple HTTP turns. Default: no body cap. Env: VGI_RPC_MAX_RESPONSE_BYTES.
  • --max-externalized-response-bytes N — HTTP-only. Cap the total bytes uploaded to external storage during one HTTP response. Default: unbounded. Env: VGI_RPC_MAX_EXTERNALIZED_RESPONSE_BYTES.
  • --max-stream-response-bytes N — Deprecated; alias for --max-response-bytes.

Without --http the server runs over stdin/stdout pipes (the default, suitable for SubprocessTransport).

PARAMETER DESCRIPTION
protocol_or_server

A Protocol class (requires implementation) or an already-constructed RpcServer.

TYPE: type | RpcServer

implementation

The implementation object. Required when protocol_or_server is a Protocol class; must be None when passing an RpcServer.

TYPE: object | None DEFAULT: None

RAISES DESCRIPTION
TypeError

On invalid argument combinations.

Source code in vgi_rpc/rpc/__init__.py
def run_server(protocol_or_server: type | RpcServer, implementation: object | None = None) -> None:
    """Serve RPC requests, defaulting to stdin/stdout pipe transport.

    This is the recommended entry point for subprocess workers.  Accepts
    either a ``(protocol, implementation)`` pair or a pre-built ``RpcServer``.

    The function parses ``sys.argv`` and supports the following CLI flags:

    - ``--http``  — Serve over HTTP instead of stdin/stdout (requires
      ``vgi-rpc[http]``).
    - ``--unix PATH`` — Serve raw framing over a Unix domain socket.
    - ``--tcp [HOST:]PORT`` — Serve raw framing over a TCP socket.  Host
      defaults to ``127.0.0.1`` (loopback only); ``PORT`` may be ``0`` to
      auto-select.  No auth/TLS — use ``--http`` for untrusted networks.
    - ``--host HOST`` — HTTP bind address (default ``127.0.0.1``).
    - ``--port PORT`` — HTTP port (default ``0``, auto-select).
    - ``--describe`` — Enable the ``__describe__`` introspection method.
    - ``--access-log PATH`` — Append JSONL access log records to ``PATH``.
      The cross-language conformance contract requires every worker to
      accept this flag; see ``docs/access-log-spec.md``.
    - ``--max-response-bytes N`` — HTTP-only.  Cap the outgoing HTTP body
      of every method response at ``N`` bytes (including IPC framing).
      For producer streams, controls when the framework mints a
      continuation token to split the response across multiple HTTP
      turns.  Default: no body cap.  Env:
      ``VGI_RPC_MAX_RESPONSE_BYTES``.
    - ``--max-externalized-response-bytes N`` — HTTP-only.  Cap the
      total bytes uploaded to external storage during one HTTP response.
      Default: unbounded.  Env:
      ``VGI_RPC_MAX_EXTERNALIZED_RESPONSE_BYTES``.
    - ``--max-stream-response-bytes N`` — **Deprecated**; alias for
      ``--max-response-bytes``.

    Without ``--http`` the server runs over stdin/stdout pipes (the
    default, suitable for ``SubprocessTransport``).

    Args:
        protocol_or_server: A Protocol class (requires *implementation*) or
            an already-constructed ``RpcServer``.
        implementation: The implementation object.  Required when
            *protocol_or_server* is a Protocol class; must be ``None`` when
            passing an ``RpcServer``.

    Raises:
        TypeError: On invalid argument combinations.

    """
    parser = argparse.ArgumentParser(description="vgi-rpc server")
    transport_mode = parser.add_mutually_exclusive_group()
    transport_mode.add_argument(
        "--http", action="store_true", default=False, help="Serve over HTTP instead of stdin/stdout"
    )
    transport_mode.add_argument(
        "--unix",
        metavar="PATH",
        default=None,
        help=(
            "Bind to this Unix domain socket path instead of stdin/stdout. "
            "Mutually exclusive with --http.  When set, --threaded defaults to True "
            "and --idle-timeout governs self-shutdown."
        ),
    )
    transport_mode.add_argument(
        "--tcp",
        metavar="[HOST:]PORT",
        default=None,
        help=(
            "Serve raw framing over a TCP socket instead of stdin/stdout. "
            "HOST defaults to 127.0.0.1 (loopback only); PORT may be 0 to auto-select. "
            "Mutually exclusive with --http/--unix.  When set, --threaded defaults to True "
            "and --idle-timeout governs self-shutdown.  No auth/TLS — use --http for untrusted networks."
        ),
    )
    parser.add_argument(
        "--idle-timeout",
        type=float,
        default=float(os.environ.get("VGI_RPC_IDLE_TIMEOUT", "300")),
        help=(
            "Self-terminate after this many seconds with zero active connections. "
            "Only meaningful with --unix/--tcp.  A startup-grace period of max(idle_timeout, 60) "
            "protects against shutdown before the first client.  0 disables.  "
            "Env: VGI_RPC_IDLE_TIMEOUT (default 300)."
        ),
    )
    parser.add_argument(
        "--threaded",
        action=argparse.BooleanOptionalAction,
        default=None,
        help=(
            "Serve each socket connection in a separate daemon thread "
            "(only meaningful with --unix/--tcp; default True when either is set)."
        ),
    )
    parser.add_argument(
        "--max-connections",
        type=int,
        default=int(os.environ.get("VGI_RPC_MAX_CONNECTIONS", "64")),
        help=(
            "Maximum concurrent --unix/--tcp connections; excess clients wait in the socket backlog. "
            "Use 0 for unlimited. Env: VGI_RPC_MAX_CONNECTIONS (default 64)."
        ),
    )
    parser.add_argument("--host", default="127.0.0.1", help="HTTP bind address (default: 127.0.0.1)")
    parser.add_argument("--port", type=int, default=0, help="HTTP port (default: auto-select)")
    parser.add_argument(
        "--describe", action="store_true", default=False, help="Enable __describe__ introspection method"
    )
    parser.add_argument(
        "--access-log",
        metavar="PATH",
        default=os.environ.get("VGI_RPC_ACCESS_LOG"),
        help=(
            "Append JSONL access log records to PATH (vgi_rpc.access logger at INFO). "
            "PATH may contain {pid} and {server_id} placeholders. "
            "Env: VGI_RPC_ACCESS_LOG."
        ),
    )
    parser.add_argument(
        "--access-log-max-bytes",
        type=int,
        default=int(os.environ.get("VGI_RPC_ACCESS_LOG_MAX_BYTES", "0")),
        help="Rotate the access log at this size in bytes (0 = no rotation). Env: VGI_RPC_ACCESS_LOG_MAX_BYTES.",
    )
    parser.add_argument(
        "--access-log-backup-count",
        type=int,
        default=int(os.environ.get("VGI_RPC_ACCESS_LOG_BACKUP_COUNT", "5")),
        help="Number of rotated access-log files to retain. Env: VGI_RPC_ACCESS_LOG_BACKUP_COUNT.",
    )
    parser.add_argument(
        "--access-log-when",
        default=os.environ.get("VGI_RPC_ACCESS_LOG_WHEN"),
        help=(
            "Time-based rotation interval (e.g. 'H', 'D', 'midnight'); mutually exclusive with "
            "--access-log-max-bytes. Env: VGI_RPC_ACCESS_LOG_WHEN."
        ),
    )
    parser.add_argument(
        "--access-log-max-record-bytes",
        type=int,
        default=int(os.environ.get("VGI_RPC_ACCESS_LOG_MAX_RECORD_BYTES", "1048576")),
        help=(
            "Maximum size in bytes of one access-log record (default 1048576 = 1 MiB). "
            "Env: VGI_RPC_ACCESS_LOG_MAX_RECORD_BYTES."
        ),
    )
    parser.add_argument(
        "--access-log-async",
        action="store_true",
        default=os.environ.get("VGI_RPC_ACCESS_LOG_ASYNC", "").lower() in ("1", "true", "yes"),
        help=(
            "Write access-log records from a background thread so disk latency stays "
            "out of the request path. The queue is bounded and full means drop; the "
            "next record through carries 'dropped_records'. Trades durability: a crash "
            "loses queued records. Env: VGI_RPC_ACCESS_LOG_ASYNC."
        ),
    )
    parser.add_argument(
        "--access-log-queue-size",
        type=int,
        default=int(os.environ.get("VGI_RPC_ACCESS_LOG_QUEUE_SIZE", "10000")),
        help="Bound on the async access-log queue (default 10000). Env: VGI_RPC_ACCESS_LOG_QUEUE_SIZE.",
    )
    parser.add_argument(
        "--access-log-sample",
        type=float,
        default=float(os.environ.get("VGI_RPC_ACCESS_LOG_SAMPLE", "1.0")),
        help=(
            "Fraction of successful calls to log, 0.0-1.0 (default 1.0 = all). "
            "Errors are always logged. The decision is keyed on stream/request id, "
            "so every record of one stream shares it. Sampled records carry "
            "'sample_rate' so counts can be scaled. Env: VGI_RPC_ACCESS_LOG_SAMPLE."
        ),
    )
    parser.add_argument(
        "--max-response-bytes",
        type=int,
        default=int(os.environ.get("VGI_RPC_MAX_RESPONSE_BYTES", "0")) or None,
        help=(
            "HTTP-only. Cap decoded Arrow IPC bytes for every method response "
            "(including framing, before HTTP content coding). Overshoot is a "
            "structured error for unary and stream responses. Default: no body cap. "
            "Env: VGI_RPC_MAX_RESPONSE_BYTES."
        ),
    )
    parser.add_argument(
        "--max-externalized-response-bytes",
        type=int,
        default=int(os.environ.get("VGI_RPC_MAX_EXTERNALIZED_RESPONSE_BYTES", "0")) or None,
        help=(
            "HTTP-only.  Cap the total bytes uploaded to external storage "
            "during one HTTP response (one producer turn or one unary/exchange "
            "call).  Default: unbounded.  Env: "
            "VGI_RPC_MAX_EXTERNALIZED_RESPONSE_BYTES."
        ),
    )
    parser.add_argument(
        "--max-stream-response-bytes",
        type=int,
        default=int(os.environ.get("VGI_RPC_MAX_STREAM_RESPONSE_BYTES", "0")) or None,
        help=("Deprecated alias for --max-response-bytes.  Env: VGI_RPC_MAX_STREAM_RESPONSE_BYTES (deprecated)."),
    )
    args = parser.parse_args()

    # Deprecation alias
    if args.max_stream_response_bytes is not None:
        if args.max_response_bytes is not None:
            raise SystemExit("Pass either --max-response-bytes or --max-stream-response-bytes, not both")
        print(
            "warning: --max-stream-response-bytes is deprecated; use --max-response-bytes instead.",
            file=sys.stderr,
        )
        args.max_response_bytes = args.max_stream_response_bytes

    if args.access_log_max_bytes and args.access_log_when:
        raise SystemExit("--access-log-max-bytes and --access-log-when are mutually exclusive")

    if not 0.0 <= args.access_log_sample <= 1.0:
        # Fail at startup rather than at the first request: a rate of 100
        # meaning "100%" would otherwise silently log everything, and a
        # negative one would silently log nothing.
        raise SystemExit(f"--access-log-sample must be between 0.0 and 1.0, got {args.access_log_sample}")

    if args.unix is not None or args.tcp is not None:
        raw_flag = "--unix" if args.unix is not None else "--tcp"
        http_only_violations: list[str] = []
        if args.access_log:
            http_only_violations.append("--access-log")
        if args.max_response_bytes is not None:
            http_only_violations.append("--max-response-bytes")
        if args.max_externalized_response_bytes is not None:
            http_only_violations.append("--max-externalized-response-bytes")
        if http_only_violations:
            raise SystemExit(
                f"{', '.join(http_only_violations)} {'are' if len(http_only_violations) > 1 else 'is'} "
                f"HTTP-only and cannot be combined with {raw_flag}"
            )

    if isinstance(protocol_or_server, RpcServer):
        if implementation is not None:
            raise TypeError("implementation must be None when passing an RpcServer")
        server = protocol_or_server
    elif isinstance(protocol_or_server, type):
        if implementation is None:
            raise TypeError("implementation is required when passing a Protocol class")
        server = RpcServer(protocol_or_server, implementation, enable_describe=args.describe)
    else:
        raise TypeError(f"Expected a Protocol class or RpcServer, got {type(protocol_or_server).__name__}")

    if args.access_log:
        _configure_access_log(
            path=args.access_log,
            max_bytes=args.access_log_max_bytes,
            backup_count=args.access_log_backup_count,
            when=args.access_log_when,
            max_record_bytes=args.access_log_max_record_bytes,
            server_id=server.server_id,
            sample_rate=args.access_log_sample,
            use_async=args.access_log_async,
            queue_size=args.access_log_queue_size,
        )

    if args.http:
        try:
            from vgi_rpc.http import serve_http
        except ImportError:
            print("HTTP transport requires vgi-rpc[http]: pip install vgi-rpc[http]", file=sys.stderr)
            sys.exit(1)
        serve_http(
            server,
            host=args.host,
            port=args.port,
            max_response_bytes=args.max_response_bytes,
            max_externalized_response_bytes=args.max_externalized_response_bytes,
        )
    elif args.unix is not None:
        threaded = True if args.threaded is None else args.threaded
        idle_timeout: float | None = args.idle_timeout if args.idle_timeout > 0 else None
        max_connections: int | None = args.max_connections if args.max_connections > 0 else None

        if sys.platform == "win32":
            # CPython has no AF_UNIX on Windows; --unix carries a named-pipe name
            # (\\.\pipe\...) which the launcher constructs. Serve over a named
            # pipe and advertise it with the PIPE: discovery prefix (the C++
            # launcher matches on that prefix per docs/launcher-protocol.md).
            pipe_name = args.unix

            def _emit_discovery_line(bound_path: str) -> None:
                print(f"PIPE:{bound_path}", flush=True)

            serve_named_pipe(
                server,
                pipe_name,
                threaded=threaded,
                max_connections=max_connections,
                idle_timeout=idle_timeout,
                on_bound=_emit_discovery_line,
            )
        else:
            absolute_path = os.path.abspath(args.unix)

            def _emit_discovery_line(bound_path: str) -> None:
                # Mirrors the PORT:<n> convention used by the HTTP transport.  After
                # this line the worker MUST NOT write further data to stdout — see
                # the cross-language launcher contract.
                print(f"UNIX:{bound_path}", flush=True)

            serve_unix(
                server,
                absolute_path,
                threaded=threaded,
                max_connections=max_connections,
                idle_timeout=idle_timeout,
                on_bound=_emit_discovery_line,
            )
    elif args.tcp is not None:
        threaded = True if args.threaded is None else args.threaded
        idle_timeout = args.idle_timeout if args.idle_timeout > 0 else None
        max_connections = args.max_connections if args.max_connections > 0 else None

        if ":" in args.tcp:
            host_part, _, port_part = args.tcp.rpartition(":")
            tcp_host = host_part or "127.0.0.1"
        else:
            tcp_host = "127.0.0.1"
            port_part = args.tcp
        try:
            tcp_port = int(port_part)
        except ValueError:
            raise SystemExit(f"--tcp expects [HOST:]PORT, got {args.tcp!r}") from None

        def _emit_tcp_discovery_line(bound_host: str, bound_port: int) -> None:
            # Mirrors the UNIX:<path> / PORT:<n> discovery conventions.  After
            # this line the worker MUST NOT write further data to stdout — see
            # the cross-language launcher contract.
            print(f"TCP:{bound_host}:{bound_port}", flush=True)

        serve_tcp(
            server,
            tcp_host,
            tcp_port,
            threaded=threaded,
            max_connections=max_connections,
            idle_timeout=idle_timeout,
            on_bound=_emit_tcp_discovery_line,
        )
    else:
        serve_stdio(server)

connect

connect(
    protocol: type[P],
    cmd: list[str],
    *,
    on_log: Callable[[Message], None] | None = None,
    external_location: ExternalLocationConfig | None = None,
    stderr: StderrMode = INHERIT,
    stderr_logger: Logger | None = None,
    ipc_validation: IpcValidation = FULL
) -> Iterator[P]

Connect to a subprocess RPC server.

Context manager that spawns a subprocess, yields a typed proxy, and cleans up on exit.

PARAMETER DESCRIPTION
protocol

The Protocol class defining the RPC interface.

TYPE: type[P]

cmd

Command to spawn the subprocess worker.

TYPE: list[str]

on_log

Optional callback for log messages from the server.

TYPE: Callable[[Message], None] | None DEFAULT: None

external_location

Optional ExternalLocation configuration for resolving and producing externalized batches.

TYPE: ExternalLocationConfig | None DEFAULT: None

stderr

How to handle the child's stderr stream (see :class:StderrMode).

TYPE: StderrMode DEFAULT: INHERIT

stderr_logger

Logger for StderrMode.PIPE output; ignored for other modes. Defaults to logging.getLogger("vgi_rpc.subprocess.stderr").

TYPE: Logger | None DEFAULT: None

ipc_validation

Validation level for incoming IPC batches.

TYPE: IpcValidation DEFAULT: FULL

YIELDS DESCRIPTION
P

A typed RPC proxy supporting all methods defined on protocol.

Source code in vgi_rpc/rpc/__init__.py
@contextlib.contextmanager
def connect[P](
    protocol: type[P],
    cmd: list[str],
    *,
    on_log: Callable[[Message], None] | None = None,
    external_location: ExternalLocationConfig | None = None,
    stderr: StderrMode = StderrMode.INHERIT,
    stderr_logger: logging.Logger | None = None,
    ipc_validation: IpcValidation = IpcValidation.FULL,
) -> Iterator[P]:
    """Connect to a subprocess RPC server.

    Context manager that spawns a subprocess, yields a typed proxy, and
    cleans up on exit.

    Args:
        protocol: The Protocol class defining the RPC interface.
        cmd: Command to spawn the subprocess worker.
        on_log: Optional callback for log messages from the server.
        external_location: Optional ExternalLocation configuration for
            resolving and producing externalized batches.
        stderr: How to handle the child's stderr stream (see :class:`StderrMode`).
        stderr_logger: Logger for ``StderrMode.PIPE`` output; ignored for
            other modes.  Defaults to
            ``logging.getLogger("vgi_rpc.subprocess.stderr")``.
        ipc_validation: Validation level for incoming IPC batches.

    Yields:
        A typed RPC proxy supporting all methods defined on *protocol*.

    """
    transport = SubprocessTransport(cmd, stderr=stderr, stderr_logger=stderr_logger)
    try:
        with RpcConnection(
            protocol, transport, on_log=on_log, external_location=external_location, ipc_validation=ipc_validation
        ) as rpc_proxy:
            yield rpc_proxy
    finally:
        transport.close()

serve_pipe

serve_pipe(
    protocol: type[P],
    implementation: object,
    *,
    on_log: Callable[[Message], None] | None = None,
    external_location: ExternalLocationConfig | None = None,
    ipc_validation: IpcValidation | None = None
) -> Iterator[P]

Start an in-process pipe server and yield a typed client proxy.

Useful for tests and demos — no subprocess needed. A background thread runs RpcServer.serve() on the server side of a pipe pair.

PARAMETER DESCRIPTION
protocol

The Protocol class defining the RPC interface.

TYPE: type[P]

implementation

The implementation object.

TYPE: object

on_log

Optional callback for log messages from the server.

TYPE: Callable[[Message], None] | None DEFAULT: None

external_location

Optional ExternalLocation configuration for resolving and producing externalized batches.

TYPE: ExternalLocationConfig | None DEFAULT: None

ipc_validation

Validation level for incoming IPC batches. When None (the default), both components use IpcValidation.FULL.

TYPE: IpcValidation | None DEFAULT: None

YIELDS DESCRIPTION
P

A typed RPC proxy supporting all methods defined on protocol.

Source code in vgi_rpc/rpc/__init__.py
@contextlib.contextmanager
def serve_pipe[P](
    protocol: type[P],
    implementation: object,
    *,
    on_log: Callable[[Message], None] | None = None,
    external_location: ExternalLocationConfig | None = None,
    ipc_validation: IpcValidation | None = None,
) -> Iterator[P]:
    """Start an in-process pipe server and yield a typed client proxy.

    Useful for tests and demos — no subprocess needed.  A background thread
    runs ``RpcServer.serve()`` on the server side of a pipe pair.

    Args:
        protocol: The Protocol class defining the RPC interface.
        implementation: The implementation object.
        on_log: Optional callback for log messages from the server.
        external_location: Optional ExternalLocation configuration for
            resolving and producing externalized batches.
        ipc_validation: Validation level for incoming IPC batches.
            When ``None`` (the default), both components use
            ``IpcValidation.FULL``.

    Yields:
        A typed RPC proxy supporting all methods defined on *protocol*.

    """
    client_transport, server_transport = make_pipe_pair()
    server = RpcServer(
        protocol,
        implementation,
        external_location=external_location,
        ipc_validation=ipc_validation if ipc_validation is not None else IpcValidation.FULL,
    )
    thread = threading.Thread(target=server.serve, args=(server_transport,), daemon=True)
    thread.start()
    try:
        with RpcConnection(
            protocol,
            client_transport,
            on_log=on_log,
            external_location=external_location,
            ipc_validation=ipc_validation if ipc_validation is not None else IpcValidation.FULL,
        ) as proxy:
            yield proxy
    finally:
        client_transport.close()
        thread.join(timeout=5)
        server_transport.close()

describe_rpc

describe_rpc(
    protocol: type,
    *,
    methods: Mapping[str, RpcMethodInfo] | None = None
) -> str

Return a human-readable description of an RPC protocol's methods.

Source code in vgi_rpc/rpc/__init__.py
def describe_rpc(protocol: type, *, methods: Mapping[str, RpcMethodInfo] | None = None) -> str:
    """Return a human-readable description of an RPC protocol's methods."""
    if methods is None:
        methods = rpc_methods(protocol)
    lines: list[str] = [f"RPC Protocol: {protocol.__name__}", ""]

    for name, info in sorted(methods.items()):
        lines.append(f"  {name}({info.method_type.value})")
        lines.append(f"    params: {info.params_schema}")
        if info.method_type == MethodType.UNARY:
            lines.append(f"    result: {info.result_schema}")
        if info.doc:
            lines.append(f"    doc: {info.doc.strip()}")
        lines.append("")

    return "\n".join(lines)

Constants

REQUEST_VERSION module-attribute

REQUEST_VERSION = b'1'