Skip to content

External Storage

When RPC payloads grow large — multi-megabyte query results, bulk data imports, large model artifacts — they can exceed the capacity of the underlying transport. HTTP servers and reverse proxies typically enforce request and response body size limits (e.g., 10 MB). Even pipe-based transports become inefficient when serializing very large batches inline.

External storage solves this by transparently offloading oversized Arrow IPC batches to cloud storage (S3, GCS, or any compatible backend). The batch is replaced with a lightweight pointer batch — a zero-row batch carrying a download URL in its metadata. The receiving side resolves the pointer automatically, fetching the actual data in parallel chunks.

This works in both directions:

  • Large outputs — the server externalizes result batches that exceed a size threshold
  • Large inputs — the client uploads input data to a pre-signed URL and sends a pointer in the RPC request

From the caller's perspective, nothing changes. The proxy returns the same typed results; the externalization is invisible.

How It Works

Large Outputs (Server to Client)

When a server method returns a batch that exceeds externalize_threshold_bytes, the framework intercepts the response before writing it to the wire:

  1. The batch is serialized to Arrow IPC format
  2. Optionally compressed with zstd
  3. Uploaded to cloud storage via the configured ExternalStorage backend
  4. The storage backend returns a download URL (typically a pre-signed URL with an expiry)
  5. A pointer batch replaces the original — a zero-row batch with vgi_rpc.location metadata containing the URL
  6. The pointer batch is written to the wire instead of the full data

On the client side, pointer batches are detected and resolved transparently:

  1. The client reads the pointer batch and extracts the URL from vgi_rpc.location metadata
  2. A HEAD probe determines the content size and whether the server supports range requests
  3. For large payloads (>64 MB by default), parallel range-request fetching splits the download into chunks with speculative hedging for slow chunks
  4. For smaller payloads, a single GET request fetches the data
  5. If the data was compressed, it is decompressed automatically
  6. The Arrow IPC stream is parsed back into a RecordBatch and returned to the caller

This works for both unary results and streaming outputs. For streaming, all batches in a single output cycle (including any log batches) are serialized as one IPC stream, uploaded together, and resolved as a unit.

Large Inputs (Client to Server)

For the reverse direction — when a client needs to send large data to the server — the framework provides upload URLs. This is particularly important for HTTP transport where request body size limits are common.

  1. The client calls request_upload_urls() to get one or more pre-signed URL pairs from the server
  2. Each UploadUrl contains an upload_url (for PUT) and a download_url (for GET) — these may be the same or different depending on the storage backend
  3. The client serializes the large batch to Arrow IPC and uploads it directly to the upload_url via HTTP PUT, bypassing the RPC server entirely
  4. The client sends the RPC request with a pointer batch referencing the download_url
  5. The server resolves the pointer batch transparently, fetching the data from storage

This pattern lets clients send arbitrarily large payloads even when the HTTP server enforces strict request size limits — the large data goes directly to cloud storage, and only a small pointer crosses the RPC boundary.

The upload URL endpoint is enabled by passing upload_url_provider to make_wsgi_app(). The server advertises this capability via the VGI-Upload-URL-Support: true response header, and clients can discover it with http_capabilities().

Server Configuration

S3Storage requires pip install vgi-rpc[s3] (and GCSStorage requires [gcs]); the symbols are only exported from vgi_rpc when their backend dependency is installed.

from vgi_rpc import Compression, ExternalLocationConfig, FetchConfig, RpcServer, S3Storage

storage = S3Storage(bucket="my-bucket", prefix="rpc-data/")
config = ExternalLocationConfig(
    storage=storage,
    externalize_threshold_bytes=1_048_576,  # 1 MiB — batches above this go to S3
    compression=Compression(),              # zstd level 3 by default
    fetch_config=FetchConfig(
        chunk_size_bytes=8 * 1024 * 1024,   # 8 MiB parallel chunks
        max_parallel_requests=8,
    ),
)

server = RpcServer(MyService, MyServiceImpl(), external_location=config)

The same ExternalLocationConfig is passed to the client so it can resolve pointer batches in responses:

from vgi_rpc import connect

with connect(MyService, ["python", "worker.py"], external_location=config) as proxy:
    result = proxy.large_query()  # pointer batches resolved transparently

For HTTP transport with upload URL support:

from vgi_rpc import make_wsgi_app

app = make_wsgi_app(
    server,
    upload_url_provider=storage,       # enables __upload_url__ endpoint
    max_upload_bytes=100 * 1024 * 1024,  # advertise 100 MiB max upload
)

Client Upload URLs

from vgi_rpc import UploadUrl, request_upload_urls

