Transport-agnostic RPC framework built on Apache Arrow IPC serialization — the Java port of vgi-rpc.
Built by 🚜 Query.Farm
Define RPC interfaces as ordinary Java interfaces. The framework derives Apache Arrow schemas from your method signatures and record component types, and hands you a typed client proxy with automatic serialization/deserialization. There are no .proto files or codegen steps — your Java types are the schema. Unlike JSON-over-HTTP, structured data stays in Arrow's columnar format for efficient transfer, which pays off for large or batch-oriented workloads.
This is a port of the Python reference implementation, vgi-rpc, and is wire-compatible with it: the same calls interoperate across the Python, Java, Go, and C++ peers (the conformance suite runs the Python driver against this Java worker over every transport).
- Interface-based services — define a service as a typed Java interface; the client proxy preserves that interface for full IDE autocompletion.
- Apache Arrow IPC wire format — columnar serialization for structured data.
- Two method types — unary calls and streaming (producer and exchange patterns).
- Transport-agnostic — stdio pipe, subprocess, Unix domain socket, raw TCP socket (trusted networks — no auth/TLS), shared memory, or HTTP.
- Automatic schema inference — Java types and
recordcomponents map to Arrow types;@ArrowFieldrefines them. - Pluggable authentication —
AuthContext+ authenticators for HTTP (bearer, mTLS/XFCC; JWT/OAuth in the optionalvgirpc-oauthmodule). - Runtime introspection — the
vgi_rpc.Reflection.v1protocol (list_protocols, thendescribe) for dynamic service discovery —Introspect.listProtocols/describeProtocolover any held connection (see Discovering what a server hosts) — with a canonical protocol hash every port agrees on. The old__describe__RPC is retired; a request for it is refused with a message naming its replacement. - Shared-memory transport — zero-copy batch transfer between co-located processes (auto-negotiated on JDK 22+ via a multi-release overlay; transparent pipe fallback otherwise).
- Large-batch externalization — oversized batches transparently spilled to S3 (
vgirpc-s3) or GCS (vgirpc-gcs). - Pre-published results — a unary method can answer with an
ExternalRefto an object published once, so a large, rarely-changing result is never re-serialized or re-uploaded per call.
- Java 21+ at runtime. The shared-memory side-channel additionally requires JDK 22+ (where
java.lang.foreignis GA); on 21 it transparently falls back to inline transfer.
Artifacts are published to Maven Central under the farm.query group.
Gradle (Kotlin DSL):
dependencies {
implementation("farm.query:vgirpc:0.31.0") // core: protocol, transports, HTTP, schema
implementation("farm.query:vgirpc-iroh:0.31.0") // optional: official native Iroh binding
implementation("farm.query:vgirpc-oauth:0.31.0") // optional: JWT / OAuth / PKCE auth
implementation("farm.query:vgirpc-s3:0.31.0") // optional: S3 external storage
implementation("farm.query:vgirpc-gcs:0.31.0") // optional: GCS external storage
}Maven:
<dependency>
<groupId>farm.query</groupId>
<artifactId>vgirpc</artifactId>
<version>0.31.0</version>
</dependency>The core depends on Apache Arrow and SLF4J (API only — bring your own logging backend).
1. Define a service as a Java interface (shared by client and server):
public interface Calculator {
double add(double a, double b);
String greet(String name);
}2. Implement it and serve it. A worker typically serves over stdio so a parent process can drive it as a subprocess:
import farm.query.vgirpc.RpcServer;
import farm.query.vgirpc.transport.StdioTransport;
public final class CalculatorWorker {
public static void main(String[] args) {
Calculator impl = new Calculator() {
public double add(double a, double b) { return a + b; }
public String greet(String name) { return "Hello, " + name + "!"; }
};
RpcServer server = new RpcServer(Calculator.class, impl);
try (StdioTransport transport = new StdioTransport()) {
server.serve(transport);
}
}
}3. Call it through a typed proxy. The client launches the worker and gets back something that is a Calculator:
import farm.query.vgirpc.RpcConnection;
import farm.query.vgirpc.transport.SubprocessTransport;
import java.util.List;
var transport = new SubprocessTransport(List.of(
"java", "--add-opens=java.base/java.nio=ALL-UNNAMED", "-Dio.netty.noUnsafe=false",
"-cp", "worker.jar", "CalculatorWorker"));
try (RpcConnection conn = new RpcConnection(transport)) {
Calculator calc = conn.proxy(Calculator.class);
double sum = calc.add(2.0, 3.0); // 5.0
String hello = calc.greet("World"); // "Hello, World!"
}Three things to get right:
- Run with
--add-opens=java.base/java.nio=ALL-UNNAMEDon every JVM that touches the library (both the worker and the client above) — Apache Arrow accessesjava.niointernals and throws on startup without it. Notice it's passed both to the client JVM and, in theSubprocessTransportcommand, to the spawned worker.- On Java 25+, also pass
-Dio.netty.noUnsafe=false. Arrow allocates through Netty 4.2, which switchessun.misc.Unsafeoff by default on Java 25+, and Arrow's allocator cannot start without it (apache/arrow-java#728). The library sets this default itself when it is the first thing in the process to touch Arrow, but a JVM that allocates Arrow memory some other way first needs the flag. It is harmless on Java 21–24.- Compile services with
-parameters— the framework binds call arguments by parameter name (matching the Python reference's keyword-argument wire semantics).
| Module | Purpose |
|---|---|
vgirpc |
Core library — wire protocol, transports, HTTP server/client (Jetty 12), schema derivation, marshalling, external-location support, shared-memory primitive. |
vgirpc-iroh |
Optional official Kotlin/JVM Iroh provider for raw Arrow mux and HTTP semantics. |
vgirpc-oauth |
Optional OAuth/JWT support (JWKS validation, PKCE, signed cookies). Split out so core users don't pull nimbus-jose-jwt. |
vgirpc-s3 |
Amazon S3 ExternalStorage backend for large-batch externalization. |
vgirpc-gcs |
Google Cloud Storage ExternalStorage backend. |
| Transport | Use case |
|---|---|
stdio (StdioTransport) |
Worker process driven over stdin/stdout by a parent. |
subprocess (SubprocessTransport) |
Client spawns and talks to a worker subprocess. |
Unix socket (UnixSocketTransport) |
Co-located processes over a domain socket. |
| shared memory | Zero-copy batch transfer for co-located processes; auto-negotiated on JDK 22+, transparent pipe fallback otherwise. |
HTTP (HttpServer / Jetty 12) |
Networked, stateless-server streaming; auth via authenticators. |
Iroh Arrow mux (IrohTransports) |
Stateful, authenticated QUIC to an iroh:// worker. |
HTTP over Iroh (HttpRpcConnection.irohBuilder) |
Existing HTTP state/continuation semantics carried over authenticated iroh-http/2. |
Both Iroh modes use the official computer.iroh JVM binding from the optional
vgirpc-iroh module. HTTP over Iroh retains the ordinary HTTP client’s OPTIONS
discovery, response budgets, headers, authentication, and continuation logic:
try (var connection = HttpRpcConnection.irohBuilder(
"httpi://<64-lowercase-hex-endpoint-id>/vgi",
IrohTransportOptions.defaults())
.bearerToken(accessToken)
.buildIroh()) {
Calculator calculator = connection.proxy(Calculator.class);
double result = calculator.add(2.0, 3.0);
}The provider owns one authenticated Iroh connection, opens one bidirectional
stream per HTTP request, and closes it with the HttpRpcConnection. It does not
start or download a helper executable.
- Unary — request batch in, one result (or error) batch out.
- Streaming — a
RpcStream<S extends StreamState>whose state'sprocess(input, out, ctx)runs once per tick, in two flavours: producer (server emits a sequence of output batches) and exchange (lockstep input batch → output batch).
Some unary results are large and change rarely (a worker's whole catalog, say). Rather than letting
the externalizer serialize and upload them on every call, publish the result once with
Externalizer.publishExternal and answer later calls with the cached reference:
public interface CatalogService {
String catalog(CallContext ctx); // CallContext is off-wire
}
final class CatalogImpl implements CatalogService {
private volatile ExternalRef ref;
@Override public String catalog(CallContext ctx) {
ExternalRef r = ref;
if (r == null) {
Schema schema = ServiceIntrospector.describe(CatalogService.class).get("catalog").resultSchema();
// Builds {result: [value]}, serializes it exactly like the per-call externalizer,
// hashes the raw bytes, compresses (zstd here) and uploads once.
r = ref = Externalizer.publishExternal(schema, buildCatalog(), storage,
ExternalLocationConfig.Compression.zstd(), /* includeSha256 */ true);
}
ctx.respondWithExternalRef(r);
return null; // ignored once a ref is set
}
}ctx.respondWithExternalRef(ref) makes the dispatcher write the ExternalLocation pointer batch
(vgi_rpc.location, plus vgi_rpc.location.sha256 only when the ref has a digest) directly: the
returned value is ignored and not validated, nothing is serialized or uploaded, the inline and
shared-memory routes are never taken, and this holds whether or not the server has external
storage configured and regardless of the externalization threshold. A ref is not charged to
max_externalized_response_bytes. It works on every transport (pipe, unix, TCP, HTTP); clients
resolve it like any other pointer and need no change. Unary methods with a result only — a stream
method's context refuses the call.
publishExternal also takes a ready-built 1-row VectorSchemaRoot. A ref built with
includeSha256 = false (or ExternalRef.of(url) for an object published out of band) carries no
digest, so clients skip the content check — use that for an object rewritten in place. You own the
ref's cache 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, a pre-signed URL expires (re-sign or rebuild the
ref before then), and a ref must only be returned to callers who are all entitled to the same
content.
A service's wire name is its routing key: it rides every request as vgi_rpc.protocol, and
over HTTP it is also the protocol path segment. By default it is the interface's simple name.
Declare it explicitly when the name is a cross-implementation contract:
@ProtocolName("orders.v2")
@ProtocolVersion("2.0.0")
public interface OrderService { ... }Put the major version in the name. An incompatible major then becomes a different protocol and
an unroutable request 404s — an answer every proxy and load balancer understands without an Arrow
parser — and orders.v1 and orders.v2 can be served side by side while clients migrate. A Java
simple name cannot express that shape at all, since no Java identifier contains a dot.
The declaration is read from the interface's own annotations. An interface that extends a declared protocol and does not redeclare gets its own simple name rather than silently answering to its parent's routing key.
A request naming a protocol this server does not host is refused with ProtocolNotSupportedError
(protocol_not_supported), distinct from protocol_not_specified for a request that named none
and from method_not_implemented for a hosted protocol missing the method. A client probing for
an optional protocol depends on telling those apart.
The same answer comes back from vgi_rpc.Reflection.v1's describe, which asks the same question
with the name as an argument rather than as a routing key. Both check the name against the grammar
before the lookup, so a name that cannot be a protocol name is refused without being echoed back
in the message.
Introspect.listProtocols(target) and Introspect.describeProtocol(target, name) ask
vgi_rpc.Reflection.v1 over a connection you already hold. target is an RpcConnection (pipe,
subprocess, Unix socket, TCP, Iroh), an HttpRpcConnection (HTTP, HTTP over Iroh), a typed proxy
from either one's proxy(Class) — bound to any protocol — or a raw RpcTransport. The connection is
reused and never closed; nothing new is opened.
try (RpcConnection conn = new RpcConnection(transport)) {
Calculator calc = conn.proxy(Calculator.class);
for (Introspect.HostedProtocol p : Introspect.listProtocols(calc)) {
// Server order: application protocols (primary first), then vgi_rpc.Reflection.v1.
System.out.println(p.name() + " " + p.version() + " " + p.hash());
}
Introspect.ServiceDescription d = Introspect.describeProtocol(conn, "Calculator");
calc.add(1, 2); // the connection is still yours
}HostedProtocol is an immutable record: name, version, hash, deprecated (default false),
deprecationMessage (default "") and features (default empty). describeProtocol lists first,
then describes, so the two failures stay apart: a server that does not host reflection throws
ReflectionNotSupportedError (an RpcError carrying the server's error fields), while an unknown
name is an ordinary RpcError with errorKind() "protocol_not_supported". No listing is ever
inferred, and the connection stays usable after either.
A Java RpcServer always hosts reflection. The Python reference hosts it only when built with
enable_describe=True (the default is false; vgi-rpc-conformance --describe). Against a server
without it, or one that predates reflection (method_not_implemented, UNIMPLEMENTED, or a bare
HTTP 404), these throw ReflectionNotSupportedError. On a byte-stream transport, don't call them while
a stream is open on the same connection: the calls would interleave on one channel.
When the Python and Java implementations disagree, Python is the reference. Wire format, metadata keys, error semantics, and stream-state token layout match byte-for-byte so the two interoperate. See the Python project's README for the higher-level protocol design.
Proxy proof lets a worker refuse any request that did not arrive through a trusted proxy. The proxy mints a per-request HMAC-SHA256 over a timestamp, a fresh nonce and the worker's own identifier, keyed by a secret shared only with that worker. Unlike a forwarded assertion about what happened at a TLS terminator, a proof cannot be produced by someone who merely reaches the worker directly — without the secret there is nothing to replay.
var secrets = ProxyProof.parseSecrets("prod-use1:" + hexSecret);
var config = ProxyProof.Config.of(ProxyProof.Mode.REQUIRE, "worker-a", secrets);
HttpServer.Config.builder()
.authenticator(ProxyProof.require(config, existingAuthenticator)) // inner may be null
.proxyProofRequired(true) // REQUIRE mode only
.build();It composes as an AND, not an alternative: do not pass the gate to Authenticator.chain,
whose first-authenticated-wins semantics would let any later credential bypass it.
proxyProofRequired(true) advertises VGI-Proxy-Proof-Required: true on every response,
GET /health and OPTIONS included, so an operator or proxy can confirm the worker really does
reject unproofed requests — otherwise a misconfiguration turns the whole feature into a silent
no-op. Set it in REQUIRE mode only: off and allow never deny, so they must not claim to. It
is a separate knob because the gate arrives as an opaque Authenticator the server cannot
introspect, and it advertises only — enforcement is entirely the gate's.
The key id doubles as the calling proxy's label, so AuthContext.claims().get("vgi_proxy_proof")
records which proxy served each request — derived from the secret that verified, never from the
transmitted field. OPTIONS, /.well-known/ and {prefix}/health stay reachable without a proof
in every mode.
Needs no dependency beyond the JDK. The normative cross-language contract is
docs/proxy-proof-spec.md in the vgi-rpc repository.
Apache License 2.0 — Copyright 2026 Query Farm LLC · https://query.farm
