Skip to main content

grainlift

Build ADBC services in Python. Applications load the existing native Grainlift ADBC driver; your worker supplies query behavior and lazy Arrow batches over VGI-RPC. No downstream ADBC driver is required.

The toolkit exposes the Grainlift protocol 0.4.0 ADBC operation surface over HTTP and authenticated TCP/mTLS: transactions, statements, preparation, typed options, parameter batches and streams, updates and ingestion, metadata, partitioned results, and Substrait plans. Your backend implements each capability through Connection and Statement hooks. The toolkit manages authentication, ownership, quotas, Arrow transport, and cleanup; it does not emulate database semantics. Unimplemented backend hooks return ADBC NOT_IMPLEMENTED.

Existing Connection.execute(sql) workers remain supported through a statement adapter. Those workers retain their original query-only capabilities and should use autocommit=True with the Python ADBC driver manager.

Installation

Python 3.13 or newer is required.

pip install grainlift

To use the optional supervised Granian HTTP host, install the granian extra:

pip install "grainlift[granian]"

This package is the service side only. Client applications connect through the native Grainlift ADBC driver, which must speak the same Grainlift protocol version (0.4.0).

Hosting

TcpServer provides bounded TCP/mTLS admission, verified certificate URI identity, read/write deadlines and explicit draining/shutdown. Plain TCP requires an explicit local principal and a loopback bind address. serve_granian is an optional supervised HTTP host; install grainlift[granian]. It creates the worker inside one serving process and preserves ADBC session affinity. The existing serve() Waitress entry point remains available.

from grainlift import Service, TcpServer, TLSConfig

# Use the AnswerWorker defined below, or your own Worker implementation.
tls = TLSConfig(
    certificate="server.pem",
    private_key="server-key.pem",
    client_ca="clients-ca.pem",
    principals={"spiffe://example.org/query-client": "analytics"},
)
with TcpServer(Service(AnswerWorker()), host="0.0.0.0", port=8443, tls=tls) as host:
    shutdown_event.wait()  # Your process supervisor signals this event.

See hosting and lifecycle configuration for a complete signal-handling example, Granian deployment, limits, credential rotation and the distinction between draining, forced transport shutdown and worker cleanup.

See grainlift-hello-world-python for a complete worker and ordinary ADBC client.

Authoring statements

Override Connection.new_statement() to create independent backend statements. This small example implements one query; add only capabilities the backend can perform correctly. See the API contract for every hook.

import pyarrow as pa
from grainlift import AdbcError, Connection, QueryResult, Statement, Worker

SCHEMA = pa.schema([("answer", pa.int64())])

class AnswerStatement(Statement):
    def __init__(self) -> None:
        self.sql: str | None = None

    def set_sql_query(self, sql: str) -> None:
        self.sql = sql

    def execute(self) -> QueryResult:
        if self.sql != "SELECT 42":
            raise AdbcError("Expected SELECT 42", "invalid_arguments", sqlstate="42000")
        batch = pa.record_batch([[42]], schema=SCHEMA)
        return QueryResult(SCHEMA, iter([batch]))

class AnswerConnection(Connection):
    def new_statement(self) -> Statement:
        return AnswerStatement()

class AnswerWorker(Worker):
    def connect(self, principal: str) -> Connection:
        return AnswerConnection()

For a database-backed service, retain the real backend statement, delegate prepare, bind, bind_stream, execute_update, and other supported hooks, and close that backend object in Statement.close(). Implement transactions in the connection's option/commit/rollback hooks. Ingestion uses typed statement options such as adbc.ingest.target_table and adbc.ingest.mode, parameter binding, and execute_update(); the backend performs the actual ingestion. Substrait plans are opaque bytes passed to the backend, not translated into SQL.

Return QueryResult(schema, batch_iterator) for queries and metadata. Iterators should expose close() to release resources. Generate bounded batches lazily, and preserve the same schema, including metadata, throughout each result.

Serializable results

Instead of an iterator, a result can be a ResultProducer: a dataclass whose fields are the complete resumable state, with a produce() method that returns the next batch or None. Over HTTP the service serializes the producer into the encrypted, principal-bound continuation token after every batch, which is the same mechanism VGI-RPC streams use. The server therefore retains no iterator or replay batch between fetches, and a retried fetch recomputes its batch from the token. Process isolation and TCP drive the same object in memory.