# Get pre-signed upload URLs from the server
urls: list[UploadUrl] = request_upload_urls(
    base_url="http://localhost:8080",
    count=1,
)

# Upload large data directly to storage (bypasses RPC server)
import httpx2
ipc_data = serialize_large_batch(my_batch)
httpx2.put(urls[0].upload_url, content=ipc_data)

# Send RPC request with pointer to uploaded data
# The server resolves the pointer transparently

Clients can discover whether the server supports upload URLs:

from vgi_rpc import http_capabilities

caps = http_capabilities("http://localhost:8080")
if caps.upload_url_support:
    urls = request_upload_urls("http://localhost:8080", count=5)

Compression

When compression is enabled, batches are compressed with zstd before upload and decompressed transparently on fetch:

from vgi_rpc import Compression

# Default: zstd level 3
config = ExternalLocationConfig(
    storage=storage,
    compression=Compression(),  # algorithm="zstd", level=3
)

# Custom compression level (higher = smaller but slower)
config = ExternalLocationConfig(
    storage=storage,
    compression=Compression(level=9),
)

# gzip — stdlib zlib, the right choice when consumers can't do zstd
config = ExternalLocationConfig(
    storage=storage,
    compression=Compression(algorithm="gzip"),  # level defaults to 6
)

# Disable compression
config = ExternalLocationConfig(
    storage=storage,
    compression=None,
)

Compression(algorithm=...) accepts "zstd" (default; needs zstandard on both writer and reader) or "gzip" (stdlib zlib, no extra dependency). The storage backend stores the Content-Encoding header alongside the object; on fetch the Content-Encoding: zstd / gzip header triggers automatic decompression.

zstd requires pip install vgi-rpc[external] (installs zstandard); gzip works with no extra dependency.

Parallel Fetching

For large externalized batches, the client fetches data in parallel chunks using range requests. This significantly reduces download time for multi-megabyte payloads:

  • A HEAD probe determines the total size and whether the server supports Accept-Ranges: bytes
  • If the payload exceeds parallel_threshold_bytes (default 64 MB) and ranges are supported, it is split into chunk_size_bytes chunks (default 8 MB) fetched concurrently
  • Speculative hedging: if a chunk takes longer than speculative_retry_multiplier times the median chunk time, a hedge request is launched in parallel — the first response wins
  • For smaller payloads or servers without range support, a single GET request is used
from vgi_rpc import FetchConfig

fetch_config = FetchConfig(
    parallel_threshold_bytes=64 * 1024 * 1024,  # 64 MiB
    chunk_size_bytes=8 * 1024 * 1024,            # 8 MiB chunks
    max_parallel_requests=8,                      # concurrent fetches
    timeout_seconds=60.0,                         # overall deadline
    max_fetch_bytes=256 * 1024 * 1024,            # 256 MiB hard cap
    speculative_retry_multiplier=2.0,             # hedge at 2x median
    max_speculative_hedges=4,                     # max hedge requests
)

Pre-published References

Per-call externalization re-serializes, re-compresses and re-uploads a result on every call. When a large result rarely changes — a worker's whole catalog, say — publish it once with publish_external(), keep the returned ExternalRef, and return that ref from the unary method on later calls. The server writes the pointer batch directly: no serialization, compression or upload happens during the call.

import threading
from typing import Protocol

import pyarrow as pa

from vgi_rpc import Compression, ExternalRef, ExternalStorage, publish_external, rpc_methods


class CatalogService(Protocol):
    def catalog(self) -> str: ...


class CatalogServiceImpl:
    def __init__(self, storage: ExternalStorage, compression: Compression | None) -> None:
        self._storage, self._compression = storage, compression
        self._ref: ExternalRef | None = None
        self._lock = threading.Lock()

    def catalog(self) -> str | ExternalRef:
        with self._lock:
            if self._ref is None:  # publish once per process / catalog version
                schema = rpc_methods(CatalogService)["catalog"].result_schema
                batch = pa.RecordBatch.from_pydict({"result": [build_catalog()]}, schema=schema)
                self._ref = publish_external(batch, self._storage, self._compression)
            return self._ref

A unary method may return an ExternalRef instead of its declared value. The return annotation may say so (-> str | ExternalRef); ExternalRef is stripped when the result schema is derived, so the wire contract (and protocol hash) is that of str. Returning a ref:

  • always answers with a pointer — whether or not the server has external storage configured, and regardless of externalize_threshold_bytes (a ref is never inlined, nor routed through shared memory);
  • skips building and validating the result value;
  • is not counted toward max_externalized_response_bytes (nothing is uploaded during the call); the tiny pointer still goes through the wire-body budget;
  • works on every transport (pipe, subprocess, Unix, TCP, HTTP). Stream methods are not supported.

