Skip to content

Repository files navigation

aerolite ☄

CI License: MIT Rust TypeScript

A Meteor-style realtime server in Rust — write your app in TypeScript. The classic DDP surface (live collections, methods, publications, and the wire protocol stock Meteor clients already speak) plus two primitives Meteor never had, durable objects and workflows, on a Raft-replicated core with automatic failover. One binary: server, CLI, an htop-style cluster dashboard, and built-in OpenTelemetry — traces and metrics as local Parquet, queried by the embedded DuckDB (aerolite otel top), mirrored to S3 for fleets. Faster than Meteor 3.5 on every scenario of Meteor's own benchmark suite, on roughly a quarter of the CPU and memory — against Meteor's best configuration, while fsyncing every acked write.

┌────────────┐  DDP over WebSocket   ┌──────────────────────────────────┐
│ any DDP    │ ◄───────────────────► │     aerolite (axum + tokio)      │
│ client     │  sub/method/added/…   │                                  │
└────────────┘                       │  ┌────────────┐ ┌─────────────┐  │
┌────────────┐  HTTP /rpc, /stats    │  │ Collections│ │  Methods /  │  │
│ aerolite   │ ◄───────────────────► │  │ (live)     │ │ Publications│  │
│ CLI / top  │                       │  └─────┬──────┘ └─────────────┘  │
└────────────┘                       │  ┌─────┴──────┐ ┌─────────────┐  │
                                     │  │  Durable   │ │  Workflow   │  │
                                     │  │  Objects   │ │  Engine     │  │
                                     │  └─────┬──────┘ └──────┬──────┘  │
                                     │  ┌─────┴───────────────┴──────┐  │
                                     │  │ Raft log (cluster mode)    │  │
                                     │  └────────────┬───────────────┘  │
                                     │       sled (embedded db)         │
                                     └──────────────────────────────────┘

Quickstart

cargo install --path .    # installs the `aerolite` binary

aerolite init my-app      # scaffold: manifest + TS backend + React client
cd my-app
aerolite run              # boots the server AND your backend (server/app.ts)

aerolite call greet '"world"'
aerolite workflow start signup '{"email":"ada@example.com"}'
aerolite otel top         # query your own telemetry

cd client && npm install && npm run dev    # the React client (Vite + DDP)

Like Meteor, aerolite run serves the project containing the current directory (it finds aerolite.toml by walking up) and refuses to run elsewhere. Your backend is plain TypeScript — server/app.ts, hosted as the server's app worker with the SDK vendored in (zero npm installs, Node >= 23.6). The scaffolded client/ is a standard Vite + React app on the vendored aerolite client library — live queries in, optimistic writes out — start a signup workflow from the CLI and watch its welcome message stream into the page. Server state lives project-locally in .aerolite/data, telemetry in otel/. (aerolite init --server-only skips the client; the built-in demo backend used throughout this README stays available anywhere via aerolite run --demo.)

The client library

client/ is the isomorphic client — dependency-free ESM that runs identically in browsers and Node >= 21, vendored into every scaffold (client/src/lib/). A first-party DDP client, a local minimongo cache, and React hooks:

useSubscribe('messages');
const messages = useFind<Message>('messages', { room: 'general' });
// Optimistic, no stub: renders instantly, reconciles on the server's ack.
client.collection('messages').insert({ text, sentAt: Date.now() });

Three things Meteor's client stack doesn't give you:

  • Provable query parity. The client cache runs a JS port of the server's query engine, and both are pinned to one shared vector file (tests/fixtures/query-vectors.json, 109 cases) that both CI suites must pass — the selector that filters your publication server-side is the selector that filters your useFind client-side, guaranteed by test, not by convention. (Known, vector-excluded divergences: the regex x option and non-ASCII string ordering.)
  • Latency compensation with zero stubs. Collection writes are protocol builtins the client understands, so ANY insert/update/remove applies to the local cache before the wire send and rolls back on rejection — no hand-written method simulation, which is where Meteor apps usually give up on optimistic UI. The server guarantees a method's updated message follows its own writes (the updated barrier), so an optimistic document reconciles without ever flickering.
  • External-store reactivity, no Tracker. Live queries are refcounted, deduped by content, maintained per-document (no re-scan per change), and read through useSyncExternalStore with identity-stable snapshots — React-native reactivity, StrictMode-safe.