from dataclasses import dataclass

from grainlift import QueryResult, ResultProducer

@dataclass
class Countdown(ResultProducer):
    remaining: int

    def produce(self) -> pa.RecordBatch | None:
        if self.remaining == 0:
            return None
        self.remaining -= 1
        return pa.record_batch([[self.remaining]], schema=SCHEMA)

# In Statement.execute():
return QueryResult.from_producer(SCHEMA, Countdown(3))

Fields must be Arrow-serializable, and Limits.producer_state_bytes (64 KiB by default) bounds the encoded state. Keep sockets, files and backend cursors out of producers; use an iterator for those results. Sessions, and therefore result handles, still belong to one process, so producers do not make results survive a restart.

Development host

grainlift serve module:Factory (or grainlift.cli.run("module:Factory") from your own console script) serves a worker on loopback with --host waitress, granian or mtls. HTTP hosts read the bearer token from GRAINLIFT_TOKEN, or generate and print one when it is unset.

Anonymous access

Authentication is required by default. A service that only exposes public, read-only data can also accept clients without credentials:

app = service.app(anonymous_principal="anonymous")                        # anonymous only
app = service.app(tokens={token: "analyst"}, anonymous_principal="anonymous")  # both

Requests without an Authorization header act as the anonymous principal. The native driver sends no header when grainlift.auth.bearer_token is unset. A request with a wrong token is rejected, never downgraded to anonymous. All anonymous clients share one principal, so it should reach only public capabilities: the toolkit does not decide which statements are read-only, and the worker can check the principal passed to Worker.connect. Anonymous continuation tokens are sealed in a separate authentication domain, and the anonymous principal must differ from every token principal. serve(), serve_granian() and grainlift serve --auth anonymous accept the same option.

Lifecycle and resource contract

Each service owns its sessions in one process. Route all calls, including continuations, to that process. Restart invalidates every handle. There is no transparent multi-replica behavior. Calls on one session serialize under its own lock; independent sessions and connection factories can run concurrently. Worker factories must therefore be thread-safe. Opening connections reserve session quota before invoking the factory. The reaper skips busy sessions.

Session lock waits default to five seconds. In-process callbacks must still be bounded and cooperative: Python threads cannot safely preempt arbitrary calls. Service.close() rejects new work and requests cancellation of busy connections. It waits up to Limits.shutdown_seconds for busy session locks, then reports ADBC TIMEOUT if callbacks remain. Their eventual return triggers cleanup. This wait budget does not preempt in-process cancellation or cleanup hooks.

Defaults: 64 sessions, 32 statements and 32 result handles per session, one query result per statement, 1 MiB of referenced Arrow buffers per batch, 1 MiB schema descriptor, 2 MiB HTTP request body, 64 KiB SQL, and 300 seconds idle lifetime. Metadata and partition readers consume the same result-handle quota. Typed option responses and serialized partition responses, including their schema, have the 1 MiB response budget. Partition execution returns at most 1,024 descriptors. HTTP responses have a 2 MiB transport budget, including protocol overhead. Oversized results are rejected. Parameter uploads use the bounded spool described below instead of collecting all parameter batches in memory. Transport/schema overhead can therefore reject a batch below its buffer limit.

Results retain at most one batch for replay. Reading the immediately previous sequence repeats that batch; skipped/older sequences fail. EOF closes the iterator. Explicit result release, statement close/reuse, connection close, stream cancel, idle expiry, and Service.close() release resources. A broken HTTP connection does not necessarily mean the logical cursor is abandoned: idle expiry cleans up clients that disappear. Use Service as a context manager. The same rule applies to TCP: control and result sockets can belong to the same logical session, so closing one socket does not implicitly destroy that session.

Parameter binding uses an anonymous Arrow IPC spool with a 64 MiB cumulative Limits.bind_bytes budget, including schema, dictionaries, and stream framing. Each batch also obeys Limits.batch_bytes. bind accepts one batch; bind_stream accepts a sequence, including an empty stream with a known schema. An explicit finish turn distinguishes successful end-of-input from disconnect. The backend receives parameters only after that finish is validated.

