Clients

Embedded Clients

Polyglot SDKs for Java, TypeScript, Python, and Rust that keep a local mirror in sync for network-free reads.

Why embedded clients

A cache that still costs a network round-trip on every read is only half a cache. The Aeron Cache embedded clients close that gap: each client maintains a local copy of the cache data, kept in sync with the server over streaming updates, so reads are served from memory with no round-trip. Writes go to the server; the resulting ADD_ITEM / REMOVE_ITEM / CLEAR_CACHE / DELETE_CACHE / PATCH_ITEM events flow back over the stream and update the local map automatically.

That is the embedded cache pattern, and it is identical across all four languages — Java, TypeScript, Python, and Rust — with minimal external dependencies. The same high-level API also drives counters, TTL entries, JSON merge patch, bulk operations, and live subscriptions.

Transports

The same API is offered over several transports. Pick by language and latency budget:

TransportEndpointLanguagesNotes
HTTP + WebSocket:7070 (HTTP) / :7071 (WS)Java, TypeScript, Python, RustCommands over HTTP; streaming subscriptions over WebSocket.
Bidirectional WebSocket/api/ws/v1/bidiJava, TypeScript, Python, RustFull command + subscription surface multiplexed over one JSON WebSocket (AeronBidiClient).
Aeron SBE GatewayUDP :7075 / :7076 (or IPC)Java, RustCommands and streaming over Aeron using the SBE wire protocol (AeronGatewayClient).

The Aeron SBE gateway is Java- and Rust-only; the bidirectional WebSocket is the JSON analogue for TypeScript and Python (and is available in Java and Rust too).

Install

Java (GitHub Packages)

repositories {
    maven {
        url = uri("https://maven.pkg.github.com/bhf/aeron-cache-embedded")
    }
}
dependencies {
    implementation("com.aeron.cache:aeron-cache-embedded-client:1.0.0")
}

TypeScript (GitHub Packages npm registry)

Configure your .npmrc to use GitHub Packages for the @bhf scope, then install:

npm install @bhf/aeron-cache-embedded-client

Python (pip + git)

pip install "git+https://github.com/bhf/aeron-cache-embedded.git@py-v1.0.0#subdirectory=libraries/python"

Prebuilt .whl and .tar.gz artifacts are also attached to the GitHub Releases.

Rust (Cargo + git)

[dependencies]
aeron-cache-embedded-client = { git = "https://github.com/bhf/aeron-cache-embedded", tag = "rust-v1.0.0" }

Create a cache and use the embedded mirror

Construct a client with the HTTP and WebSocket URLs, create a cache, then obtain an embedded handle. Subscribe once to hydrate and keep a live local mirror; writes go to the server, and getLocal reads are served from that mirror with no network round-trip.

var client = new AeronCacheClient("http://localhost:7070", "ws://localhost:7071");
client.createCache("sample-cache");

EmbeddedAeronCache cache = client.getCache("sample-cache");
cache.subscribe(event -> {}, true);          // hydrate a live local mirror
cache.put("stay", "tuned");
System.out.println(cache.getLocal("stay"));  // local read, no network call
const client = new AeronCacheClient("http://localhost:7070", "ws://localhost:7071");
await client.createCache("sample-cache");

const cache = new EmbeddedAeronCache(client, "sample-cache");
cache.subscribe(() => {}, undefined, undefined, true); // hydrate a live local mirror
await cache.put("stay", "tuned");
console.log(cache.getLocal("stay")); // sync local read, no network call
import asyncio
from aeron_cache.client import AeronCacheClient
from aeron_cache.embedded_cache import EmbeddedAeronCache

async def main():
    client = AeronCacheClient("http://localhost:7070", "ws://localhost:7071")
    client.create_cache("sample-cache")

    cache = EmbeddedAeronCache(client, "sample-cache")
    await cache.subscribe(lambda e: None, hydrate=True)  # hydrate a live local mirror
    cache.put("stay", "tuned")
    print(cache.get_local("stay"))  # local read, no network call

asyncio.run(main())
use aeron_cache_embedded_client::AeronCacheClient;

let client = AeronCacheClient::new(
    "http://localhost:7070".into(), "ws://localhost:7071".into());
client.create_cache("sample-cache")?;

let cache = client.get_cache("sample-cache");
let _sub = cache.subscribe_ext(true)?;      // hydrate a live local mirror
cache.insert("stay", "tuned")?;
println!("{:?}", cache.get_local("stay")); // Option<String>, local read

Counters

Counter caches hold 64-bit integers and expose the same lifecycle operations plus increment, decrement, and set. An EmbeddedCounterCache shadows counter values locally over the stream.

await client.createCounterCache("counter-cache");
const counters = client.getCounterCache("counter-cache");
await counters.put("requests", 10);
await counters.increment("requests", 5); // -> 15
await counters.decrement("requests", 3); // -> 12
await counters.set("requests", 100);     // -> 100