publish_external(batch, storage, compression=None, *, include_sha256=True) serializes the 1-row result batch exactly as the per-call externalizer would, hashes the raw IPC bytes, compresses if asked, and calls storage.upload() once. With include_sha256=False the ref carries no digest, so clients skip the content check — for an object rewritten in place, or one too large to hash. Clients need no change: a ref's pointer is indistinguishable from any other.

The caller owns the ref and the object's lifecycle:

  • A long-lived ref must not point at an object under the short-TTL lifecycle rule used for per-call uploads (see Object Lifecycle) — publish under a different prefix.
  • A pre-signed URL expires; re-sign or rebuild the ref before it does.
  • Only return a ref to callers who are all entitled to the same content.

Object Lifecycle

Uploaded objects persist indefinitely — vgi-rpc does not delete them. Configure storage-level cleanup policies to auto-expire old data.

S3 lifecycle rule:

aws s3api put-bucket-lifecycle-configuration \
  --bucket MY_BUCKET \
  --lifecycle-configuration '{
    "Rules": [{
      "ID": "expire-vgi-rpc",
      "Filter": {"Prefix": "rpc-data/"},
      "Status": "Enabled",
      "Expiration": {"Days": 1}
    }]
  }'

GCS lifecycle rule:

gsutil lifecycle set <(cat <<EOF
{"rule": [{"action": {"type": "Delete"},
           "condition": {"age": 1, "matchesPrefix": ["rpc-data/"]}}]}
EOF
) gs://MY_BUCKET

URL Validation

By default, external location URLs are validated with https_only_validator, which rejects non-HTTPS URLs. This prevents pointer batches from being used to probe internal networks. You can provide a custom validator via ExternalLocationConfig.url_validator.

API Reference

Configuration

ExternalLocationConfig module-attribute

ExternalLocationConfig = ServerExternalConfig

Compression dataclass

Compression(
    algorithm: Literal["zstd", "gzip"] = "zstd",
    level: int = 3,
)

Compression settings for externalized data.

ATTRIBUTE DESCRIPTION
algorithm

Compression algorithm. "zstd" (default) needs zstandard on both the writer and the reader; "gzip" uses stdlib zlib and is the right choice when consumers (e.g. browsers without a zstd polyfill, generic HTTP tooling) can't do zstd.

TYPE: Literal['zstd', 'gzip']

level

Compression level. Codec-specific — 1-22 for zstd (default 3), 1-9 for gzip (default 6 when level is left at the zstd-shaped default of 3, since gzip-3 produces noticeably worse ratios).

TYPE: int

FetchConfig dataclass

FetchConfig(
    parallel_threshold_bytes: int = 64 * 1024 * 1024,
    chunk_size_bytes: int = 8 * 1024 * 1024,
    max_parallel_requests: int = 8,
    timeout_seconds: float = 60.0,
    max_fetch_bytes: int = 256 * 1024 * 1024,
    max_decompressed_bytes: int | None = None,
    max_redirects: int = 5,
    speculative_retry_multiplier: float = 2.0,
    max_speculative_hedges: int = 4,
)

Configuration for parallel range-request fetching.

Maintains a persistent aiohttp.ClientSession backed by a daemon thread. Use as a context manager or call close() to release resources.

ATTRIBUTE DESCRIPTION
parallel_threshold_bytes

Below this size, use a single GET.

TYPE: int

chunk_size_bytes

Size of each Range request chunk.

TYPE: int

max_parallel_requests

Semaphore limit on concurrent requests.

TYPE: int

timeout_seconds

Overall deadline for the fetch.

TYPE: float

max_fetch_bytes

Hard cap on encoded/on-wire download size.

TYPE: int

max_decompressed_bytes

Hard cap after content decoding. None preserves the historical effective limit of 16 * max_fetch_bytes.

TYPE: int | None

max_redirects

Maximum redirects followed by each probe or data request. Every target is passed through the caller's URL validator before it is requested. Set to 0 to reject redirects.

TYPE: int

speculative_retry_multiplier

Launch a hedge request for chunks taking longer than multiplier * median. Set to 0 to disable hedging.

TYPE: float

max_speculative_hedges

Maximum number of hedge requests per fetch. Prevents runaway request amplification under adversarial or slow-backend conditions. 0 means no limit (bounded only by chunk count).

TYPE: int

close

close() -> None

Close the pooled session, stop the event loop, and join the thread.