Subscriptions are refcounted and deduped by (name, params); reconnect re-subscribes everything and re-sends unacked method calls (methods are therefore at-least-once across reconnects — same trade-off Meteor makes). Durables and workflows are first-class: client.durable('counter', 'c1').call('increment', 2), client.workflow.start('signup', {email}).

CLI

aerolite init my-app        # scaffold a new project
aerolite run                # this project's backend, foreground, 127.0.0.1:3000
aerolite run --workers 4    # scale app code across 4 worker processes (round-robin)
aerolite run --no-worker    # host workers yourself; lock with AEROLITE_WORKER_TOKEN
aerolite run --demo         # the built-in demo backend (works anywhere)
aerolite run 3              # managed local 3-node Raft cluster on :3000-:3002
                            #   (Ctrl+C stops it; `aerolite scale`/`stop` from elsewhere)
aerolite scale 5            # resize the managed cluster (data persists)
aerolite stop               # stop it (data kept)

aerolite status             # node role, term, leader
aerolite top                # live htop-style cluster dashboard
aerolite top --snapshot     # one plain-text frame (scripts/CI)

aerolite call tasks.insert '"buy milk"'          # invoke any server method
aerolite object call counter demo increment 5    # durable object call
aerolite workflow start order.fulfill '{"sku":"x"}'
aerolite workflow status <run-id>
aerolite workflow list
aerolite workflow remove <run-id>                # abort a stuck run, delete record+journal

aerolite collection list                         # generic collection CRUD
aerolite collection insert todos '{"text":"hi"}'
aerolite collection docs todos
aerolite collection find todos '{"done":false,"score":{"$gte":10}}' --project '{"secret":0}'
aerolite collection update todos <id> '{"done":true}'
aerolite collection remove todos <id>

aerolite otel top                                # telemetry: latency/error summary
aerolite otel query "SELECT ... FROM spans"      # SQL over your telemetry (embedded DuckDB)
aerolite otel sync --to s3://bucket/otel         # mirror telemetry to S3 now
aerolite otel prune --older-than-days 1          # delete old telemetry files

Data directory: project-local by default — aerolite run inside a project keeps server state in <project>/.aerolite/data (the .meteor/local analog). With --demo (and for managed clusters) it's the platform dir — ~/Library/Application Support/aerolite on macOS, $XDG_DATA_HOME/aerolite (or ~/.local/share/aerolite) elsewhere; override with --data or AEROLITE_DATA. A managed cluster keeps each node in node-<i>/ plus node-<i>.log and a cluster.json state file under that root.

aerolite top targets any node (--addr) and discovers the rest of the cluster from it. It shows per-node role, term, commit/applied/log indexes, latency, uptime, write rate with a burst-activity strip (idle polls are blank), a durable objects panel (live flag, call count, last call), and a live workflow table. A running run that no node is executing is flagged STUCK in red — select it with ↑/↓ and press x twice to abort and remove it. Keys: ↑↓ select, x remove, p pause, r refresh, q quit.

The bottom logs panel merges every node's recent log events — timestamped, node-tagged, level-colored, with structured fields — tailed live. Each node keeps a ring buffer of its own tracing events (down to DEBUG for aerolite's own modules) served at GET /logs?since=<seq>; messages are written to be actionable ("commit timed out — check that a majority of members are reachable"). Press Tab to focus the panel, then ↑↓/PgUp/PgDn to scroll, End to re-follow the tail, Home for the top, and v to cycle the minimum level (DEBUG → INFO → WARN). Stream discontinuities are marked explicitly: "node restarted — log stream reset" and "… N log lines missed (node's buffer wrapped)".

aerolite scale note: Raft membership is static, so scaling restarts the cluster with the new member list (a few hundred ms of unavailability). Data dirs persist and new nodes catch up from the leader's log.