JSON merge patch

patchItem / patch_item merges a JSON fragment into a stored value instead of replacing it, following RFC 7386: nested objects merge recursively, scalars and arrays replace, and a null field deletes it.

await client.putItem("sample-cache", "doc", '{"a":1}');
await client.patchItem("sample-cache", "doc", '{"b":2}'); // -> {"a":1,"b":2}

See JSON Merge Patch for the full semantics.

Bulk operations

Submit a batch of cache and/or counter operations in one request. A batch may freely mix regular-cache and counter ops; each operation carries its own requestId, echoed on the matching per-operation response, and results stream back in request order.

const response = await client.bulkOps({
  requestId: "batch-1",
  operations: [
    { operationType: "ADD_ITEM", requestId: "op-1", cacheId: "sample-cache", key: "k1", value: "v1" },
    { operationType: "INCREMENT_COUNTER", requestId: "op-2", cacheId: "counter-cache", key: "requests", counterValue: 5 },
    { operationType: "GET_ITEM", requestId: "op-3", cacheId: "sample-cache", key: "k1" },
  ],
});

TTL and inspection

Timed entries schedule a TTL removal you can cancel; the inspection calls list caches, aggregate stats, and all pending timers (cache + counter).

await client.putTimedItem("sample-cache", "session", "active", 60000);
await client.cancelItemRemoval("sample-cache", "session");

await client.getCaches();  // list caches
await client.getStats();   // aggregate statistics
await client.getTimers();  // all pending TTL removal timers (cache + counter)

Bidirectional WebSocket

AeronBidiClient carries the full command surface — including bulk and getTimers — plus dynamic subscribe/unsubscribe over a single persistent WebSocket connection. It is the JSON path to the full command surface for every language.

const client = new AeronBidiClient("ws://localhost:7071");
await client.createCache("bidi-cache");
await client.putItem("bidi-cache", "k", "v");
const timers = await client.getTimers();

Aeron SBE gateway

For the lowest-overhead transport, Java and Rust can speak the native Aeron SBE gateway (UDP or IPC) with AeronGatewayClient. The high-level API is transport-neutral — identical to HTTP+WS, including the embedded mirror. Enable it on the backend first (AERON_GATEWAY_ENABLED=true).

// UDP with a self-contained embedded media driver:
AeronGatewayClient client =
    AeronGatewayClient.connect(embeddedDriver.aeronDirectoryName(), "127.0.0.1");
client.awaitConnected(10, TimeUnit.SECONDS);

client.createCache("aeron-sample-cache");
client.putItem("aeron-sample-cache", "aeron-key", "aeron-value");

EmbeddedAeronCache embedded = client.getCache("aeron-sample-cache");
embedded.subscribe(event ->
    System.out.println(event.getEventType() + " " + event.getItemKey() + "=" + event.getItemValue()));

IPC skips the network entirely: the client and server share one media driver directory and must run on the same host (AeronGatewayClient.connectIpc(aeronDir)), matching the backend’s GATEWAY_TRANSPORT_MEDIA=ipc and AERON_DIR.

EmbeddedObjectCache (deep-merge)

A plain EmbeddedAeronCache stores opaque strings, so it cannot merge PATCH_ITEM deltas and ignores them (and rejects patch-mode subscriptions). When your values are JSON objects, use the embedded object cache, which deep-merges PATCH_ITEM deltas into the stored object with RFC 7386 semantics — so patch-mode subscriptions that stream only changed fields reconstruct the full object locally without losing untouched fields.

It is available in all four libraries: EmbeddedObjectCache (Java: Jackson ObjectNode, with getLocalAs to deserialize into a POJO; Python: dict; TypeScript: Record<string, any>) and EmbeddedObjects (Rust: serde_json::Value). Obtain one via getObjectCache / get_object_cache / embedded_object_cache.

EmbeddedObjectCache flags = client.getObjectCache("feature-flags");
flags.subscribe(/* patch-mode */);
MyFlags current = flags.getLocalAs("checkout", MyFlags.class);

This is the pattern behind feature-flag and dynamic-config distribution: push a one-field patch, and every embedded object cache reconstructs the merged state locally. See streaming subscriptions for keyed and patch-mode subscription details.

Reading operationStatus

Every response carries an operationStatus business status (SUCCESS, CACHE_EXISTS, UNKNOWN_KEY, UNKNOWN_CACHE, …). The clients surface these as a field to branch on rather than throwing transport-level exceptions for HTTP 400 — so a “cache already exists” or “unknown key” is ordinary control flow, not an error path.

Takeaway: write through the client, read from the mirror, and switch transports — HTTP+WS, bidi WebSocket, or native Aeron — without changing your code.