Upload turns carry sequence numbers; only an identical immediate replay is acknowledged again. A statement can own one completed binding and one pending replacement, each separately capped. Failure, cancellation, or idle expiry of the pending upload leaves the previous completed binding intact. Execution and preparation reject unfinished uploads. Backend parameter readers remain valid until successful replacement, query/plan replacement, or statement close. Statement close invokes backend cleanup before closing its retained reader.

Security

Every request and continuation is authenticated. Opaque session handles are bound to a configured principal; child handles are scoped to that session. Targets are server-selected. Service(database_options=..., connection_options=...) defines authoritative settings: callers may neither supply those keys in either opening option scope nor mutate them afterward. Other caller options are decoded strictly and delegated to the worker. Option values are str, bytes, signed 64-bit int, or IEEE 754 float (including NaN and infinities); booleans and duplicate keys are rejected.

Worker.open_connection(principal, database_options, connection_options) is the database-factory hook. Its default rejects database options, calls the existing connect(principal), and applies connection options with cleanup on failure. Override it for an actual configured database factory. get_option results are client-visible: never expose credentials through a backend getter.

The CLI binds to loopback and requires a bearer token. For deployment, host Service.app(tokens={token: principal}) behind HTTPS and enforce process affinity. Do not expose plain HTTP with bearer tokens on an untrusted network. The supported TCP and mTLS hosts have explicit admission, I/O, and shutdown bounds; see hosting. Iroh serving is not implemented here.

The WSGI wrapper filters VGI-RPC transport logs during Grainlift requests because diagnostics can contain SQL, credentials, Arrow values, and raw exceptions. It preserves logger levels, handlers, propagation, and unrelated applications' logs. It installs filters on loaded VGI-RPC loggers at app creation and request entry; re-audit this boundary when adding transports or upgrading VGI-RPC. Grainlift's own access event contains only status and duration. AdbcError preserves client-facing status, SQLSTATE, vendor code and binary details; unexpected worker errors become a generic INTERNAL error. Do not include secrets in client-facing AdbcError messages. Worker-authored logging is the worker author's responsibility.

Partition descriptors are signed wrappers bound to the service instance, target, and authenticated principal. The same principal can read one through another connection until its Limits.idle_seconds expiry. Tampering, another principal, another service instance, and expired descriptors are rejected before backend execution. A restart invalidates wrappers. The toolkit creates no durable partition store and makes no cross-replica or transaction-recovery guarantee.

Optional process isolation

from grainlift import IsolatedWorker, Limits, Service

worker = IsolatedWorker(
    "my_worker:MyWorker",  # Importable Worker subclass, constructed in the child.
    worker_options={"example_setting": "value"},
    timeout_seconds=5,
    startup_timeout_seconds=10,
    max_message_bytes=2 * 1024 * 1024,
    max_results=32,
    max_statements=32,
    max_bind_bytes=64 * 1024 * 1024,
)
service = Service(worker, limits=Limits(sessions=16))
app = service.app(tokens={"replace-with-a-secret": "configured-principal"})

Each connection owns a spawned process and two anonymous pipes. Worker options must be JSON-compatible. Results use capped JSON/Arrow IPC messages; SQL executes once and each fetch pulls one batch. Schema and transport overhead count toward the IPC cap. Overflow returns INVALID_DATA and releases the affected result. The child independently caps live statements, results, and parameter spools. The full statement/connection capability surface is forwarded to its backend; the selected worker still determines which capabilities are implemented. Use an importable module and guard executable entry points with if __name__ == "__main__": for multiprocessing spawn.

Opening, statement/connection operations, upload turns, fetching, and cleanup have finite deadlines. Timeout terminates the child and returns TIMEOUT; a crash returns IO. Termination waits up to 0.7 additional seconds for the child to exit. Cancellation also terminates the child. All existing statements/results on that connection become unusable; the client must close it and reconnect explicitly. No transaction recovery or rollback behavior is emulated. Per-operation deadlines do not constitute a single aggregate service-shutdown deadline.

In-process workers may override Connection.cancel() and Statement.cancel() with thread-safe, nonblocking hooks. Otherwise cancellation returns NOT_IMPLEMENTED. Legacy statements delegate cancellation to their connection. Connection cancellation authenticates the session owner; statement cancellation additionally requires an operation currently active on that statement. Closing an HTTP result stream releases that cursor and is separate from ADBC cancellation.

Isolation contains callback hangs and crashes; it is not a sandbox for untrusted code. Session quotas bound child count and IPC limits bound transferred data. Apply OS/container memory and CPU quotas to bound allocations inside a worker. Workers must not spawn unmanaged descendants. Child stdout/stderr are discarded; worker-managed logging destinations remain the author's responsibility.

Validation and native-client limitation

Tests cover the actual native Grainlift C ABI through adbc-driver-manager, including multiple batches, empty results with schemas, schema inference, authentication, early release, errors, and handle cleanup. Toolkit tests exercise principal isolation (including HTTP continuations), immediate replay, idle expiry, stream cancellation, invalid schemas, and limits below, at, and above boundaries. Focused feature tests cover typed options, authoritative configuration, metadata filters and quotas, statement hooks, transactions, binding replay and spool cleanup, and signed partition ownership. Isolated-worker and native-driver fixtures exercise backend delegation separately; see the API coverage map.

The wire preserves vendor_code, but the current Rust adbc_ffi dependency in Grainlift overwrites the C error's vendor-code slot with the ADBC 1.1 private-data sentinel. Consequently adbc-driver-manager 1.12.0 exposes vendor_code=None in these tests, even when the worker sent a code. Status, SQLSTATE, and binary details survive. This is an existing native-client limitation, not a claim of complete end-to-end error fidelity.

Development

Python 3.13+ and uv are required. Clone this repository and run:

uv sync --locked --extra granian
./check_quality.sh
uv run --no-sync pytest

The toolkit uses published vgi-rpc[http]>=0.47.1; the lockfile pins its index release and dependencies. No modified VGI runtime is required. Unary service methods return frozen typed response dataclasses, which VGI serializes through its standard binary result envelope. Opening connections, setting options, and filtered metadata discovery use named request dataclasses with native Arrow fields. Option values use a nested typed record, and partition descriptors use an Arrow binary list with signed typed claims. Parameter uploads use a fixed one-row envelope so empty and zero-column Arrow data remain unambiguous. Protocol 0.4.0 requires a matching Grainlift client; earlier wire clients must be upgraded together with the service. See the API guide for details.

The CI workflow checks Linux/macOS with Python 3.13/3.14, runs Ruff, formatting, strict mypy and isolated pydoclint, then installs the built wheel and runs its tests. A configured workflow is not evidence that every matrix job has passed; consult the repository's Actions results for the revision being deployed.

Metadata

Release files for grainlift 0.2.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for grainlift 0.2.0
File Size Uploaded
grainlift-0.2.0.tar.gz 107.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for grainlift 0.2.0
File Interpreter ABI Platform
grainlift-0.2.0-py3-none-any.whl Python 3 none any Details

Total release size: 170.7 kB

Release files / grainlift-0.2.0.tar.gz

Download URL grainlift-0.2.0.tar.gz
Size 107.2 kB
Tags Source
SHA-256 checksum
How to use checksums
65aca6031daa7c7469b57bc1244e7eadbe55f1d89a2563e83240c446562a64a4
BLAKE2b-256 checksum
How to use checksums
80d8ede375ce18ad72237ba4e3a04cdc59b21fef913bae1ca74afb04523d0c03
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 29, 2026.

Transparency log

Release files / grainlift-0.2.0-py3-none-any.whl

Download URL grainlift-0.2.0-py3-none-any.whl
Size 63.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
26885b8b91202f3ac55b245b064b885c9528442060aa1eb94b6cf7c251e0d92d
BLAKE2b-256 checksum
How to use checksums
52fef5d720ce0e971e0d94439c96950c398faa772f29b7fad69cfb32c4929bda
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 29, 2026.

Transparency log

Release history Release notifications | RSS feed

0.2.1

2 release files

This release

0.2.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page