Source code in vgi_rpc/external_fetch.py
def close(self) -> None:
    """Close the pooled session, stop the event loop, and join the thread."""
    pool = self._pool
    with pool.lock:
        loop = pool.loop
        thread = pool.thread
        if pool.session is not None and loop is not None and not loop.is_closed():
            asyncio.run_coroutine_threadsafe(pool.session.close(), loop).result(timeout=5)
        pool.session = None
        if loop is not None and not loop.is_closed():
            loop.call_soon_threadsafe(loop.stop)
        if thread is not None:
            thread.join(timeout=5)
            if thread.is_alive():
                raise RuntimeError("external fetch event-loop thread did not stop within 5s")
        # ``loop.stop()`` and joining its thread do not release the event
        # loop's selector or self-pipe descriptors.  Close it explicitly
        # before dropping the last owned reference; repeated short-lived
        # client configs otherwise exhaust macOS's common 256-FD limit.
        if loop is not None and not loop.is_closed():
            loop.close()
        pool.loop = None
        pool.thread = None

__del__

__del__() -> None

Safety net: close pool on garbage collection.

Source code in vgi_rpc/external_fetch.py
def __del__(self) -> None:
    """Safety net: close pool on garbage collection."""
    with contextlib.suppress(Exception):
        self.close()

__enter__

__enter__() -> FetchConfig

Enter context manager.

Source code in vgi_rpc/external_fetch.py
def __enter__(self) -> FetchConfig:
    """Enter context manager."""
    return self

__exit__

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

Exit context manager, closing the pool.

Source code in vgi_rpc/external_fetch.py
def __exit__(
    self,
    exc_type: type[BaseException] | None,
    exc_val: BaseException | None,
    exc_tb: TracebackType | None,
) -> None:
    """Exit context manager, closing the pool."""
    self.close()

Pre-published References

ExternalRef dataclass

ExternalRef(url: str, sha256: str | None = None)

A reference to an already-published unary result.

A unary method may return an ExternalRef in place of its declared result value. The server then answers with the ExternalLocation pointer batch for url directly — no result serialization, compression, or upload happens during the call, and the ref is used whether or not the server has external storage configured and regardless of externalize_threshold_bytes. Clients resolve it like any other pointer, so they need no change.

Build one with :func:publish_external (or by hand for an object published out of band). The object at url must be an Arrow IPC stream (optionally Content-Encoding-compressed) whose schema is the method's result schema and which holds exactly one 1-row data batch.

The caller owns caching the ref and the object's lifecycle: a long-lived ref must not point at an object under the short-TTL lifecycle rule used for per-call uploads, and a pre-signed URL expires — re-sign or rebuild the ref before then. Only return a ref to callers who are all entitled to the same content.

ATTRIBUTE DESCRIPTION
url

Where the published IPC stream lives.

TYPE: str

sha256

Lowercase hex SHA-256 of the raw (pre-compression) IPC stream bytes, sent as vgi_rpc.location.sha256. None omits the key, so clients skip the content check — use this for an object rewritten in place or one too large to hash.

TYPE: str | None

__post_init__

__post_init__() -> None

Validate the URL and digest shape.

RAISES DESCRIPTION
ValueError

If url is empty or sha256 is not 64 lowercase hex characters.

Source code in vgi_rpc/external.py
def __post_init__(self) -> None:
    """Validate the URL and digest shape.

    Raises:
        ValueError: If ``url`` is empty or ``sha256`` is not 64
            lowercase hex characters.

    """
    if not self.url:
        raise ValueError("ExternalRef.url must be non-empty")
    if self.sha256 is not None and (
        len(self.sha256) != 64 or any(c not in "0123456789abcdef" for c in self.sha256)
    ):
        raise ValueError("ExternalRef.sha256 must be 64 lowercase hex characters (or None)")

pointer_batch

pointer_batch(
    schema: Schema,
) -> tuple[RecordBatch, KeyValueMetadata]

Build the zero-row pointer batch announcing this ref.

PARAMETER DESCRIPTION
schema

The method's result schema.

TYPE: Schema

RETURNS DESCRIPTION
RecordBatch

(batch, custom_metadata) as from

KeyValueMetadata

func:make_external_location_batch.

Source code in vgi_rpc/external.py
def pointer_batch(self, schema: pa.Schema) -> tuple[pa.RecordBatch, pa.KeyValueMetadata]:
    """Build the zero-row pointer batch announcing this ref.

    Args:
        schema: The method's result schema.

    Returns:
        ``(batch, custom_metadata)`` as from
        :func:`make_external_location_batch`.

    """
    return make_external_location_batch(schema, self.url, sha256=self.sha256)

publish_external

publish_external(
    batch: RecordBatch,
    storage: ExternalStorage,
    compression: Compression | None = None,
    *,
    include_sha256: bool = True
) -> ExternalRef

Publish a unary result batch once and return a reusable reference.

Serializes batch exactly as the per-call externalizer does (an IPC stream of the schema plus this one batch), hashes the raw bytes, compresses when compression is given, and calls storage.upload once. Cache the returned :class:ExternalRef and return it from the unary method on later calls; the server writes the pointer directly.

Build batch against the method's result schema, e.g.::

schema = rpc_methods(MyService)["catalog"].result_schema
ref = publish_external(pa.RecordBatch.from_pydict({"result": [value]}, schema=schema), storage)
PARAMETER DESCRIPTION
batch

The 1-row result batch (single result column).

TYPE: RecordBatch

storage

Storage backend to upload to.

TYPE: ExternalStorage

compression

Optional compression applied before upload (pass the server's ServerExternalConfig.compression to match it).

TYPE: Compression | None DEFAULT: None

include_sha256

When False the ref carries no digest, so clients skip the content check.

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
An

class:ExternalRef for the uploaded object.

TYPE: ExternalRef

RAISES DESCRIPTION
ValueError

If batch does not have exactly one row.

Source code in vgi_rpc/external.py
def publish_external(
    batch: pa.RecordBatch,
    storage: ExternalStorage,
    compression: Compression | None = None,
    *,
    include_sha256: bool = True,
) -> ExternalRef:
    """Publish a unary result batch once and return a reusable reference.

    Serializes *batch* exactly as the per-call externalizer does (an IPC
    stream of the schema plus this one batch), hashes the raw bytes,
    compresses when *compression* is given, and calls ``storage.upload``
    once.  Cache the returned :class:`ExternalRef` and return it from the
    unary method on later calls; the server writes the pointer directly.

    Build *batch* against the method's result schema, e.g.::

        schema = rpc_methods(MyService)["catalog"].result_schema
        ref = publish_external(pa.RecordBatch.from_pydict({"result": [value]}, schema=schema), storage)

    Args:
        batch: The 1-row result batch (single ``result`` column).
        storage: Storage backend to upload to.
        compression: Optional compression applied before upload (pass the
            server's ``ServerExternalConfig.compression`` to match it).
        include_sha256: When ``False`` the ref carries no digest, so
            clients skip the content check.

    Returns:
        An :class:`ExternalRef` for the uploaded object.

    Raises:
        ValueError: If *batch* does not have exactly one row.

    """
    if batch.num_rows != 1:
        raise ValueError(f"publish_external expects a 1-row result batch, got {batch.num_rows} rows")
    url, data_sha256, _raw_size = _upload_ipc_bytes(
        _serialize_single_batch(batch, None), batch.schema, storage, compression
    )
    return ExternalRef(url=url, sha256=data_sha256 if include_sha256 else None)

Storage Protocol

ExternalStorage

Bases: Protocol

Pluggable storage interface for externalizing large batches.

Implementations must be thread-safe — upload() may be called concurrently from different server threads.

.. important:: Object lifecycle — vgi-rpc does not manage the lifecycle of uploaded objects. Data uploaded via upload() or generate_upload_url() persists indefinitely unless the server operator configures cleanup. Use storage-level lifecycle rules (e.g. S3 Lifecycle Policies, GCS Object Lifecycle Management) or bucket-level TTLs to automatically expire and delete stale objects.

upload

upload(
    data: bytes,
    schema: Schema,
    *,
    content_encoding: str | None = None
) -> str

Upload serialized IPC data and return a URL for retrieval.

The uploaded object is not automatically deleted — server operators are responsible for configuring object cleanup via storage lifecycle rules or TTLs.

PARAMETER DESCRIPTION
data

Complete Arrow IPC stream bytes.

TYPE: bytes

schema

The schema of the data being uploaded.

TYPE: Schema

content_encoding

Optional encoding applied to data (e.g. "zstd"). Backends should store this so that fetchers can decompress correctly.

TYPE: str | None DEFAULT: None

RETURNS DESCRIPTION
str

A URL (typically pre-signed) that can be fetched to retrieve

str

the uploaded data.

Source code in vgi_rpc/external.py
def upload(self, data: bytes, schema: pa.Schema, *, content_encoding: str | None = None) -> str:
    """Upload serialized IPC data and return a URL for retrieval.

    The uploaded object is not automatically deleted — server
    operators are responsible for configuring object cleanup via
    storage lifecycle rules or TTLs.

    Args:
        data: Complete Arrow IPC stream bytes.
        schema: The schema of the data being uploaded.
        content_encoding: Optional encoding applied to *data*
            (e.g. ``"zstd"``).  Backends should store this so that
            fetchers can decompress correctly.

    Returns:
        A URL (typically pre-signed) that can be fetched to retrieve
        the uploaded data.

    """
    ...

UploadUrlProvider

Bases: Protocol

Generates pre-signed upload URL pairs.

Implementations must be thread-safe — generate_upload_url() may be called concurrently from different server threads.

.. important:: Object lifecycle — vgi-rpc does not manage the lifecycle of uploaded objects. Data uploaded via these URLs persists indefinitely unless the server operator configures cleanup. Use storage-level lifecycle rules (e.g. S3 Lifecycle Policies, GCS Object Lifecycle Management) or bucket-level TTLs to automatically expire and delete stale objects.

generate_upload_url

generate_upload_url(schema: Schema) -> UploadUrl

Generate a pre-signed upload/download URL pair.

The caller receives time-limited PUT and GET URLs for a new storage object. The uploaded object is not automatically deleted — server operators are responsible for configuring object cleanup via storage lifecycle rules or TTLs.

PARAMETER DESCRIPTION
schema

The Arrow schema of the data to be uploaded. Backends may use this for content-type or metadata hints.

TYPE: Schema

RETURNS DESCRIPTION
UploadUrl

An UploadUrl with PUT and GET URLs for the same object.

Source code in vgi_rpc/external.py
def generate_upload_url(self, schema: pa.Schema) -> UploadUrl:
    """Generate a pre-signed upload/download URL pair.

    The caller receives time-limited PUT and GET URLs for a new
    storage object.  **The uploaded object is not automatically
    deleted** — server operators are responsible for configuring
    object cleanup via storage lifecycle rules or TTLs.

    Args:
        schema: The Arrow schema of the data to be uploaded.
            Backends may use this for content-type or metadata hints.

    Returns:
        An ``UploadUrl`` with PUT and GET URLs for the same object.

    """
    ...

UploadUrl dataclass

UploadUrl(
    upload_url: str, download_url: str, expires_at: datetime
)

Pre-signed URL pair for client-side data upload.

S3/GCS pre-signed URLs are signed per HTTP method, so a PUT URL cannot be used for GET — hence two URLs for the same storage object.

ATTRIBUTE DESCRIPTION
upload_url

Pre-signed PUT URL for uploading data.

TYPE: str

download_url

Pre-signed GET URL for downloading the uploaded data.

TYPE: str

expires_at

Expiration time for the pre-signed URLs (UTC).

TYPE: datetime

Validation

https_only_validator

https_only_validator(url: str) -> None

Reject URLs that do not use the https scheme.

This is the default url_validator for ExternalLocationConfig. It prevents the client from issuing requests over plain HTTP or other schemes (ftp, file, etc.) when resolving external-location pointers.

PARAMETER DESCRIPTION
url

The URL to validate.

TYPE: str

RAISES DESCRIPTION
ValueError

If the URL scheme is not https.

Source code in vgi_rpc/external.py
def https_only_validator(url: str) -> None:
    """Reject URLs that do not use the ``https`` scheme.

    This is the default ``url_validator`` for ``ExternalLocationConfig``.
    It prevents the client from issuing requests over plain HTTP or other
    schemes (``ftp``, ``file``, etc.) when resolving external-location
    pointers.

    Args:
        url: The URL to validate.

    Raises:
        ValueError: If the URL scheme is not ``https``.

    """
    from urllib.parse import urlparse

    parsed = urlparse(url)
    if parsed.scheme != "https":
        raise ValueError(f"URL scheme '{parsed.scheme}' not allowed (only 'https' is permitted)")

S3 Backend

Requires pip install vgi-rpc[s3].

S3Storage dataclass

S3Storage(
    bucket: str,
    prefix: str = "vgi-rpc/",
    presign_expiry_seconds: int = 3600,
    region_name: str | None = None,
    endpoint_url: str | None = None,
)

S3-backed ExternalStorage using boto3.

.. important:: Object lifecycle — uploaded objects persist indefinitely. Configure an S3 Lifecycle Policy_ on the bucket to expire objects under prefix (default vgi-rpc/) after a suitable retention period. See :mod:vgi_rpc.external for full details and examples.

ATTRIBUTE DESCRIPTION
bucket

S3 bucket name.

TYPE: str

prefix

Key prefix for uploaded objects.

TYPE: str

presign_expiry_seconds

Lifetime of pre-signed GET URLs.

TYPE: int

region_name

AWS region (None uses boto3 default).

TYPE: str | None

endpoint_url

Custom endpoint for S3-compatible services (e.g. MinIO, LocalStack).

TYPE: str | None

.. _S3 Lifecycle Policy: https://docs.aws.amazon.com/AmazonS3/latest/userguide/object-lifecycle-mgmt.html

upload

upload(
    data: bytes,
    schema: Schema,
    *,
    content_encoding: str | None = None
) -> str

Upload IPC data to S3 and return a pre-signed GET URL.

PARAMETER DESCRIPTION
data

Serialized Arrow IPC stream bytes.

TYPE: bytes

schema

Schema of the data (unused but required by protocol).

TYPE: Schema

content_encoding

Optional encoding applied to data (e.g. "zstd").

TYPE: str | None DEFAULT: None

RETURNS DESCRIPTION
str

A pre-signed URL that can be used to download the data.

Source code in vgi_rpc/s3.py
def upload(self, data: bytes, schema: pa.Schema, *, content_encoding: str | None = None) -> str:
    """Upload IPC data to S3 and return a pre-signed GET URL.

    Args:
        data: Serialized Arrow IPC stream bytes.
        schema: Schema of the data (unused but required by protocol).
        content_encoding: Optional encoding applied to *data*
            (e.g. ``"zstd"``).

    Returns:
        A pre-signed URL that can be used to download the data.

    """
    client = self._get_client()

    ext = ".arrow.zst" if content_encoding == "zstd" else ".arrow"
    key = f"{self.prefix}{uuid.uuid4().hex}{ext}"

    put_kwargs: dict[str, str | bytes] = {
        "Bucket": self.bucket,
        "Key": key,
        "Body": data,
        "ContentType": "application/octet-stream",
    }
    if content_encoding is not None:
        put_kwargs["ContentEncoding"] = content_encoding

    t0 = time.monotonic()
    try:
        client.put_object(**put_kwargs)
    except Exception as exc:
        _logger.error(
            "S3 upload failed: bucket=%s key=%s",
            self.bucket,
            key,
            exc_info=True,
            extra={"bucket": self.bucket, "key": key, "error_type": type(exc).__name__},
        )
        raise
    duration_ms = (time.monotonic() - t0) * 1000

    url: str = client.generate_presigned_url(
        "get_object",
        Params={"Bucket": self.bucket, "Key": key},
        ExpiresIn=self.presign_expiry_seconds,
    )
    _logger.debug(
        "S3 upload completed: bucket=%s key=%s (%d bytes, %.1fms)",
        self.bucket,
        key,
        len(data),
        duration_ms,
        extra={"bucket": self.bucket, "key": key, "size_bytes": len(data), "duration_ms": round(duration_ms, 2)},
    )
    return url

generate_upload_url

generate_upload_url(schema: Schema) -> UploadUrl

Generate pre-signed PUT and GET URLs for client-side upload.

The created S3 object is not automatically deleted. Configure S3 Lifecycle Policies on the bucket to expire objects after a suitable retention period.

PARAMETER DESCRIPTION
schema

The Arrow schema of the data to be uploaded (unused but available for metadata hints).

TYPE: Schema

RETURNS DESCRIPTION
UploadUrl

An UploadUrl with PUT and GET pre-signed URLs for the

UploadUrl

same S3 object.

Source code in vgi_rpc/s3.py
def generate_upload_url(self, schema: pa.Schema) -> UploadUrl:
    """Generate pre-signed PUT and GET URLs for client-side upload.

    The created S3 object is not automatically deleted.  Configure
    S3 Lifecycle Policies on the bucket to expire objects after
    a suitable retention period.

    Args:
        schema: The Arrow schema of the data to be uploaded
            (unused but available for metadata hints).

    Returns:
        An ``UploadUrl`` with PUT and GET pre-signed URLs for the
        same S3 object.

    """
    client = self._get_client()
    key = f"{self.prefix}{uuid.uuid4().hex}.arrow"
    params = {"Bucket": self.bucket, "Key": key}
    expires_at = datetime.now(UTC) + timedelta(seconds=self.presign_expiry_seconds)

    put_url: str = client.generate_presigned_url(
        "put_object",
        Params=params,
        ExpiresIn=self.presign_expiry_seconds,
    )
    get_url: str = client.generate_presigned_url(
        "get_object",
        Params=params,
        ExpiresIn=self.presign_expiry_seconds,
    )

    _logger.debug(
        "S3 upload URL generated: bucket=%s key=%s",
        self.bucket,
        key,
        extra={"bucket": self.bucket, "key": key},
    )
    return UploadUrl(upload_url=put_url, download_url=get_url, expires_at=expires_at)

GCS Backend

Requires pip install vgi-rpc[gcs].

GCSStorage dataclass

GCSStorage(
    bucket: str,
    prefix: str = "vgi-rpc/",
    presign_expiry_seconds: int = 3600,
    project: str | None = None,
)

GCS-backed ExternalStorage using google-cloud-storage.

.. important:: Object lifecycle — uploaded objects persist indefinitely. Configure Object Lifecycle Management_ on the bucket to delete objects under prefix (default vgi-rpc/) after a suitable retention period. See :mod:vgi_rpc.external for full details and examples.

ATTRIBUTE DESCRIPTION
bucket

GCS bucket name.

TYPE: str

prefix

Key prefix for uploaded objects.

TYPE: str

presign_expiry_seconds

Lifetime of signed GET URLs.

TYPE: int

project

GCS project ID (None uses Application Default Credentials default project).

TYPE: str | None

.. _Object Lifecycle Management: https://cloud.google.com/storage/docs/lifecycle

upload

upload(
    data: bytes,
    schema: Schema,
    *,
    content_encoding: str | None = None
) -> str

Upload IPC data to GCS and return a signed GET URL.

PARAMETER DESCRIPTION
data

Serialized Arrow IPC stream bytes.

TYPE: bytes

schema

Schema of the data (unused but required by protocol).

TYPE: Schema

content_encoding

Optional encoding applied to data (e.g. "zstd").

TYPE: str | None DEFAULT: None

RETURNS DESCRIPTION
str

A signed URL that can be used to download the data.

Source code in vgi_rpc/gcs.py
def upload(self, data: bytes, schema: pa.Schema, *, content_encoding: str | None = None) -> str:
    """Upload IPC data to GCS and return a signed GET URL.

    Args:
        data: Serialized Arrow IPC stream bytes.
        schema: Schema of the data (unused but required by protocol).
        content_encoding: Optional encoding applied to *data*
            (e.g. ``"zstd"``).

    Returns:
        A signed URL that can be used to download the data.

    """
    client = self._get_client()
    bucket = client.bucket(self.bucket)
    ext = ".arrow.zst" if content_encoding == "zstd" else ".arrow"
    blob_name = f"{self.prefix}{uuid.uuid4().hex}{ext}"
    blob = bucket.blob(blob_name)
    if content_encoding is not None:
        blob.content_encoding = content_encoding
    t0 = time.monotonic()
    try:
        blob.upload_from_string(data, content_type="application/octet-stream")
    except Exception as exc:
        _logger.error(
            "GCS upload failed: bucket=%s key=%s",
            self.bucket,
            blob_name,
            exc_info=True,
            extra={"bucket": self.bucket, "key": blob_name, "error_type": type(exc).__name__},
        )
        raise
    duration_ms = (time.monotonic() - t0) * 1000

    url: str = blob.generate_signed_url(
        version="v4",
        expiration=timedelta(seconds=self.presign_expiry_seconds),
        method="GET",
    )
    _logger.debug(
        "GCS upload completed: bucket=%s key=%s (%d bytes, %.1fms)",
        self.bucket,
        blob_name,
        len(data),
        duration_ms,
        extra={
            "bucket": self.bucket,
            "key": blob_name,
            "size_bytes": len(data),
            "duration_ms": round(duration_ms, 2),
        },
    )
    return url

generate_upload_url

generate_upload_url(schema: Schema) -> UploadUrl

Generate signed PUT and GET URLs for client-side upload.

The created GCS object is not automatically deleted. Configure GCS Object Lifecycle Management on the bucket to expire objects after a suitable retention period.

PARAMETER DESCRIPTION
schema

The Arrow schema of the data to be uploaded (unused but available for metadata hints).

TYPE: Schema

RETURNS DESCRIPTION
UploadUrl

An UploadUrl with PUT and GET signed URLs for the

UploadUrl

same GCS object.

Source code in vgi_rpc/gcs.py
def generate_upload_url(self, schema: pa.Schema) -> UploadUrl:
    """Generate signed PUT and GET URLs for client-side upload.

    The created GCS object is not automatically deleted.  Configure
    GCS Object Lifecycle Management on the bucket to expire objects
    after a suitable retention period.

    Args:
        schema: The Arrow schema of the data to be uploaded
            (unused but available for metadata hints).

    Returns:
        An ``UploadUrl`` with PUT and GET signed URLs for the
        same GCS object.

    """
    client = self._get_client()
    bucket = client.bucket(self.bucket)
    blob_name = f"{self.prefix}{uuid.uuid4().hex}.arrow"
    blob = bucket.blob(blob_name)
    expiration = timedelta(seconds=self.presign_expiry_seconds)
    expires_at = datetime.now(UTC) + expiration

    put_url: str = blob.generate_signed_url(
        version="v4",
        expiration=expiration,
        method="PUT",
        content_type="application/octet-stream",
    )
    get_url: str = blob.generate_signed_url(
        version="v4",
        expiration=expiration,
        method="GET",
    )

    _logger.debug(
        "GCS upload URL generated: bucket=%s key=%s",
        self.bucket,
        blob_name,
        extra={"bucket": self.bucket, "key": blob_name},
    )
    return UploadUrl(upload_url=put_url, download_url=get_url, expires_at=expires_at)