The collection.* methods are an open write surface (like Meteor's insecure package) — great for dev, gate or remove before exposing publicly.

CLI commands hit the node's HTTP POST /rpc endpoint ({method, params}), which — like DDP sessions — transparently forwards to the cluster leader.

aerolite run --demo serves the built-in demo app: open http://127.0.0.1:3000 for a live task list, a persistent counter, and a workflow runner over DDP. Kill the server and restart it: everything comes back. (The demo backend is also what the benchmarks and the React test app talk to.)

cargo test

The classic Meteor surface

  • DDP 1 wire protocol at ws://host/websocket: connect/connected, ping/pong, method/result/updated, sub/ready/nosub, added/changed/removed. Existing DDP clients can connect.
  • Collections (server.collection("tasks")): persistent document stores. Every insert/update/remove is broadcast to subscribed sessions in real time, and find(selector) queries with the same Mongo-subset language as publication filters (from Rust, TS via tasks.find({score: {$gte: 10}}, {secret: 0}), or the CLI) — with malformed selectors rejected loudly, since a query API should error rather than silently match nothing.
  • Methods (server.method(name, handler)): async RPC handlers, errors in Meteor.Error {error, reason} shape.
  • Publications (server.publish(name, handler)): map subscription params to what a client receives — whole collections, or per-document filtered views: Mongo-style matches with dot-path field access ("profile.address.city", numeric segments index arrays), equality or operator conditions ($eq $ne $in $nin $gt $gte $lt $lte $exists $regex (+$options, linear-time — safe for untrusted patterns) $elemMatch $not, AND-ed within a field), and nestable $or/$and/$nor combinators — plus field projections ({"secret": 0}, {"text": 1}, or dotted subtrees like {"profile.city": 1}; overlapping subs deep-merge their subtrees). Malformed conditions and projections fail closed — including inside negations, where arguments are validated first so a malformed inner can never invert into match-everything. Visibility is tracked per session by a merge-box: a document updated out of a filter arrives as removed, one updated into it as added with (projected) full fields, and overlapping subscriptions merge — the client sees the union of their projections, broadened by changed when a sub widens it and narrowed by changed {cleared} when one goes away; removed only when the last provider does. A patch touching only hidden fields sends nothing.
let server = Server::new("./data");
server.publish("tasks", |_s, _p| async { Ok(vec!["tasks".into()]) });
server.publish("my.tasks", |_s, params| async move {
    let owner = str_param(&params, 0)?;
    Ok(vec![PublishSpec::filtered("tasks", json!({"owner": owner}))])
});
server.method("tasks.insert", |server, params| async move {
    let text = str_param(&params, 0)?;
    Ok(json!(server.collection("tasks").insert(json!({"text": text}))))
});
server.method("tasks.open", |server, _params| async move {
    // find(): the same selector language publications filter with
    Ok(json!(server.collection("tasks").find(&json!({"done": false}))))
});
server.listen("127.0.0.1:3000").await

New primitive: durable objects

A durable object is a named, stateful, single-threaded actor (in the spirit of Cloudflare Durable Objects):

  • At most one live instance per (class, id); all calls go through a mailbox, so handlers never race — no locks needed in your code.
  • After every successful call the object's snapshot() is persisted to sled; the factory rebuilds the object from that snapshot on next use, including after a process restart.
struct Counter { count: i64 }

#[async_trait]
impl DurableObject for Counter {
    async fn call(&mut self, method: &str, params: Value) -> Result<Value, Error> {
        match method {
            "increment" => { self.count += params.as_i64().unwrap_or(1); Ok(json!(self.count)) }
            "get" => Ok(json!(self.count)),
            other => Err(Error::new("durable-no-method", other)),
        }
    }
    fn snapshot(&self) -> Value { json!({"count": self.count}) }
}

server.durable_class("counter", |_id, state| Box::new(Counter {
    count: state.get("count").and_then(Value::as_i64).unwrap_or(0),
}));

Clients reach objects through the built-in method durable.call(class, id, method, args).

New primitive: workflows

A workflow is a journaled, replayable async function (in the spirit of Temporal). Each ctx.step(name, fut)'s result is journaled exactly once: when a run is replayed — e.g. resumed after a crash — journaled steps return their recorded result without re-executing. ctx.sleep(duration) journals its wake deadline, so a resumed run sleeps only for the remaining time.

Delivery guarantee — read this before calling external systems from a step. Step side effects are at-least-once, not exactly-once: the step's future runs before its journal entry commits, so a crash or leadership loss in that window re-executes the step on resume, after its external effect may already have happened (the same contract as Temporal activities). Make step bodies idempotent, using ctx.step_key(name) — stable across replays, unique per (run, step position) — as the idempotency key for external calls:

let key = ctx.step_key("charge-payment");
let charge: Value = ctx.step("charge-payment", async move {
    payments.charge(4200).idempotency_key(key).send().await
}).await?;
server.workflow("order.fulfill", |mut ctx, input| async move {
    let reserve: Value = ctx.step("reserve-inventory", async { /* … */ Ok(json!({"reserved": true})) }).await?;
    let charge: Value  = ctx.step("charge-payment",   async { /* … */ Ok(json!({"charged": true})) }).await?;
    ctx.sleep(Duration::from_secs(60 * 60)).await?;   // durable pause
    let ship: Value    = ctx.step("ship-order",       async { /* … */ Ok(json!({"tracking": "TRK-1"})) }).await?;
    Ok(json!({"reserve": reserve, "charge": charge, "ship": ship}))
});

Run state lives in the _workflowRuns live collection, so clients can subscribe to run progress (workflow.runs publication) like any other data. Built-in methods: workflow.start(name, input) -> runId and workflow.status(runId). On startup the engine re-spawns every run that was running when the process stopped, replaying its journal.

Write your app in TypeScript

You don't have to write Rust. An app worker is an external Node process that hosts your methods, durable object classes, and workflow functions in TypeScript (or plain JS), registered over a WebSocket to /worker:

// app.ts — run with: node app.ts   (Node >= 23.6; SDK in sdk/)
import { app } from "aerolite";

const trips = app.collection("trips");

app.method("greet", async (name: string) => `hello ${name}`);

app.publish("trips", ["trips"]);                 // whole collection
app.publish("my.trips", (owner) => [             // per-document filter
  { collection: "trips", filter: { owner } },
]);

app.durable("visits", class {
  count = 0;                       // instance fields ARE the persisted state
  hit() { return ++this.count; }
});

app.workflow("trip.plan", async (ctx, input) => {
  const flight = await ctx.step("book-flight", (key) => bookFlight(input.city, key));
  await ctx.sleep(60 * 60 * 1000);            // durable pause
  const hotel = await ctx.step("book-hotel", () => bookHotel(input.city));
  return { tripId: await ctx.step("save", () => trips.insert({ ...flight, ...hotel })) };
});

await app.serve("ws://127.0.0.1:3000/worker"); // pass every node's URL in a cluster

The division of labor is strict: the worker hosts code, the core keeps every guarantee. Durable object state is core-owned (shipped out per call, snapshotted back through the same mailbox/lease path as Rust objects); workflow ctx.step / ctx.sleep are round-trips to the core, which owns the journal, cursor, and epoch/lease fences — so replay, exactly-once results, at-least-once side effects (with key as the idempotency key), and failover semantics are identical to native workflows. If the worker process crashes mid-run, the run stays running and resumes with replay as soon as a worker re-registers; with no worker connected, its methods fail fast with retryable worker-unavailable.

Lifecycle and scale. aerolite run hosts your worker with two containment guarantees: each child holds a stdin pipe that closes when the server dies — by any means, SIGKILL included — so the SDK exits on EOF and orphans cannot exist; and the server mints a per-run token its workers must present, so a stale worker from another project is rejected loudly instead of silently replacing your app. --workers N spawns a pool and round-robins calls across it — safe by construction, because methods are stateless and durable serialization / workflow journals live in the core, so any worker may execute any turn (a Node process is one core; the pool is the single-box scaling knob). For deployment, --no-worker skips spawning so workers run and scale as their own processes — anywhere — with registration open on a trusted network or locked to a shared AEROLITE_WORKER_TOKEN on both ends.

Clustering — Raft with failover

Run 3+ nodes as a Raft cluster with static membership — each node lists its peers, and every node must be configured with the same total member set:

# the SAME --cluster list on every node; each node filters itself out
CLUSTER=10.0.0.1:3000,10.0.0.2:3000,10.0.0.3:3000
aerolite run --addr 10.0.0.1:3000 --cluster $CLUSTER   # node 1
aerolite run --addr 10.0.0.2:3000 --cluster $CLUSTER   # node 2
aerolite run --addr 10.0.0.3:3000 --cluster $CLUSTER   # node 3

(--peers also exists for explicit per-node peer lists. --advertise overrides the address peers use to reach a node if it differs from the bind address — with --cluster, the advertise address is what must appear in the member list. Every flag also has an AEROLITE_* env var. Watch the whole cluster with aerolite top. Programmatically: Server::new_clustered(data_dir, ClusterConfig { self_addr, peers }). GET /raft/status on any node shows its role, term, and current leader.)

The cluster runs a built-in Raft implementation (src/raft.rs): leader election plus a replicated command log persisted in each node's sled db. Any minority of nodes can crash without losing state or stranding execution — kill -9 the leader and the survivors elect a new one that carries on.

How it works:

  • Everything is a committed log entry. Collection changes, durable object snapshots, and workflow journal writes are proposed to the Raft leader and applied — in the same order, idempotently — on every node once a quorum has them. A mutation that returns success has survived failover by construction.
  • The leader executes, followers replicate and serve reads. DDP clients can connect to any node: subscriptions are served from the local replica, method calls are transparently proxied to the leader (with retry across elections).
  • Durable objects: live instances exist only on the leader; each call's snapshot must commit through Raft before the caller sees success. After a failover the new leader rebuilds instances from replicated snapshots — the single-instance, serialized-calls guarantee holds across crashes.
  • Workflows: the run record commits before execution starts, and each step's journal entry commits before the workflow proceeds. A newly elected leader resumes every in-flight run, replaying journaled steps (they are not re-executed) and honoring the remainder of durable sleeps. A deposed leader can no longer commit journal entries, so it cannot corrupt a run the new leader has taken over.

Test app (React)

webapp/ is a Vite + React front-end that exercises the whole surface over plain DDP using the off-the-shelf simpleddp client — no custom protocol code. Live task list (pub/sub), the counter durable object, and a workflow runner with a live run table, plus a node picker and connection status. Methods called against any node transparently reach the Raft leader, so you can kill -9 the leader mid-click and writes keep working.

cd webapp && npm install && npm run dev   # aerolite cluster on :3000 assumed

Persistence

Everything — collection documents, durable object state, workflow journals and run records — lives in one embedded sled database under the data directory. No external database required.

Collection writes are stored per document and acknowledged through a group-committed write-ahead log: concurrent writers share one sync instead of each paying their own, and --durability picks what an acknowledged write survives:

Mode An acked write survives Sync call (macOS) Cost/sync
full sudden power loss F_FULLFSYNC ~4ms
barrier (default) process/OS crash, writes ordered to the drive F_BARRIERFSYNC ~0.4ms
os process/OS crash fsync ~0.02ms

barrier is stronger than MongoDB's or Postgres's defaults (both ack at OS-crash durability). On startup the WAL is replayed into sled and truncated; a torn tail from a crash mid-append is detected by checksum and discarded. In full mode the WAL is bypassed and sled itself is flushed — the pre-WAL behavior, byte for byte.

Benchmarks

Measured with Meteor's own benchmark suite

The strongest version of this comparison wasn't ours to design, so we didn't: bench/official runs github.com/meteor/performance — Meteor's own suite — against both servers. Its load generators, load profiles, CPU/RAM collector and result format are used unmodified, and aerolite is driven by simpleddp, the same DDP client the suite points at Meteor. All we supply is an app implementing the same DDP surface (a literal port of the suite's tasks-3.x app code) and a driver that can start either server — one driver, so the load and the measurement are shared code paths rather than two implementations kept honest by hand.

