On this page
Proposal This post describes a design under review. No jobs code exists yet.
Every AI application grows background work: embedding backfills, document parsing, email, nightly feature builds, streaming aggregations. Teams already have code for it, written against Celery, BullMQ, PySpark or Flink, and a broker to run it: Redis, RabbitMQ, a Spark cluster, a Flink cluster. We do not want them to rewrite that code. The proposal is to run each family on Loam with as little change as the family allows, on one Rust core, and to make durability opt-in with Resonate's own decorators.
| Workload | How it would run on Loam | Change to your code |
|---|---|---|
| Celery | A kombu transport, a result backend and a beat scheduler, in a Python package | The broker URL and result backend become loam://… |
| BullMQ v6 | BullMQ's own classes, bound to a Loam implementation of BullMQ's pluggable backend | The import changes from bullmq to @loam/bullmq |
| PySpark | Sail, run by Loam as a Spark Connect server per namespace | SparkSession.builder.remote("sc://…") |
| Flink SQL | RisingWave, Loam's companion stream processor | A dialect port. DataStream jobs run unmodified on Apache Flink |
One core: operon-jobs
Under all four sits one Rust crate with one protobuf surface, loam.jobs.v1, served over connect-rust like the Live sync API. Python and TypeScript clients are generated from it, and the Celery and BullMQ adapters are thin layers over those clients.
#[async_trait]
pub trait Jobs: Send + Sync {
async fn enqueue(&self, cx: &Ctx, queue: &QueueId, task: TaskSpec, opts: EnqueueOpts)
-> Result<Enqueued, JobsError>;
async fn lease(&self, cx: &Ctx, req: LeaseRequest) -> Result<Vec<Lease>, JobsError>;
async fn extend(&self, cx: &Ctx, token: &LeaseToken, by: Duration) -> Result<Deadline, JobsError>;
async fn complete(&self, cx: &Ctx, token: &LeaseToken, outcome: Outcome)
-> Result<Completed, JobsError>;
async fn schedule(&self, cx: &Ctx, spec: ScheduleSpec) -> Result<ScheduleId, JobsError>;
async fn flow(&self, cx: &Ctx, graph: FlowSpec) -> Result<FlowHandle, JobsError>;
async fn submit_engine_job(&self, cx: &Ctx, job: EngineJob, opts: SubmitOpts)
-> Result<RunId, JobsError>;
fn watch(&self, cx: &Ctx, what: Selector, from: Cursor)
-> BoxStream<'static, Result<Event, JobsError>>;
// … plus bulk enqueue, progress reports, engine control and admin calls
}The namespace always comes from the caller's credential (Ctx), never from a request body. Every write that depends on holding a job carries a LeaseToken with a fencing epoch.
Where the state lives
| State | Store |
|---|---|
| Job records, ready and delayed indexes, leases, dedup keys, rate buckets, chord counters | A TiKV keyspace. A local redb file in operon dev only; every other deployment needs TiKV |
| Payloads and results over 16 KiB | The object store |
| The queue's event log | A Loam stream per queue, written through an outbox with an idempotent producer |
| Schedules, flows and durable-mode tasks | The embedded Resonate server |
Every queue operation is a small multi-key transaction: dequeue means read the ready index, lock the job and write its lease. TiKV's transactions give exactly that. No Redis, RabbitMQ or ZooKeeper is involved.
Semantics, stated plainly
- Delivery is at least once. A job goes to one worker at a time. It is delivered again if its lease expires (a crashed or stalled worker) or the worker gives it back.
- The outcome is recorded exactly once. Completion is fenced by the lease epoch. If a stalled worker and its replacement both finish a job, only the holder of the current epoch can record the result; the other gets
Fenced. A result, a retry count or a parent's child counter is never applied twice. - Effectively-once needs idempotent tasks or durable mode. A task's own side effects happen at least once in queue mode, exactly as they do on Redis or RabbitMQ today. Loam never claims exactly-once execution.
Celery
Celery has three plug points, and the package implements all three. The transport maps kombu calls onto the queue: _put is enqueue, _get is a long-polled lease, basic_ack is complete, basic_reject is release or dead-letter. The lease is a server-side visibility timeout, and the transport extends it every third of the timeout while a task runs, so a long acks_late task on a live worker is never redelivered, which removes Redis's "a task longer than the visibility timeout runs twice" problem. ETA tasks go into the delayed index instead of sitting unacked in a worker's memory. The result backend implements Celery's key-value primitives including incr, so chords use Celery's native counter in one TiKV transaction instead of the polling chord_unlock task. Beat upserts every schedule entry as a server-side schedule, so running beat twice cannot duplicate a tick.
# celeryconfig.py: the only change for a queue-mode app
broker_transport = "loam_celery.transport:Transport"
broker_url = "loam://jobs.example:7720/my-namespace"
result_backend = "loam://jobs.example:7720/my-namespace"
beat_scheduler = "loam"Some Celery features are switched on by a hard-coded list of broker types rather than by capabilities. Remote control and events work; mingle and gossip stay off.
BullMQ v6
We expected to have to emulate Redis for BullMQ. Reading the source changed the plan: BullMQ v6 has a pluggable backend interface, IQueueBackend, about 81 methods that describe every operation the Queue, Worker and Job classes need, with a Redis and a PostgreSQL implementation. So @loam/bullmq is a backend factory for the unmodified bullmq package, plus a module that re-exports BullMQ's classes already bound to it. The code users run is BullMQ's own, so the API cannot drift. Emulating Redis instead would mean matching 49 Lua scripts and 49 Redis commands, forever, and inviting every other Redis workload.
The gate is BullMQ's own backend-neutral test suite passing against Loam. It supports v6 only (v5 has no backend interface), and not BullMQ Pro's commercial features.
Durable mode, with Resonate's decorators
Queue mode runs tasks as the framework always has. Durable mode is opt-in per task, using Resonate's own SDKs, not a new decorator system. The framework keeps admission (queue, priority, rate limit, delay), and the task body becomes a Resonate function whose promise id is derived from the job id. A redelivered job then resumes the same execution: finished steps return their recorded results.
@resonate.register(name="charge", version=1)
async def charge(ctx, order_id: str):
hold = await ctx.run(reserve, order_id)
receipt = await ctx.run(capture, hold) # checkpointed: skipped on resume once recorded; capture must still be idempotent
await ctx.run(send_receipt, receipt)
return receipt
@app.task(bind=True, acks_late=True) # still a Celery task
def charge_task(self, order_id):
handle = resonate.run(f"celery:{self.request.id}", charge, order_id)
return asyncio.run(handle.result())Queue-mode jobs create no promises, which keeps millions of short jobs out of the durable store. Schedules and flows (Celery's chain, group and chord, BullMQ's parent and child flows) run on Resonate either way: a flow is a durable function whose joins are promises, so a crash mid-flow resumes from its last finished step.
PySpark on Sail
Sail is a Spark Connect server written in Rust on DataFusion. PySpark clients connect to it unchanged. Loam would run one Sail server per namespace (Sail runs Python UDFs in its own process, so sharing one across tenants would run one tenant's Python beside another's data), start it on the first connection, scale it to zero when idle, and route sc:// connections to it through an authenticating proxy. Tables live on the object store as Iceberg through Lakekeeper.
Sail does not cover everything. Spark Connect does not carry RDDs or SparkContext at all. Sail's own CI reports about 91% of Spark 3.5's Connect tests passing and about 71% of Spark 4.2's. Structured Streaming is not ready. So there is a fallback: jobs Sail cannot run go to Apache Spark on Kubernetes, unmodified, reading the same Iceberg tables. A loam spark check command would run a job on Sail and report which calls it could not handle, before a team moves anything. Resonate only wraps whole Spark actions as steps; it never enters Spark's own task retries.
Flink SQL on RisingWave
For streaming SQL, the proposal moves Flink SQL jobs to RisingWave, which Loam already plans as its companion stream processor: Rust, Apache-2.0, state in object storage, Iceberg sources and sinks. RisingWave speaks Postgres-dialect SQL, not Flink SQL, so "little change" means a porting table with a worked example per construct, not a translator. DataStream jobs have no Rust replacement; they run on Apache Flink with the Flink Kubernetes Operator, checkpointing to the object store, and Loam manages only their lifecycle (deploy, savepoint, upgrade, rollback) as durable workflows. Flink's checkpoints stay the job's only durable state. The Sail and RisingWave deep dive covers both projects.
Where this stands
| Part | Status |
|---|---|
loam.jobs.v1, the queue core, the Celery transport, result backend and beat, TiKV job store | Planned |
@loam/bullmq on BullMQ v6's backend interface, flows on Resonate promises | Planned |
| PySpark on Sail per namespace, the Spark Connect proxy, the Spark fallback | Planned |
| Flink SQL on RisingWave, DataStream lifecycle on the Flink operator | Planned |
The first phase is Celery, because it exercises the whole core: leases, fencing, delays, retries, dead letters, schedules on Resonate and the TiKV store. Its gate is Celery's own integration suite passing against operon dev. Jobs depend on the embedded durable server and the TiKV client layer, both of which exist. The next post covers the other databases teams bring: Postgres and MySQL.