Meteor is measured at its best. Meteor 3.5 defaults to Mongo change streams; its classic oplog tailing is far faster at fanout, so both are reported and the comparison column is oplog.

Apple M1 Max, single node, loopback. Median of 3 repetitions — 36 runs, zero failed virtual users. Best value bolded:

Scenario (the suite's own) Meteor 3.5 (change streams) Meteor 3.5 (oplog) aerolite
fanout-light — p50, 50 subscribers 109.44 ms 5.06 ms 3.65 ms
fanout-heavy — p50, 200 subscribers 109.06 ms 13.32 ms 8.85 ms
ddp-non-reactive-light — VU session p50 90.9 ms 90.9 ms 79.1 ms
ddp-reactive-light — VU session p50 4,316.6 ms 117.9 ms 76.0 ms
CPU, whole stack under reactive load 18.3 % 25.3 % 7.0 %
RAM, whole stack under reactive load 411 MB 339 MB 93 MB

Both stacks are two processes and both are measured whole — Meteor's node bundle plus mongod, aerolite's core plus its node app worker.

Read honestly: aerolite leads every scenario against Meteor's best configuration, by 1.4–1.6× on latency and ~3.6× on both CPU and memory, while group-committing an fsync before acking every write where Mongo at w:1 acknowledges from memory. The change-streams column is Meteor's shipping default rather than a strawman — its ~109 ms fanout is flat whether 50 or 200 subscribers are listening (a fixed batching interval, not a scaling limit), and the suite's own README quotes the same ~4.2 s reactive session we measured. Where the observer driver cannot matter — methods with no subscription — Meteor's two configurations land on the identical 90.9 ms, which is the control that says the rest of the table is measuring what it claims to.

Reproduction steps, and the caveats that would move these numbers, are in bench/official/README.md; the charted report is bench/official/report.html.

Our own microbenchmarks

Single node, same machine, same client library (simpleddp), identical workloads, three ways: Meteor 3.5 + Mongo, the native Rust demo backend, and the same surface hosted in a TypeScript app worker (bench/ts-app) — the same app-code-in-JS shape as Meteor, with the worker-protocol hop priced in.

We benchmark Meteor's best configuration. Meteor 3.5's default reactivity engine (Mongo change streams) measured 97 ms fanout p50 on this workload — a ~25× regression against its classic oplog tailing (3.8 ms), matching reports upstream — so the numbers below pin Meteor to oplog (--settings settings-oplog.json; both configs ship in bench/meteor-app so you can reproduce either). Full harness and interactive report in bench/; MacBook (Apple silicon), best value bolded:

Scenario Meteor 3.5 (oplog) aerolite (rust) aerolite (TS worker)
echo latency p50 1.58 ms 1.41 ms 1.47 ms
echo throughput (32 clients) 17,291 ops/s 20,361 ops/s 18,960 ops/s
insert throughput (16 clients, durable) 5,683 ops/s 7,884 ops/s 6,492 ops/s
pub/sub fanout p50 (20 subscribers) 3.82 ms 4.38 ms 3.32 ms
subscribe-ready p50 (1000 docs) 10.13 ms 10.91 ms 11.71 ms

Read honestly: aerolite wins method latency, echo throughput, and durable insert throughput (the last durability-adjusted in aerolite's favor — both aerolite variants group-commit a barrier fsync before acking every insert, while Mongo at w:1 acknowledges from memory); fanout and subscribe-ready are at parity, trading a millisecond either way run to run. Against Meteor's stock 3.5 defaults, aerolite's fanout advantage is ~25× — but that's Meteor's regression, not our win, so it lives in this footnote rather than the table. The report also charts CPU and RSS over the run (the TS variant's node worker included) — native aerolite peaks around 100 MB where Meteor's node + mongod hold ~350 MB.

Observability — OpenTelemetry, queryable with DuckDB

The server is fully instrumented with OpenTelemetry, and the pipeline is Parquet-native: instead of shipping to a collector, traces and metrics land as complete Parquet files in a project-local directory you can query directly:

otel/
  spans/spans-<ts>-<pid>-<n>.parquet      # one row per span
  metrics/metrics-<ts>-<pid>-<n>.parquet  # one row per metric interval (DELTA)

Spans: aerolite.method (every RPC, with rpc.method + status), aerolite.durable.call, aerolite.workflow.run, aerolite.workflow.step, aerolite.worker.request. Histograms mirror every span, plus aerolite.collection.commit and aerolite.storage.flush for the write path, aerolite.wal.bytes, and gauges for active sessions, in-flight workflow runs, worker connectivity, Raft term, and process RSS/CPU. Attributes are a JSON column; files are written every ~30s (or 8k rows) and on shutdown, so they are always complete and always queryable.

aerolite run                       # telemetry on by default -> ./otel
aerolite run --otel-dir /tmp/otel  # elsewhere
aerolite run --no-otel             # off

DuckDB is embedded — in-process, like SQLite — so the same binary queries its own telemetry with no external install. spans and metrics are pre-registered views over the parquet globs:

$ aerolite otel top          # per-operation latency/errors + latest gauges
name                    calls  avg_ms  p50_ms  p99_ms  errors
----------------------  -----  ------  ------  ------  ------
aerolite.method         24     0.51    0.27    4.12    1
aerolite.workflow.step  3      4.8     4.37    6.01    0
aerolite.durable.call   1      4.89    4.89    4.89    0
aerolite.workflow.run   1      768.11  768.11  768.11  0

$ aerolite otel query "SELECT attributes->>'rpc.method' m, count(*) n \
                       FROM spans WHERE name='aerolite.method' GROUP BY 1"

Cleanup is automatic, the way TSDBs and log rotation do it: the sink prunes files older than the retention window (default 7 days, --otel-retention-days, 0 = keep forever) as part of its flush cycle. aerolite otel prune [--older-than-days N] cleans up on demand.

Or point any external DuckDB at the same files:

-- Where does time go? (per operation, with tail)
SELECT name, count(*) calls, round(avg(duration_ms), 2) avg_ms,
       round(quantile_cont(duration_ms, 0.99), 2) p99_ms
FROM 'otel/spans/*.parquet' GROUP BY 1 ORDER BY calls DESC;

-- Slowest individual method calls, with which method
SELECT attributes->>'rpc.method' AS method, duration_ms, status
FROM 'otel/spans/*.parquet' WHERE name = 'aerolite.method'
ORDER BY duration_ms DESC LIMIT 20;

-- Error rate per method over time (1-minute buckets)
SELECT time_bucket(INTERVAL 1 minute, start_time) minute,
       attributes->>'rpc.method' AS method,
       count(*) FILTER (status = 'error') * 100.0 / count(*) AS error_pct
FROM 'otel/spans/*.parquet' WHERE name = 'aerolite.method'
GROUP BY 1, 2 ORDER BY 1;

-- fsync amortization: how many writers shared each group flush?
SELECT sum(count) flushes, round(sum(sum)/sum(count), 3) avg_flush_ms
FROM 'otel/metrics/*.parquet' WHERE name = 'aerolite.storage.flush';

Metrics use DELTA temporality, so every row is a self-contained interval — plain sum() / avg() are correct with no windowing tricks.

Deployment: mirror to S3

For deployment, the same directory mirrors to S3 (or any S3-compatible store — MinIO, R2). Files are write-once with globally unique names, so sync is a pure "upload what the remote prefix doesn't have yet": idempotent under crashes, retries, and multiple nodes.

# continuous: completed files upload right after each flush
AWS_ACCESS_KEY_ID=… AWS_SECRET_ACCESS_KEY=… AWS_DEFAULT_REGION=us-east-1 \
  aerolite run --otel-s3 s3://my-bucket/otel

# one-shot backfill (same idempotent pass)
aerolite otel sync --to s3://my-bucket/otel --node 127.0.0.1:3000

Files land under <prefix>/<node>/{spans,metrics}/, one prefix per node. For MinIO/R2 set AWS_ENDPOINT (and AWS_ALLOW_HTTP=true for plain http). A transient outage self-heals: the uploader retries every minute and on every flush. Query the whole fleet straight from the bucket with DuckDB's httpfs — the schema needs no changes:

INSTALL httpfs; LOAD httpfs;
SELECT node, name, count(*), round(avg(duration_ms), 2) AS avg_ms
FROM 's3://my-bucket/otel/*/spans/*.parquet' GROUP BY 1, 2;

Pair with a bucket lifecycle rule for remote retention; local retention (--otel-retention-days) keeps the on-disk window short.

Limitations (by design, for now)

  • Publication filters implement Mongo's core query surface (equality, comparisons, membership, $regex, $elemMatch, $not/$nor, $or/$and, dot-paths) but not its long tail ($type, $mod, $size, geo/text operators), and array fields compare by strict equality — use $elemMatch for element semantics.
  • Unindexed find() is an in-memory full scan. It scales linearly (measured 100k → 1M docs: ~6.7M docs/s equality, ~2.3M docs/s $regex, Apple silicon; cargo run --release --example find_bench) — fine for dashboards and CLI queries, but hot paths should either subscribe (publications evaluate per change, no scans) or declare a secondary index: ensure_index("owner") (Rust), tasks.ensureIndex("owner") (TS), aerolite collection ensure-index (CLI). Indexes serve equality/$in/range selectors from an ordered in-memory structure with every candidate re-verified by the matcher — identical results, and cost that scales with the RESULT SIZE instead of the collection: at 1M docs, selective equality is 166 ms scanned vs 1.3 ms indexed (125×), $in 280 ms vs 6.3 ms (45×), ranges 330 ms vs 20 ms (17×, materializing 10k matches) — ~6.7 ms end to end through the method layer. Indexes are in-memory and declared at startup (rebuilt from loaded docs, not persisted); $regex/$elemMatch/negations still scan.
  • There is no authorization layer, and the data plane is open to any client that can reach the port. This is bigger than "no accounts yet", so read it carefully before exposing a node. The core registers builtin methods — collection.list / find / docs / insert / update / remove / ensureIndex, durable.call, workflow.start / remove — and any caller may invoke them, over DDP or an unauthenticated POST /rpc. Concretely, curl -X POST .../rpc -d '{"method":"collection.find","params":["users",{}]}' reads your users collection. Three consequences people get wrong: (1) publications gate subscriptions but not collection.find, so a filtered publication does not keep a client out of the underlying documents, and the same methods reach internal collections such as _workflowRuns; (2) you cannot write the authorization yourself yet — method and publication handlers receive (server, params) and no caller identity, so a handler cannot tell who is calling it; (3) consequently per-user publications are not possible today — Meteor's Docs.find({owner: this.userId}) has no equivalent here, because a publication sees only client-supplied params, so a client can subscribe with someone else's id. Shipping accounts alone would therefore not close this: threading a caller context through the request path has to come first. Meteor's nearest equivalent (the insecure package) can at least be removed and replaced with allow/deny rules — aerolite has no such hook yet. The default bind is loopback (127.0.0.1:3000), and that is currently the only thing standing between a deployment and the open internet: do not bind a public interface without putting your own authenticating proxy in front. Also unauthenticated: GET /logs, which returns the in-memory log buffer (method names, errors, and whatever your app put in log fields). See SECURITY.md.
  • No EJSON extended types, no server-to-server DDP.
  • Workflow steps must be deterministic in sequence (same step names in the same order on replay); the engine detects and fails nondeterministic replays.
  • Cluster membership is static: adding or removing nodes requires reconfiguring and restarting every node. Writes need a quorum, so a majority of nodes must be up.
  • The Raft log has no compaction or snapshot transfer yet — it grows with every write, and a node that was down catches up by replaying the log.
  • Workflow steps are exactly-once in the journal but at-least-once in side effects: a step's body may run on a deposed leader whose journal write then fails to commit, and again on the new leader (the same semantics as Temporal activities — make step bodies idempotent).
  • The inter-node /raft/* and /cluster/* endpoints are unauthenticated — run the cluster on a trusted private network. (/worker IS authenticated: locked to a minted token under aerolite run, and to AEROLITE_WORKER_TOKEN with --no-worker when set.)
  • Cluster writes execute on the Raft leader (followers proxy), so write throughput is bounded by one node no matter how many nodes or workers you add; reads and subscription fanout scale out per node. Sharding durables/collections across Raft groups is future work.

Inspiration & prior art

aerolite stands on ideas from projects we admire:

  • Meteor — the DDP protocol, live collections, methods, publications, and the conviction that realtime should be the default, not a bolt-on.
  • Temporal — journaled workflow replay, exactly-once step results with at-least-once side effects, and the worker model that keeps application code out of the consensus core.
  • Cloudflare Durable Objects — the named single-threaded stateful actor as a first-class primitive.
  • sled — the embedded storage engine under everything.
  • The Raft paper (Ongaro & Ousterhout) — the consensus algorithm, including the §4.2.3 leader-stickiness refinement.

License

MIT — see CONTRIBUTING.md to get involved.

About

Meteor-style realtime server in Rust — DDP, live collections, durable objects, journaled workflows, Raft failover. Write your app in TypeScript.

Topics

Resources

Contributing

Security policy

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages