§26 Loam Jobs: Celery, BullMQ, PySpark and Flink Jobs on Loam

Existing jobs on Loam with minimal changes: `operon-jobs` and the `loam.jobs.v1` connect-rust surface (enqueue, lease with fencing, complete, schedule, flow, engine jobs, watch); at-least-once semantics, idempotency, priorities, rate limits, DLQs; the Celery kombu transport, result backend and beat scheduler; `@loam/bullmq` as a BullMQ v6 `IQueueBackend`; Resonate for schedules, flows and opt-in durable tasks; PySpark on Sail with a Spark fallback; Flink SQL on RisingWave and DataStream under the Flink operator; track J (D204–D216; numbers left a gap for docs 23–25)

Status: Proposed · 2026-09-29. The direction is the owner's, from 2026-09-29: "The goal is to deploy PySpark, Flink, BullMQ, Celery jobs on Loam and run it, with some modification like adding Resonate decorators ... I want to expose clean Rust API for this." The owner then approved the per-workload direction in §1 ("yes", 2026-09-29). This document turns that direction into decisions D204–D216 and open questions Q95–Q104. The owner decided Q97 and Q100 on 2026-09-29: jobs require TiKV on every deployment except operon dev, and @loam/bullmq ships for BullMQ v6 only. The choices it makes on top of the approved direction (the semantics, the state layout, the API details, the phasing) are proposals until the owner confirms them. No code is written by this document; operon-jobs is a design.

Numbering. On main the highest decision is D147 and the highest question Q44. The showcase suite's D-SC-1…16 and Q-SC-1…10 and the Postgres-wire D-PG-1 get numbers at merge, so they may take D148–D164 and Q45–Q54. Docs 23 (Neon/WeSQL) and 24/25 (CPU-time runtime, Clever Cloud stack) are being written on unpushed branches and pick numbers too, so this document leaves a gap of 20 for each: it starts at D204 (164 + 20 + 20) and Q95 (54 + 20 + 20). Renumber at merge if the gap was too small or too large.

Markers: (source) means read in the upstream repository at the revision in §17. (estimate) means computed, not measured. (verify) means the plan that builds it checks it first. "The research" is the data-processing research of 2026-09-29 (.superpowers/research/rust-data-processing-2026-09.md, not committed), which §9 and §10 follow.


1. Summary (D204)

Loam runs four families of jobs that teams already have, with as little change to their code as each family allows, on one Rust core:

WorkloadHow it runs on LoamChange to user codeWhere Resonate fits
CeleryA kombu transport (loam://), a result backend and a beat scheduler, in the Python package loam-celery, over Loam's queue (TiKV state, Loam streams as the queue's event log, Resonate for schedules and flows)The broker URL becomes loam://… and the result backend loam://…; optionally @app.task becomes a Resonate functionVery good: durable retries, beat → Resonate schedules, chain/group/chord → durable steps and fan-out
BullMQA drop-in TypeScript package @loam/bullmq: BullMQ's own Queue, Worker, FlowProducer, QueueEvents and Job classes, bound to a Loam implementation of BullMQ v6's pluggable IQueueBackend (§7.2). Redis is not emulated (BullMQ's Redis backend is 49 Lua scripts, §7.2.4)The import changes from bullmq to @loam/bullmqVery good: jobs → durable functions, flows → parent/child promises, repeatable jobs and job schedulers → Resonate schedules
PySparkSail (lakehq/sail), run by Loam as a managed Spark Connect server, with Iceberg tables on RustFS through Lakekeeper. Sail is a separate process, never linked (D51)SparkSession.builder.remote("sc://…")Job level only: step checkpoints, retries and schedules around Spark actions, never inside Spark tasks. RDD jobs and unsupported UDFs fall back to Apache Spark on Kubernetes
FlinkFlink SQL on RisingWave (default) or Arroyo. DataStream jobs run unmodified on Apache Flink (the Kubernetes operator, state on RustFS), reading Loam streams through the Kafka gateway (M5, D74)SQL: small dialect edits; DataStream: none (it stays on the JVM)Only the job lifecycle (deploy, savepoint, upgrade, rollback). Flink's checkpoints stay the only durability of a running job

Under all four sits one crate, operon-jobs, with one protobuf surface, loam.jobs.v1, served over connect-rust (D128). Python and TypeScript clients are generated from it, and the Celery and BullMQ adapters are thin layers over those clients. The Rust trait (§5) has seven operations: enqueue, lease, complete, schedule, flow, submit_engine_job and watch, plus the admin calls around them.

Decorators come from Resonate's own SDKs (Python @resonate.register, TypeScript resonate.register), pointed at the Resonate server embedded in Loam (§21). Loam adds at most a thin loam helper that fills in the URL, token and tenant. There is no new decorator system.

2. Goals and non-goals

2.1 Goals

  1. Existing jobs run with a configuration change. A Celery app runs with a new broker URL and result backend. A BullMQ app runs with a new import. A PySpark job runs with a new remote(...) URL, as long as it uses the Spark Connect surface. A Flink SQL job runs after dialect edits.
  2. Durability is opt-in, per task, with the standard decorators. Replacing @app.task with @resonate.register (or registering a BullMQ processor as a Resonate function) makes a task durable: its steps are checkpointed and a crash resumes it. Nothing forces users to rewrite tasks that work.
  3. One Rust core, one protobuf surface. Every adapter, every SDK and the console use loam.jobs.v1. Semantics are defined once, in §6, and every adapter documents where the framework it imitates differs.
  4. State on Loam's own primitives. Job state and leases on TiKV (or the local store in operon dev), payloads on the object store, the queue's event log on Loam streams, schedules and flows on the embedded Resonate server. No Redis, RabbitMQ, ZooKeeper or Postgres is needed.
  5. Buy, not build. Loam writes the queue core, the adapters and the engine lifecycle glue. It hosts Sail, RisingWave, Arroyo and Flink; it does not rewrite them (user memory: prefer buy over build).

2.2 Non-goals

  • No Redis protocol. Loam does not speak RESP or run Lua. BullMQ users change an import instead (§7.2.4).
  • No AMQP. Celery reaches Loam through a kombu transport, not an AMQP 0-9-1 server. Other AMQP clients are out of scope.
  • No Spark or Flink engine in the Loam binary. Sail, RisingWave, Arroyo and Flink run as their own processes. D51 (no embedded engines) holds.
  • No second durability layer on streaming jobs. A Flink or RisingWave job's state is its own checkpoints. Loam orchestrates the lifecycle around them and never checkpoints inside them.
  • No exactly-once side effects. Delivery is at least once. Loam gives idempotency keys and fenced completion so that the record of a job's outcome is exactly once, and documents how user code makes its side effects idempotent (§6.2).
  • No new decorator or workflow DSL. Resonate's SDKs are the durable programming model.
  • No BullMQ Pro features (groups, batches, observables) and no Celery features that need a message broker Loam is not (AMQP exchanges beyond direct, fanout and topic routing).

3. Architecture

  Celery app (Python)        BullMQ app (Node/Bun)       PySpark client        Flink SQL / jar
  kombu transport loam://    @loam/bullmq                sc://<ns>.jobs…       loam jobs submit
  result backend loam://          │                           │                      │
  beat: LoamScheduler             │                           │                      │
        │  generated Connect clients (loam.jobs.v1)           │ gRPC (Spark Connect)  │
        ▼                         ▼                           ▼                      ▼
┌───────────────────────────── operon / loam binary ─────────────────────────────────────────────┐
│  jobs listener (connect-rust, Connect + gRPC + gRPC-Web)     Spark Connect proxy (auth, route) │
│        │                                                          │                           │
│        ▼                                                          │                           │
│  operon-jobs:  JobsService ── QueueCore ─── Leases/fencing ───────┼── EngineJobs              │
│                    │            │  ▲             │                │   (Sail, RisingWave,       │
│                    │            │  │ outbox      │                │    Arroyo, Flink operator) │
│                    ▼            ▼  │ relay       ▼                ▼                           │
│  operon-durable (Resonate)   JobStore            payloads      lifecycle workflows            │
│  schedules, flows (promises) (TiKV | local)      (object store) (Resonate Rust SDK, in-proc)  │
│        │                        │                                                              │
│        ▼                        ▼                                                              │
│  durable store (§21)     queue event log: Loam stream `_jobs/<queue>` (idempotent producer)   │
└────────────────────────────────────────────────────────────────────────────────────────────────┘
        │                        │                         │                          │
   SQLite | TiKV            TiKV keyspace `loam_jobs`   RustFS / S3 bucket       Sail pods, RisingWave,
                            (local: redb in dev)        ns/<ns>/jobs/…           Flink (K8s operator)

3.1 Components

ComponentWhat it doesBuilt on
operon-jobsThe Jobs trait (§5), the queue core, lease fencing, rate limits, the outbox relay, the engine-job registryoperon-tikv (cloud), redb (dev), operon-store, operon-durable
operon-jobs-protoloam.jobs.v1 messages and the JobsService trait and client, generated in build.rs like operon-live-protobuffa, connect-rust (connectrpc-build)
Jobs listenerServes JobsService over Connect, gRPC and gRPC-Web on 127.0.0.1:7720 (loopback until the auth plan, like Live's 7710 and durable's 8001)connect-rust on axum 0.8
loam-celery (Python, PyPI)kombu transport, Celery result backend, beat schedulerThe generated Python Connect client
@loam/bullmq (TypeScript, npm)A BullMQ v6 IQueueBackend over loam.jobs.v1, and BullMQ's classes re-exported bound to itThe generated TypeScript Connect client (protobuf-es)
loam helpers (Python, TypeScript)Configure Resonate's SDKs for Loam: URL, token, tenant, group. Nothing elseResonate's SDKs
Engine runnersDeploy and supervise Sail, RisingWave, Arroyo and Flink jobs; each lifecycle is a Resonate workflowKubernetes API (kube-rs) in cloud; local processes in dev

3.2 Where each piece of state lives (D214)

StateStoreWhy
Job records, ready and delayed indexes, leases, dedup keys, rate-limit buckets, chord countersTiKV keyspace loam_jobs in cloud; a redb file in operon dev only (§6.8); every other deployment, self-hosted included, uses TiKV (Q97, owner-decided)Every state change is a small multi-key transaction (dequeue = read the ready index, lock the job, write the lease). TiKV gives that with Percolator transactions; redb gives it in one process
Payloads and results larger than 16 KiB (estimate)Object store (RustFS in cloud and self-hosted, any S3 elsewhere), ns/<ns>/jobs/<queue>/<job>/{payload,result}TiKV prefers small values (§21 Q42); payloads can be megabytes
The queue's event logA Loam stream per queue, _jobs/<queue>, written by an outbox relay with an idempotent producer (D72)QueueEvents and watch, replay, audit, and analytics through a stream → Iceberg link (M4), with exactly-once events
Schedules, flows, durable-mode tasksResonate (the embedded server, §21)Cron with a promise template, parent/child promises, replay
Engine job specs and runsTiKV (a job record of kind engine) plus Resonate (the lifecycle workflow)The run is a workflow: submit, wait, savepoint, upgrade
Engine checkpoints and savepoints, Spark and Iceberg dataObject store (RustFS); Iceberg metadata in LakekeeperThe engines write them there themselves

4. Concepts

  • Namespace. The tenant, as in §18 §6 and §21 §5.1. Every queue, schedule, flow and engine job belongs to one namespace. Keys, quotas and the event-log streams are per namespace.
  • Queue. A named, namespaced set of jobs with its own settings: priorities on or off, the default visibility timeout, the retry policy, rate limits, concurrency caps, retention and the dead-letter queue.
  • Job. One unit of work: a task name, a payload, options and a state (§6.4). A job has a server-assigned JobId (a ULID) and an optional client jobId that is unique in the queue (BullMQ's jobId, Celery's task id).
  • Task spec. The task name (Celery's name, BullMQ's name), the payload (bytes plus a content type: application/json, Celery's kombu envelope, msgpack), and headers.
  • Lease. A worker's exclusive, time-bounded claim on one job, with a fencing token (§6.3). Completing, failing, extending or moving a job requires its current token.
  • Worker. A process that leases jobs from one or more queues. Workers are identified (WorkerId: host, pid, a random suffix) so that stalled-job reports and metrics can name them.
  • Schedule. A cron or interval rule that enqueues a job or starts a durable function. It is a Resonate schedule (§8.1).
  • Flow. A graph of jobs with dependencies: Celery's chain, group and chord; BullMQ's parent/child flows; a general DAG. The graph's joins are Resonate promises (§8.2).
  • Engine job. A Spark Connect session or batch job, a streaming SQL pipeline, or a Flink deployment. It is not leased by workers; an engine runner owns it (§9, §10).
  • Mode. A queue job runs in queue mode (plain lease, run, complete, like Celery and BullMQ today) or durable mode (the task body is a Resonate function; a retry replays it from its last checkpoint). §8.3 says when each is used.

5. The Rust API (D205)

5.1 The trait

The owner's sketch had seven methods. The refined trait keeps them and adds what the adapters need: a request context that carries the tenant, a fencing token instead of a bare lease id, lease extension, typed outcomes, typed errors and the admin calls. Everything a Celery transport or a BullMQ backend does is one of these calls.

/// Who is calling, for which namespace, until when. Built by the transport
/// layer from the credential (the unified auth plan, D111) and passed to
/// every call; the trait never trusts a namespace named in a request body.
pub struct Ctx {
    pub namespace: NamespaceId,
    pub principal: Principal,          // API key, agent token or `system`
    pub deadline: Option<Instant>,
    pub request_id: RequestId,         // for logs and traces
}

pub struct QueueId(pub String);        // unique in the namespace; `[a-z0-9._:-]{1,128}`
pub struct JobId(pub Ulid);            // server-assigned
pub struct WorkerId(pub String);       // host/pid/random, chosen by the worker

/// The fencing token of a lease: the job and the lease epoch. Every write
/// that depends on holding the job (extend, complete, progress, move)
/// carries it, and the store refuses it once the epoch has moved (§6.3).
pub struct LeaseToken { pub queue: QueueId, pub job: JobId, pub epoch: u64 }

#[async_trait]
pub trait Jobs: Send + Sync {
    // Producing
    async fn enqueue(&self, cx: &Ctx, queue: &QueueId, task: TaskSpec, opts: EnqueueOpts)
        -> Result<Enqueued, JobsError>;
    async fn enqueue_bulk(&self, cx: &Ctx, queue: &QueueId, jobs: Vec<(TaskSpec, EnqueueOpts)>)
        -> Result<Vec<Result<Enqueued, JobsError>>, JobsError>;

    // Consuming
    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 report(&self, cx: &Ctx, token: &LeaseToken, report: Report)
        -> Result<(), JobsError>;              // progress, log line, data update

    // Orchestration (Resonate)
    async fn schedule(&self, cx: &Ctx, spec: ScheduleSpec) -> Result<ScheduleId, JobsError>;
    async fn unschedule(&self, cx: &Ctx, id: &ScheduleId) -> Result<(), JobsError>;
    async fn flow(&self, cx: &Ctx, graph: FlowSpec) -> Result<FlowHandle, JobsError>;

    // Engines (Sail, RisingWave, Arroyo, Flink)
    async fn submit_engine_job(&self, cx: &Ctx, job: EngineJob, opts: SubmitOpts)
        -> Result<RunId, JobsError>;
    async fn control_engine_job(&self, cx: &Ctx, run: &RunId, action: EngineAction)
        -> Result<ActionId, JobsError>;         // savepoint, suspend, resume, upgrade, rollback, cancel

    // Observing: a resumable stream of events for one job, queue, flow, run or schedule
    fn watch(&self, cx: &Ctx, what: Selector, from: Cursor)
        -> BoxStream<'static, Result<Event, JobsError>>;

    // Admin (§5.4)
    async fn queue_admin(&self, cx: &Ctx, queue: &QueueId, op: QueueOp)
        -> Result<QueueOpResult, JobsError>;
    async fn job_admin(&self, cx: &Ctx, queue: &QueueId, job: JobRef, op: JobOp)
        -> Result<JobOpResult, JobsError>;
    async fn query(&self, cx: &Ctx, q: JobQuery) -> Result<Page<JobView>, JobsError>;

    // Result store (§7.1.4): Celery's key-value result backend, TTL'd, one TiKV txn each
    async fn result_get(&self, cx: &Ctx, key: &ResultKey) -> Result<Option<Payload>, JobsError>;
    async fn result_mget(&self, cx: &Ctx, keys: Vec<ResultKey>)
        -> Result<Vec<Option<Payload>>, JobsError>;
    async fn result_set(&self, cx: &Ctx, key: &ResultKey, value: Payload, ttl: Option<Duration>)
        -> Result<(), JobsError>;
    async fn result_delete(&self, cx: &Ctx, key: &ResultKey) -> Result<(), JobsError>;
    async fn result_incr(&self, cx: &Ctx, key: &ResultKey, member: &str, ttl: Option<Duration>)
        -> Result<Counted, JobsError>;          // atomic; a member counts once (§7.1.4)
    async fn result_expire(&self, cx: &Ctx, key: &ResultKey, ttl: Duration)
        -> Result<(), JobsError>;
}

5.2 The main types

pub struct TaskSpec {
    pub name: String,                  // Celery task name, BullMQ job name
    pub payload: Payload,              // bytes + content type (json, celery-kombu, msgpack)
    pub headers: BTreeMap<String, String>,
}

pub struct EnqueueOpts {
    pub job_key: Option<String>,       // unique in the queue (BullMQ jobId, Celery task id)
    pub idempotency: Option<Idempotency>, // key + ttl (default 24 h), §6.2
    pub dedup: Option<Dedup>,          // BullMQ deduplication: simple | throttle(ttl) | debounce(ttl, replace, extend)
    pub not_before: Option<Delay>,     // delay or absolute time (Celery eta/countdown, BullMQ delay)
    pub priority: Priority,            // §6.5; None = the FIFO band
    pub lifo: bool,
    pub attempts: u32,                 // default: the queue's
    pub backoff: Option<Backoff>,      // fixed | exponential { base, max, jitter } | custom(name)
    pub visibility: Option<Duration>,  // lease length for this job, default the queue's
    pub timeout: Option<Duration>,     // hard limit: the lease cannot be extended past it
    pub expires: Option<SystemTime>,   // Celery `expires`: not started after this → discarded
    pub parent: Option<ParentRef>,     // flow membership (§8.2)
    pub retention: Option<Retention>,  // removeOnComplete / removeOnFail, result TTL
    pub mode: Mode,                    // Queue | Durable { function } (§8.3)
}

pub struct Enqueued { pub job: JobId, pub created: bool, pub state: JobState }

pub struct LeaseRequest {
    pub queues: Vec<QueueId>,          // in the worker's order of preference (Celery -Q a,b)
    pub worker: WorkerId,
    pub max: u32,                      // prefetch / free concurrency slots
    pub wait: Duration,                // long-poll, at most 30 s; 0 = return at once
    pub names: Option<Vec<String>>,    // only these task names (Celery routing by name)
}

pub struct Lease {
    pub token: LeaseToken,
    pub task: TaskSpec,                // payload inline, or a presigned object-store URL (§6.7)
    pub attempt: u32,                  // 1-based; stalled redeliveries counted separately
    pub stalled: u32,
    pub deadline: SystemTime,          // by the store's clock (§6.3)
    pub parent: Option<ParentRef>,
    pub children: Option<ChildrenSummary>, // for parents that resumed after their children
}

pub enum Outcome {
    Ok { result: Option<Payload> },
    Retry { error: JobError, after: Option<Duration> }, // after = None: the job's backoff
    Fail { error: JobError, dead_letter: bool },        // no more attempts (UnrecoverableError)
    Release { after: Option<Duration> },                // give back without spending an attempt (nack)
    Delay { until: SystemTime },                        // BullMQ moveToDelayed
    WaitChildren,                                       // BullMQ moveToWaitingChildren
    RateLimited { for_: Duration },                     // BullMQ Worker.RateLimitError: queue-wide pause
}

pub enum JobState {
    Waiting, Prioritized, Delayed, Active, WaitingChildren,
    Completed, Failed, DeadLettered, Canceled,
}

ScheduleSpec, FlowSpec and EngineJob are in §8.1, §8.2 and §9–§10.

5.3 Errors

One error enum, mapped one to one onto Connect codes, so the Python and TypeScript clients see the same classes:

JobsErrorConnect codeRetry?When
InvalidArgument(msg)invalid_argumentNoBad names, sizes, cron expressions, options that contradict each other
NotFound(what)not_foundNoUnknown queue, job, schedule, flow or run
AlreadyExists { existing }already_existsNoA job_key that exists with different parameters. With the same parameters enqueue succeeds with created: false
IdempotencyConflict { key }failed_preconditionNoAn idempotency key reused with a different request (hash mismatch, like D146 and Live's args_hash)
Fenced { current_epoch }failed_preconditionNoThe lease epoch moved: the job was reclaimed and leased again. The worker must drop its result
InvalidState { from, to }failed_preconditionNoFor example promote on a job that is not delayed
QueuePausedunavailableYes, after resumelease on a paused queue returns no jobs rather than this error; admin writes may return it
RateLimited { retry_after }resource_exhaustedYesA namespace quota (D65) or the jobs listener's own request limit. A queue's job rate limit never errors: lease returns fewer jobs
PayloadTooLarge { limit }resource_exhaustedNoOver the namespace's payload limit (default 64 MiB)
Unavailable(msg)unavailableYesTiKV region errors after the runner's retries, the durable server down, an engine runner unreachable
DeadlineExceededdeadline_exceededCaller decidesThe Ctx deadline passed
Unauthenticated, PermissionDeniedthe sameNoFrom the auth layer (D111)
Internal(msg)internalNoA bug or corrupt record; logged with the request id

A write whose outcome is unknown (a timeout after TiKV's commit, §20 §5) is retried by the store with its commit token (operon-tikv runner), so clients see either success or a definite error. An Unavailable from enqueue is safe to retry only with a job_key or idempotency key; the generated clients add an idempotency key to every enqueue by default (a UUIDv7 made once per call, reused across that call's retries).

5.4 Admin operations

QueueOp: Create(QueueConfig), Update(QueueConfig), Pause, Resume, Drain { delayed }, Clean { state, grace, limit }, Obliterate { force }, Counts, RetryAll { state }, PromoteAll, SetRateLimit, SetConcurrency, Workers, Metrics, RedriveDeadLetters { limit }. JobOp: Get, Remove, Retry, Promote, ChangePriority, ChangeDelay, UpdateData, Logs, Cancel, RemoveDedupKey. JobQuery: by queue, state, name, time range and tag, paginated, newest first; BullMQ's getJobs(types, start, end, asc) and Celery's inspect map onto it.

These cover BullMQ's IQueueBackend getters and admin methods (§7.2) and Celery's purge, inspect and control calls that do not need a broadcast channel (§7.1.5).

5.5 The protobuf surface (D206)

  • Package loam.jobs.v1, files under proto/loam/jobs/v1/ at the workspace root (like proto/loam/live/v1/): jobs.proto (the service), types.proto, engines.proto, events.proto.
  • Service JobsService: unary Enqueue, EnqueueBulk, Lease (long-poll), Extend, Complete, Report, Schedule, Unschedule, Flow, SubmitEngineJob, ControlEngineJob, QueueAdmin, JobAdmin, Query, and the result store ResultGet, ResultMGet, ResultSet, ResultDelete, ResultIncr, ResultExpire (§7.1.4); server-streaming Watch. Lease is unary with a wait, not a stream: a worker's prefetch is explicit, and connection loss cannot strand jobs in a push buffer. A server-streamed LeaseStream is a later optimization (Q98).
  • Server operon-jobs-proto: buffa messages and the connect-rust JobsService trait, generated in build.rs with connectrpc-build and the system protoc (D128, as operon-live-proto does). It speaks Connect, gRPC and gRPC-Web on one listener.
  • Clients: buf generate with protobuf-es and @connectrpc/connect for TypeScript, and connect-python (connectrpc on PyPI) for Python (§12.1). The Celery and BullMQ adapters wrap these clients; they do not hand-write HTTP.
  • Versioning: additive changes only within v1; buf breaking runs in CI against the last release.

5.6 Method → backing primitives

MethodTiKV (loam_jobs)Loam streamsResonateObject store
enqueueOne optimistic txn: idempotency record, job_key index, dedup key, job record, ready or delayed index entry, outbox rowEvent added/waiting/delayed via the outbox relay (§6.6)Durable mode: nothing at enqueue; the root promise is created when the job first runs (§8.3). With a parent: the child's completion promise (§8.2)Payload PUT before the txn when over 16 KiB (§6.7)
leaseOne pessimistic txn per shard visited: scan the ready index head, check the rate bucket and concurrency counter, lock and move up to max jobs to active with epoch + 1 and a deadlineactive events—Presigned GET URL for large payloads
extendTxn: check the epoch, move the deadline (never past timeout)———
completeTxn: check the epoch; write the result or error; move to completed, failed, delayed (retry), waiting-children or the DLQ; decrement concurrency; parent bookkeeping; outbox rowcompleted/failed/retries-exhausted/delayed eventsSettles the flow node's promise (through the outbox, §8.2)Result PUT when over 16 KiB
reportTxn: check the epoch; progress, log append (capped), data updateprogress event——
scheduleSchedule metadata (for listing and BullMQ's getters)—schedule.create with a promise template (§8.1)—
flowJob records for every node, parents in waiting-childrenadded eventsRoot promise = flow id; one child promise per node; chord joins (§8.2)—
submit_engine_jobEngine run recordRun eventsThe lifecycle workflow (§9, §10)Artifacts (jars, SQL, Python files), savepoints and checkpoints written by the engines
watchSnapshot of current state for the selectorSubscribe from the cursor (a stream offset)Promise state for flows and durable jobs—
stalled sweep (internal)Per-shard scan of expired leases → waiting (stalled + 1) or failedstalled events——
delayed promoter (internal)Per-shard scan of due delayed entries → readywaiting events——

6. Semantics (D207)

6.1 Delivery

  • At least once. A job is delivered to one worker at a time. It is delivered again if the lease expires before complete (a crashed or stalled worker), or if the worker returns Retry or Release. Duplicated deliveries are therefore possible; lost jobs are not, once enqueue has returned.
  • The outcome is recorded exactly once. complete is fenced (§6.3): of two workers that both ran a job (the first stalled, the second leased it after the lease expired), only the holder of the current epoch can record the outcome. The other gets Fenced. So 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. Loam gives two tools: the lease token's epoch as a fencing token for the task's own writes (a Loam write with a Fence, an external API with an idempotency key made from (job id, step)), and durable mode, where each step is a Resonate checkpoint and its promise id is the idempotency key of its effect (§8.3, §21 §6.7). The docs say plainly that an existing task with side effects and no idempotency key can repeat those side effects on redelivery, exactly as it can on Redis or RabbitMQ today.
  • No exactly-once claim. Loam does not say "exactly once" for job execution anywhere, in the docs or the API.

6.2 Idempotency and deduplication

Three separate mechanisms, because the frameworks have three:

MechanismScopeBehaviourFrom
job_keyUnique in the queue while the job exists (until retention removes it)A second enqueue with the same key returns the existing job (created: false) if the parameters match, else AlreadyExistsBullMQ jobId; Celery task id (the kombu message's id header)
Idempotency keyPer namespace and queue, for a TTL (default 24 h), independent of the job's retentionA retried enqueue returns the first call's Enqueued. A different request under the same key is IdempotencyConflictD146's operations API; Live's idempotency records (D118)
DedupPer queue and dedup idBullMQ's modes: simple (while the job is not finished), throttle (for ttl), debounce (replace the pending job's data, optionally extending the ttl), keep-last-if-activeBullMQ deduplication

The record layout follows Live's IdempotencyRecord: a hash of the key, the result, an expiry by the TSO clock and a hash of the canonical request.

6.3 Leases and fencing

  • Epoch per job. A job record carries epoch. lease increments it and writes (worker, epoch, deadline). Every later write for that job (extend, complete, report, moves) is one transaction that reads the job with get_for_update and refuses the write unless its epoch equals the token's. This is the rule operon-meta-tikv's check_fence already implements for metastore leases (crates/operon-meta-tikv/src/leases.rs).
  • Expiry alone does not break a fence. As in operon_common::meta::Fence: until the stalled sweep reclaims the job (epoch + 1 on the next lease), a late complete from the original worker still succeeds. That avoids failing work that finished a moment after its deadline.
  • Clock. Deadlines are judged by the store's clock (the TSO physical time in TiKV, as in leases.rs; the process clock with the local store), never by a worker's clock. The Celery transport renews at a third of the lease, to leave room for a slow TiKV write; BullMQ's own worker keeps its lockRenewTime setting, because @loam/bullmq runs BullMQ's unmodified Worker.
  • Stalled sweep. Each queue shard has one owner at a time (a metastore lease jobs/sweep/<ns>/<queue>/<shard>, placed by the router's rendezvous hashing, §18 §5). The owner scans the shard's active index by deadline. An expired job goes back to waiting with stalled + 1 and a stalled event; past the queue's max_stalled (BullMQ's maxStalledCount, default 1) it fails with "job stalled more than allowable limit", as BullMQ does. Celery's acks_late redelivery is the same path with max_stalled unlimited.
  • Visibility and ETA are separate. A delayed job is in the delayed index, not leased. It is never "invisible" while it waits, so a long countdown or eta does not cause the duplicate deliveries it causes on Redis and SQS when it exceeds visibility_timeout (§7.1).
  • Hard timeout. timeout caps the total lease; extend past it is refused, and the sweep treats the job as timed out (a retryable failure with error.code = "timeout").

6.4 States

             enqueue                    lease                    complete(Ok)
  (none) ──► waiting | prioritized ───► active ─────────────────► completed ──► (removed by retention)
      │           ▲   ▲                  │  │  │ complete(Retry), attempts left
      │ delay     │   │ due              │  │  └───────────────► delayed ──┐
      └──► delayed┘   └──────────────────┼──┘ complete(Release)            │ due
                                         │     lease expired (stalled)     ▼
                                         │  ──► waiting (stalled+1) ◄──────┘
                                         │ complete(Fail) or attempts exhausted
                                         ├──► failed ──► (dead_letter?) ──► DLQ queue: waiting
                                         └ complete(WaitChildren) ──► waiting-children ──(children done)──► waiting

waiting and prioritized are one index with two bands (§6.5), reported separately because BullMQ reports them separately. Paused queues keep their jobs in place and lease returns nothing (BullMQ v6 dropped the public paused state).

6.5 Priorities, ordering and LIFO

  • Model. A job with no priority is in the FIFO band, served before any prioritized job. A prioritized job has a priority in 1…2,097,151 and a lower number is served first; ties are FIFO. This is BullMQ's model exactly (PRIORITY_LIMIT = 2^21 − 1, score priority × 2^32 + counter) (source). Celery's priorities (0–9 on the Redis transport, where 0 is highest; 0–255 on RabbitMQ, where higher is higher) are mapped by the transport (§7.1).
  • Index key. ready/<shard>/<band><priority:u32 BE><seq:u64 BE>; LIFO jobs use u64::MAX − seq. A lease scans from the start of the key range, so the order is one range scan.
  • Sharding. A queue has S shards (default 1; up to 64) for throughput. A job's shard is chosen at enqueue (hash of the job id, or round robin). With S > 1, order holds per shard, not across shards: priorities are approximate across shards (a worker visits shards in a rotating order and takes the best head it sees in the first two it visits (estimate)). A queue that needs strict global priority keeps S = 1.
  • No ordering guarantee across workers. Like BullMQ and Celery, jobs are started in order but may finish in any order. Strict per-key ordering (BullMQ Pro's groups, SQS FIFO) is not offered (Q99).

6.6 The event log and watch

  • Every state change writes an outbox row in the same transaction: outbox/<shard>/<seq>.
  • A relay per queue shard (the same owner as the stalled sweep) reads the outbox in order and appends the events to the stream _jobs/<queue> in the namespace, partition = shard, with an idempotent producer (D72): producer id jobs-relay/<queue>/<shard>, epoch = the relay's lease epoch, sequence = the outbox sequence. A relay that crashed after appending and before deleting its outbox rows re-appends them, and the sequencer returns the original offsets without appending (§02 §7). So each event appears in the log exactly once, in order per shard.
  • watch reads the current state from TiKV and, in the same snapshot (one TiKV read at one start_ts), each shard's highest outbox sequence. Because every state change writes its outbox row in the same transaction, that sequence is exactly the set of changes the snapshot already reflects. watch then subscribes to each partition and skips records whose outbox sequence (every record carries it) is at or below the snapshot's; the returned cursor is the stream offset per partition of the first record it delivers. No change can fall between the snapshot and the subscription: a change committed after the snapshot has a higher outbox sequence and is delivered, even if the relay appends it later. A client that reconnects passes its cursor and misses nothing within the stream's retention (default 7 days for _jobs/*, estimate). QueueEvents (BullMQ) and celery events-style monitors read this.
  • The event stream can be linked to an Iceberg table (M4) for job analytics: durations, failure rates, per-task cost.
  • Until D72's idempotent producers land (M2), the relay writes events with at-least-once delivery and includes the outbox sequence in each record so consumers can drop duplicates. In operon dev without a stream engine, watch tails the outbox directly (§6.8).

6.7 Payloads and results

  • Inline in the job record up to 16 KiB (estimate); larger payloads are written to ns/<ns>/jobs/<queue>/<job>/payload before the enqueue transaction, which references the object. The object carries a Freshness (§03), so GC deletes objects whose transaction never committed.
  • Results follow the same rule. Celery's result backend reads results through query/Get; large results come back as a presigned URL that the client fetches.
  • Payloads are opaque to Loam. The docs repeat §21 §8's rule: pass references, not personal data, or encrypt payloads client-side; erasure (D68) cannot look inside a payload.

6.8 Stores: TiKV everywhere, redb in dev

operon-jobs defines a JobStore trait (the transactions above) with two implementations:

StoreWhereNotes
TikvJobStoreLoam cloud, self-hosted clusters and operon standalone: every deployment except operon dev; keyspace loam_jobs on the Live cluster, tenant = key prefix t/<ns>/ (large tenants get their own keyspace, like §21 §5.1)Uses operon-tikv's runner (classification, retries, commit tokens, fault plan). Optimistic for enqueue, pessimistic for lease (the ready-index head is contended)
LocalJobStoreoperon dev onlyOne redb file (<data>/jobs.redb), one writer; same key layout. operon standalone and operon cluster refuse it and ask for --jobs-store tikv://…

Owner decision (Q97, 2026-09-29): jobs require TiKV. Self-hosted clusters run TiKV for jobs; there is no Postgres or DynamoDB jobs backend, and redb stays a development store for operon dev only.

6.9 Rate limits and concurrency

  • Queue rate limit (max jobs per duration, BullMQ's limiter): a token bucket in one TiKV key per queue, debited inside the lease transaction by the number of jobs taken. lease returns fewer jobs, or none with a retry_after hint, when the bucket is empty. Outcome::RateLimited (BullMQ's Worker.RateLimitError, queue.rateLimit(ms)) empties the bucket for the given time.
  • Global concurrency (BullMQ's setGlobalConcurrency): a counter of active jobs per queue, checked and incremented in lease, decremented in complete and by the stalled sweep.
  • Worker concurrency stays in the worker (Celery's --concurrency, BullMQ's concurrency): the worker asks lease for at most its free slots.
  • Celery's rate_limit is enforced by the Celery worker itself (per worker, not global) (source, §7.1), so the transport does not need to implement it. The queue rate limit is offered as a global alternative.
  • Hot keys. The bucket and the counter are one key each, written by every lease. With batched leases (a worker takes up to its free slots in one call) this is one write per batch, not per job. A single queue is budgeted at about 2,000 leases/s with S = 1 and scales with shards (estimate; J1 measures it). Rate-limited queues above that need a sharded bucket (one per queue shard, each with max / S), which the queue config can turn on.
  • Namespace quotas (D65) cap enqueue rate, payload bytes, stored jobs and concurrent engine runs; they return RateLimited or PayloadTooLarge.

6.10 Dead-letter queues and retention

  • A queue may name a dead-letter queue (another queue in the namespace). A job that fails with dead_letter: true, or exhausts its attempts on a queue with dead_letter_on_exhausted, is re-enqueued there with its error history and original queue in headers. RedriveDeadLetters moves them back.
  • Without a DLQ, failed jobs stay in failed (BullMQ's behaviour) until retention removes them.
  • Retention by count or age per state (BullMQ's removeOnComplete/Fail, Celery's result_expires, default 1 day for results). A per-shard janitor deletes expired records and their payload objects. Idempotency records expire on their own TTL, independent of the job.

6.11 Observability

  • Metrics on the admin listener (/metrics): loam_jobs_enqueued_total{ns,queue}, loam_jobs_leased_total, loam_jobs_completed_total{outcome}, loam_jobs_stalled_total, loam_jobs_fenced_total, loam_jobs_waiting{band}, loam_jobs_delayed, loam_jobs_active, loam_jobs_oldest_waiting_seconds, loam_jobs_lease_latency_seconds, loam_jobs_run_seconds{task}, loam_jobs_dlq_depth, loam_engine_runs{engine,state}. BullMQ's getMetrics and exportPrometheusMetrics read the same counters.
  • Traces. Each job gets a trace context at enqueue (W3C traceparent in the headers, as Celery and BullMQ's telemetry propagate it). The lease and the completion are spans; durable steps add Resonate's spans (§21 §6.6). OTLP export follows Loam's M2 observability work; traces ingest into Loam is Q43.
  • Logs. Job log lines (report(Log), BullMQ's job.log) are capped per job (default 1,000 lines) and stored with the job.
  • The console lists queues, counts, jobs, DLQs, schedules, flows and engine runs from query and watch (after §19's console work).

7. Adapters

7.1 Celery: loam-celery (D208)

Celery 5.6.3 and kombu 5.6.2 are the current releases, both BSD-3-Clause (source). Celery splits its needs between two plug points: the broker (a kombu transport) carries task messages, and the result backend stores results, group results and chord counters. Beat is a third plug point. loam-celery implements all three over loam.jobs.v1.

# celeryconfig.py: the only change for a queue-mode app
import loam_celery                      # registers the kombu alias "loam" (§7.1.1)
# Same host only until the unified auth plan (D111): the jobs listener binds
# 127.0.0.1:7720 (§3.1, §11). After D111 the URL is the authenticated endpoint,
# jobs.<cloud-domain>:443 with TLS, and nothing else changes.
broker_url = "loam://127.0.0.1:7720/my-namespace"
result_backend = "loam://127.0.0.1:7720/my-namespace"
beat_scheduler = "loam"                 # optional: schedules live in Loam (§7.1.6)
broker_transport_options = {"visibility_timeout": 3600, "token_env": "LOAM_TOKEN"}
result_backend_transport_options = {"token_env": "LOAM_TOKEN"}

Credentials. The transport and the result backend are separate Celery objects, so each reads its own options: the transport from broker_transport_options, LoamBackend from Celery's result_backend_transport_options. Both take token_env (the environment variable holding the Loam token) and default to LOAM_TOKEN when it is absent, so the two lines above are optional. LoamBackend builds one loam.jobs.v1 client per process with that token; AsyncResult, GroupResult.save/restore and the chord counter all go through the backend and therefore that client. A token is never taken from the URL.

7.1.1 Registration

Plug pointMechanism in Celery/kombuWhat loam-celery does
Transportkombu resolves the URL scheme through the dict TRANSPORT_ALIASES (kombu/transport/__init__.py:21-49); there is no entry-point group for transports (entry points are read only for matchers and serializers). A URL scheme may not be a module:Class path (Connection._check_url_transport, connection.py:63-86). An explicit transport= (Celery's broker_transport, app/base.py:1213) wins over the URL schemeTwo supported ways: import loam_celery adds TRANSPORT_ALIASES["loam"] = "loam_celery.transport:Transport" at import time, or the app sets broker_transport = "loam_celery.transport:Transport". The docs recommend the second, because it does not depend on import order. An upstream kombu PR for a kombu.transports entry-point group is proposed, not required (§12.3)
Result backendby_name merges BACKEND_ALIASES, the loader's overrides and the celery.result_backends entry-point group (app/backends.py:15-47)Entry point loam = loam_celery.backend:LoamBackend, so result_backend = "loam://…" works with no import
Beat schedulerbeat_scheduler is resolved through symbol_by_name and the celery.beat_schedulers entry-point group (beat.py:686-690)Entry point loam = loam_celery.beat:LoamScheduler

7.1.2 The transport

kombu's virtual transport (kombu/transport/virtual/base.py) keeps unacked messages in an in-process OrderedDict (QoS._delivered) and restores them only at shutdown or channel close. After a crash, nothing is recovered unless the backend has its own visibility timeout (restore_unacked_once, lines 192-195, 743-751, 803) (source). Redis emulates one with an unacked hash and a restore sweep (redis.py:408-490); SQS and the new pgmq transport have a server-side one. Loam's lease is that server-side visibility timeout, so the transport follows the pgmq/SQS model, not the in-process one:

kombu callLoam callNotes
_put(queue, message)enqueueThe whole kombu message (body, headers, properties) is the payload, content type application/x-kombu+json. Loam reads four headers: task → TaskSpec.name; id + retries → job_key (§7.1.4); eta → not_before; expires → expires; and the priority property (§7.1.3)
_get(queue) / _get_many(queues)lease (with wait up to the polling interval, and max = the free prefetch slots)Long-polling replaces kombu's 1 s polling_interval loop. The delivery tag maps to the LeaseToken
basic_ack(tag)complete(Ok)Celery acks early by default, when the pool accepts the task (worker/request.py:391-392), and after the task with acks_late
basic_reject(tag, requeue=True)complete(Release)Used by task_reject_on_worker_lost and timeouts under acks_late (request.py:705-724)
basic_reject(tag, requeue=False)complete(Fail { dead_letter })To the queue's DLQ if it has one
basic_recover(requeue=True)complete(Release) for every unacked tag
unacked while runningextend every visibility_timeout / 3 from a transport threadLike the gcpubsub transport's modify_ack_deadline loop. A long acks_late task on a live worker is never redelivered by timeout, which removes Redis's "task longer than visibility_timeout runs twice" problem. A dead worker's leases expire and the jobs are redelivered (Loam's stalled path, with max_stalled unlimited for Celery queues)
_size, _purge, _delete, _new_queue, _has_queuequeue_admin(Counts / Drain / Obliterate / Create), queryQueues are created on first declare, as with Redis
get_table, queue_bind, exchange declareBindings stored in Loam (a small table per namespace)kombu's default get_table reads process-local state (base.py:702), so topic routing across processes needs shared bindings, as Redis keeps them in _kombu.binding.*
_put_fanout, supports_fanout = TrueA broadcast stream per fanout exchange (_celery/<exchange>); each consumer subscribes from latestNeeded for remote control (§7.1.5)

Transport.implements declares direct, topic and fanout exchange types and asynchronous = False in J1. driver_type = "loam".

7.1.3 Priorities, ETA and retries

  • Priorities. kombu's virtual transports clamp priority to 0–9 (base.py:472-473, 854-872); Redis serves lower numbers first by polling priority lists in ascending order (redis.py:1266-1272, 1451-1480), and a message with no priority counts as 0. loam-celery maps 0 or unset → Loam's FIFO band, and 1–9 → Loam priority 1–9, which reproduces the Redis order exactly. A transport option priority_order = "amqp" reverses it for apps written for RabbitMQ, where higher is higher. task_queue_max_priority only sets an AMQP queue argument (app/amqp.py:97-100) and is ignored.
  • ETA and countdown. Celery workers hold an ETA task unacked in memory until it is due (worker/strategy.py:180-208), which is why an ETA longer than the visibility timeout makes Redis and SQS redeliver it "again, and again in a loop" (Celery's Redis docs, redis.rst:303-334). loam-celery puts the job in Loam's delayed index instead (§6.3), so the worker receives it only when it is due and never holds it. This is a behaviour change for the better; the ETA semantics (not before) are unchanged.
  • Retries. Task.retry publishes a new message with the same task id (app/task.py:767-873). The transport therefore uses (id, retries) as the job_key, never the id alone, or a retry published while the original is still active (acks_late) would collide with it.
  • Rate limits and time limits are enforced by the Celery worker (rate_limit with a kombu TokenBucket, worker/consumer/consumer.py; time limits by the pool, request.py:362-373) and need nothing from the broker. Loam's queue rate limit (§6.9) is an optional global limit on top.

7.1.4 The result backend

LoamBackend subclasses Celery's BaseKeyValueStoreBackend (backends/base.py:1095) and implements its six primitives, get, mget, set, delete, incr and expire, over a small result store in loam.jobs.v1 (ResultGet, ResultMGet, ResultSet { ttl }, ResultDelete, ResultIncr, ResultExpire), kept in the job store under results/<ns>/… with a TTL (result_expires, default 1 day). Because it sets implements_incr = True, Celery uses the native chord counter (_apply_chord_incr, on_chord_part_return with incr, base.py:1100-1365): each header task's completion calls ResultIncr { key: chord counter, member: task id }. The counter is stored as the set of members that counted, and ResultIncr adds the member and returns the new size in one TiKV transaction, so each (group_id, task_id) counts at most once: a redelivered or re-run header task (a crash after incr and before the ack, or a stalled lease) adds nothing, and only the call that moves the size to the chord size gets reached: true and sends the body, so the body is sent once. A count is never skipped by a crash either: on_chord_part_return runs before the task is acked under acks_late (the setting the docs recommend for chords), so a worker that dies before incr leaves the job to be redelivered, and the re-run counts. With early acks a crash between the ack and incr loses the count, as it does on Redis; the docs say so. Without incr, Celery falls back to the celery.chord_unlock task, which polls and retries forever (app/builtins.py:37-79); Loam avoids that. AsyncResult.get polls through wait_for_pending in J1; J1.x adds push waiting over watch (as the Redis backend does with AsyncBackendMixin). GroupResult.save/restore are plain set/get.

7.1.5 What depends on the broker type, and what does not work

Some Celery features are switched on by a hard-coded list of broker driver_types, not by transport capabilities (source):

FeatureCondition in CeleryOn Loam
Remote control (inspect, control, revoke broadcast), through the pidbox fanout mailbox (app/control.py:439-441)conninfo.supports_exchange_type("fanout") (worker/consumer/control.py:31-33)Works: the transport declares fanout and implements _put_fanout
Events (celery events, Flower) on the celeryev topic exchange (events/event.py:15)Switched to fanout only for driver_type redis or gcpubsub (:58-60); topic needs shared bindingsWorks with shared bindings (§7.1.2). Loam's own job events (§6.6) are the better monitor; Flower compatibility is best-effort
Mingle (sync revoked tasks at worker start)driver_type in {amqp, redis, gcpubsub} (worker/consumer/mingle.py:25,34)Off. Revokes still reach running workers through the broadcast; a worker that starts later does not learn earlier revokes (use --statedb, or cancel waiting jobs through Loam's JobOp::Cancel, which removes them from the queue for good)
Gossipdriver_type in {amqp, redis} (gossip.py:34,79)Off
worker_disable_prefetchRedis only (consumer/tasks.py:60-67)Not needed: lease takes exactly the free slots

An upstream Celery issue proposing capability checks instead of driver_type lists is worth filing later (§12.3); nothing here depends on it.

7.1.6 Beat

Beat has no leader election: Celery's docs say to run exactly one scheduler, "otherwise you'd end up with duplicate tasks" (docs/userguide/periodic-tasks.rst:19-20). LoamScheduler removes that constraint. On start it upserts every beat_schedule entry as a Loam schedule (Jobs::schedule, id = the entry name), whose target enqueues the entry's task message into its queue. Loam's server fires the schedule, not the beat process, so running beat twice, or not at all after the first sync, cannot duplicate or miss a tick. crontab entries map to cron; timedelta entries map to every (§8.1). solar entries are refused with a message to keep the standard scheduler for them (Q101). loam-celery sync-schedules app does the same upsert from CI without a beat process.

7.1.7 Canvas

Queue mode runs Celery's canvas unchanged, because the canvas lives in the messages and the result backend: a chain carries its remaining steps in the message body (options['chain'], protocol 2, canvas.py:906-960, app/amqp.py:400-405) and the worker publishes the next step after success (app/trace.py:426-470); a group is several publishes plus a saved GroupResult; a chord uses the native counter (§7.1.4). Celery itself notes that the chain hand-off is not atomic (app/trace.py:435-437): a worker that dies between a task's success and publishing the next step leaves the chain stuck or, with acks_late, runs the step twice. Durable mode (§8.3) is the fix for chains where that matters.

7.1.8 Durable mode for Celery

The owner's "optionally swap @app.task for a Resonate decorator" is done with Resonate's Python SDK (0.8.1, the version that server 0.10.1 accepts, §21 §4) and no Loam decorator:

import asyncio
from resonate.retry import Exponential
import loam                                  # the thin helper: URL, token, tenant, group from LOAM_* env

resonate = loam.resonate(group="billing")    # = Resonate(url=…, token=…, group="billing")

@resonate.register(name="charge", version=1)
async def charge(ctx, order_id: str):
    hold = await ctx.options(retry_policy=Exponential(delay=1, max_retries=5)).run(reserve, order_id)
    # ctx.run checkpoints the result, but a crash after capture and before the
    # checkpoint runs the step again: side effects are at least once. Each
    # external call therefore gets an idempotency key made from the promise id
    # and the step, which the payment and mail APIs deduplicate on.
    receipt = await ctx.run(capture, hold, idempotency_key=f"{ctx.id}:capture")
    await ctx.run(send_receipt, receipt, idempotency_key=f"{ctx.id}:send_receipt")
    return receipt

@app.task(bind=True, acks_late=True)         # stays a Celery task: routing, priority, rate limits
def charge_task(self, order_id):
    handle = resonate.run(f"celery:{self.request.id}", charge, order_id)  # id = the job: a redelivery resumes
    return asyncio.run(handle.result())

The Celery message still carries admission (queue, priority, rate limit, ETA). The Resonate promise id is derived from the Celery task id, so a redelivered task resumes the same durable execution instead of starting it again: run with an existing id returns a handle to the same promise (py/resonate.py:574-575), finished steps return their memoized values, and only one worker executes it at a time (Resonate's task version fence). A chain rewritten as sequential ctx.run calls, a group as parallel calls awaited together, and a chord as that plus a final step, become atomic in the sense Celery's canvas is not: a crash anywhere resumes from the last finished step. Beat entries that target durable functions become Resonate schedules directly (resonate.schedule(...), py/resonate.py:793-842).

7.2 BullMQ: @loam/bullmq (D209)

Finding: BullMQ v6 already has a pluggable backend, so Loam implements it instead of imitating the API. BullMQ 6.3.9 (MIT, 2026-09-28) defines IQueueBackend, "a database-agnostic contract describing every high-level operation that the Queue, Worker and Job classes need" (src/interfaces/queue-backend.ts:26-61), with about 81 methods, and ships a Redis backend and a PostgreSQL backend (src/postgres/, 81 SQL command files, LISTEN/NOTIFY for wake-ups) that claims full parity: flows, schedulers, rate limiting, priorities, delays, deduplication, metrics and events (docs/gitbook/guide/postgresql.md:349-363) (source). A backend is injected as the last constructor argument (queue-base.ts:38-82, flow-producer.ts:121-151) or process-wide with setDefaultBackendFactory (src/utils/create-backend.ts:120); custom backends are documented (docs/gitbook/guide/connections.md:205-267).

So @loam/bullmq is a BackendFactory for the unmodified bullmq package (a peer dependency bounded to the BullMQ minors whose backend-neutral suite passed against Loam: >=6.3.9 <6.4 at J2; each later minor is added to the range only after the suite passes on it, J-R1), plus a module that re-exports BullMQ's classes already bound to it with withBackend (src/utils/with-backend.ts):

// Option A: change the import (the owner's "drop-in")
import { Queue, Worker, FlowProducer, QueueEvents } from "@loam/bullmq";
// Option B: keep `from "bullmq"` and set the backend once at startup
import { setDefaultBackendFactory } from "bullmq";
import { loamBackend } from "@loam/bullmq";
setDefaultBackendFactory(loamBackend({ url: process.env.LOAM_JOBS_URL, token: process.env.LOAM_TOKEN }));

The Queue, Worker, Job, FlowProducer and QueueEvents code users run is BullMQ's own, so the API is not re-implemented and cannot drift. Loam owns only the backend.

7.2.1 The backend mapping

IQueueBackend groupLoam call
addJob, addJobsenqueue, enqueue_bulk (jobId → job_key; deduplication → dedup; delay, priority, lifo, attempts, backoff, removeOnComplete/Fail → the same options)
addFlowflow (§8.2)
waitForJob(blockTimeout)lease with wait (the doc comment on waitForJob, queue-backend.ts:768, expects non-Redis backends to use notifications or polling)
moveToActivelease. BullMQ's worker token (a string it generates) maps to the LeaseToken in the backend's memory
extendLock, extendLocksextend
moveToFinished (completed / failed), moveToDelayed, moveToWaitingChildren, retryJobcomplete with Ok, Retry/Fail, Delay, WaitChildren
moveStalledJobsToWaitA no-op returning nothing: Loam's stalled sweep runs on the server (§6.3) and emits the same stalled events
setRateLimit, global rate limit and concurrencycomplete(RateLimited), queue_admin(SetRateLimit / SetConcurrency)
getCounts, getRanges, job getters, logs, metrics, workersquery, job_admin, queue_admin
Job scheduler operationsschedule, unschedule and their listing (§8.1)
publishEvent, readEvents(id, blockTimeout) (queue-backend.ts:752)watch on the queue's event stream; BullMQ's event id is the stream cursor. With S = 1 it is the partition's offset. With S > 1 it is an opaque vector cursor (every partition's offset, encoded in one string): the backend merges the partitions by each record's commit timestamp (the TSO of the transaction that wrote its outbox row, ties broken by shard), which keeps per-shard order, and readEvents(id) resumes every partition after its own offset, so no event is skipped or repeated. J2 verifies that QueueEvents passes lastEventId back unparsed; if it does not, queues with S > 1 keep their BullMQ events on one partition
pause, resume, drain, clean, obliterate, promote, retry-allqueue_admin

7.2.2 Semantics to preserve, and how

BullMQ semanticBullMQ (source)Loam
Statescompleted, failed, active, delayed, prioritized, waiting, waiting-children (src/types/job-type.ts); v6 dropped the public paused state§6.4, the same set
Priority0 = none, served before prioritized jobs; otherwise lower first, up to 2^21 − 1; score priority × 2^32 + counter (job.ts:43, getPriorityScore.lua)§6.5, the same model
Stalled jobsLock with the worker's token for lockDuration (30 s default); the checker moves a job with no lock back to wait and fails it past maxStalledCount (1 by default) with "job stalled more than allowable limit" (moveStalledJobsToWait-9.lua:48-118); jobs from schedulers are exempt§6.3, server-side, the same counts, message and exemption. The worker's stalledInterval and skipStalledCheck become no-ops
DeliveryAt least once§6.1
Delayed jobsDelayed zset scored by timestamp, promoted by markers§6.3 delayed index and promoter
Rate limiterWorker limiter {max, duration}, Worker.RateLimitError, queue.rateLimit, global rate limit and concurrency§6.9
Deduplicationsimple, throttle (ttl), debounce (extend, replace), keepLastIfActive (src/types/deduplication-options.ts)§6.2 dedup, the same four modes
Job schedulersupsertJobScheduler with pattern (cron) or every, limit, startDate/endDate, tz, offset, immediately (src/interfaces/repeat-options.ts)§8.1, Resonate schedules
FlowsParent waits in waiting-children; getChildrenValues; failParentOnFailure, continueParentOnFailure, ignoreDependencyOnFailure, removeDependencyOnFailure§8.2
Events18 event names read by QueueEvents (queue-events.ts:29-247); stream trimmed to 10,000 entries by default§6.6: the same names from the event log; retention by time instead of count (a trimEvents call maps to a retention change)
RetentionremoveOnComplete/Fail by count or age§6.10

7.2.3 Conformance

BullMQ's backend-neutral tests (tests/*.test.ts, as opposed to *.redis.test.ts) exist so that "the existing test suite can run unchanged against another backend" (create-backend.ts:108-116) (source). J2's gate is that suite, at the pinned BullMQ version, passing against operon dev with setDefaultBackendFactory(loamBackend(...)), with every excluded test listed and justified in the plan. Whether every backend-neutral test also passes on BullMQ's own Postgres backend is unverified; J2 Task 0 runs the suite on Postgres first to learn which tests are Redis-specific in practice.

7.2.4 Why not emulate Redis

The owner's direction already rejected Redis emulation; the source confirms it, and v6 makes it unnecessary:

Redis emulation (RESP + Lua in Loam)IQueueBackend (chosen)
What must match49 Lua scripts plus 67 includes, 5,243 lines, using 49 distinct Redis commands (most often EXISTS, XADD 36 times, ZSCORE, ZREM, HGET/HSET, RPOPLPUSH, LPOS, RENAME), cmsgpack in 11 scripts and cjson, key names built at runtime inside scripts (which breaks cluster slot rules without {} hash tags), BZPOPMIN on a marker zset for blocking, streams with XADD/XREAD BLOCK; a minimum of Redis 5.0 and maxmemory-policy noeviction (source)About 81 high-level methods with documented semantics
Moving targetEvery BullMQ release may change scripts; the emulator must track Lua-level behaviour, not an APIThe interface is BullMQ's public contract, with a Redis and a Postgres implementation holding it steady
Tenancy, fencing, quotasRedis has none of Loam's concepts; they would be bolted on under a key-prefix schemeNative: every call carries the namespace, and locks are Loam's fenced leases
Other usersA Redis-compatible server invites every Redis workload, which Loam does not want to supportOnly BullMQ
Build cost (estimate)A RESP server, a Lua VM with Redis semantics, the keyspace, streams and blocking commands: months, foreverA TypeScript backend of about 2,000–3,000 lines (the Redis one is 3,097) and the Rust core it needs anyway

7.2.5 Versions, Pro features and Python

  • v6 only. BullMQ v5 has no backend interface. v6 removed legacy repeatable jobs, Job#discard, debounce and the public paused state (docs/gitbook/changelog.md:193-218); v5 apps upgrade to v6 first. Owner decision (Q100, 2026-09-29): ship for BullMQ v6 only; no v5 shim.
  • BullMQ Pro is out of scope. Pro is commercial and closed (@taskforcesh/bullmq-pro); its features (groups and their rate limits and concurrency, batches, observables and cancellation) are not implemented and not imitated.
  • Python. BullMQ's Python package (3.2.7, MIT) has pluggable backends too (python/bullmq/backends). A Python loam backend is a small follow-up once the TypeScript one passes (J2.x).

7.2.6 Durable mode for BullMQ

import { Worker } from "@loam/bullmq";
import { loamResonate } from "@loam/durable";        // thin helper: new Resonate({ url, token, group })
const resonate = loamResonate({ group: "emails" });

function* sendCampaign(ctx, campaignId: string) {
  const list = yield* ctx.run(loadRecipients, campaignId);
  // At least once: a crash after sendBatch and before its checkpoint resends the
  // batch, so the mail API deduplicates on a key made from the promise id and step.
  for (const [i, batch] of chunk(list, 500).entries())
    yield* ctx.run(sendBatch, campaignId, batch, `${ctx.id}:batch:${i}`);
  return list.length;
}
const campaign = resonate.register("sendCampaign", sendCampaign);

new Worker("campaigns", async (job) =>
  (await campaign.beginRun(`bullmq:campaigns:${job.id}`, job.data.campaignId)).result());

The same rule as Celery: BullMQ keeps admission (priorities, limiter, delays, schedulers), and the promise id derived from the job id makes a redelivered or stalled job resume. FlowProducer flows map to Loam flows (§8.2), whose joins are promises; job schedulers are Resonate schedules (§8.1).

8. How Resonate is used (D210)

8.1 Schedules

  • Jobs::schedule creates a Resonate schedule (schedule.create) with a cron expression and a promise template. When it fires, the server expands {{.id}} and {{.timestamp}} in the promise id and inserts the promise with INSERT OR IGNORE, so a tick fires once even if the schedule is processed twice (process_schedule_timeout, resonate-server-sqlite lib.rs:3725-3790) (source).
  • Target. A schedule that enqueues a job targets Loam's own group (inproc://any@loam, §21 §3.5): the tick runs a small Loam durable function that calls enqueue with the tick's promise id as the idempotency key, so a tick enqueues exactly one job. A schedule that runs a durable function targets the user's group directly (poll://any@<group>).
  • every intervals (BullMQ's every, Celery's timedelta entries) that a cron expression cannot express become a Loam durable function that loops ctx.sleep(interval) → enqueue with the tick number in the idempotency key. limit, startDate/endDate, offset and immediately are fields of the schedule record that the firing function checks. Time zones (tz) are applied by Loam when it converts the rule to the server's UTC cron (verify that Resonate's cron parser has no time-zone field).
  • Upsert. schedule with an existing id and a different rule replaces it (BullMQ's upsertJobScheduler; Celery beat on restart): delete and create in one Loam operation, idempotent by id.

8.2 Flows

A flow (a Celery canvas submitted through Jobs::flow, a BullMQ FlowProducer tree, or a DAG from the SDK) is a Loam durable function in group loam, run by the Rust SDK in process (§21 §3.5):

pub struct FlowSpec {
    pub id: Option<String>,            // idempotency: the flow id = the root promise id
    pub nodes: Vec<FlowNode>,          // each: queue, task, opts, depends_on: Vec<NodeId>
    pub on_child_failure: FailurePolicy, // Fail | Continue | Ignore | Remove (BullMQ's four flags), per node
}
  1. The root promise id is the flow id. The function enqueues every node whose dependencies are met, with the step's promise id as the idempotency key, and records parent on each job.
  2. Each node has a completion promise, flow:<id>:<node>, in the root's origin (so it commits atomically with the parent, §14 §1). complete of a flow job writes an outbox action that settles that promise; the relay retries the settle until Resonate accepts it, and promise.settle of a settled promise is idempotent.
  3. The function awaits the promises of the next wave and continues: a chain is one node per wave, a group is one wave of N, a chord is a group followed by a node that depends on all of them, and BullMQ's tree is waves from the leaves up (the parent in waiting-children until its children's promises settle).
  4. A crash of the node running the flow function resumes it elsewhere after Resonate's task retry timeout (30 s by default, serve.rs:143-145), from its memoized steps.
  5. Large flows are capped per origin (§21 Q42) and split into child flows with their own origins.

BullMQ's getChildrenValues reads the children's results from the job store; the promises only carry completion.

8.3 Queue mode and durable mode

Queue modeDurable mode
What runs the taskThe framework's worker (Celery, BullMQ)The framework's worker, which calls a Resonate function
What a retry doesRuns the task again from the startResumes the function from its last finished step
Resonate state per jobNoneOne root promise (id derived from the job) plus one per step
CostOne TiKV transaction per enqueue, lease and completionPlus the durable writes of every step (SQLite or TiKV, §21 §3.3, D261), and the promise retention question (§21 Q40)
When to useShort, idempotent tasks; most jobsMulti-step tasks with side effects that must not repeat (payments, emails, paid model calls, long pipelines)

Queue-mode jobs never create Resonate promises. That keeps the per-job cost at a few TiKV writes and keeps millions of short jobs out of the durable store, whose settled promises are never pruned today (§21 §8). Durable mode is chosen per task by the user, by writing the task as a Resonate function; Loam does not decide it.

Why the queue is not built on Resonate tasks. Resonate's tasks have a fenced lease (the task version, compare-and-set on acquire, resonate-server-sqlite lib.rs:2983-2999) and a retry timeout, which is most of a queue. They lack priorities, rate limits, global concurrency, dedup modes, queue listing and counts by state, and they deliver through best-effort transports whose recovery is the 30 s retry timeout (§21 §3.4). Building those on promise tags and searches (scans on SQLite and blob backends, §21 §6.4) would be slower and harder than a purpose-built index in TiKV. So the queue is Loam's, and Resonate does what it is good at: schedules, flows and step durability.

8.4 The loam helpers

loam.resonate(group=…) (Python) and loamResonate({ group }) (TypeScript, @loam/durable) return a Resonate client configured from LOAM_URL, LOAM_TOKEN and LOAM_NAMESPACE: the durable endpoint, the bearer token (every Resonate SDK has a token option: Python token= or RESONATE_TOKEN, py/resonate.py:293; TypeScript token/tokenProvider, ts/resonate.ts:102-129), and the group. They add nothing to Resonate's API: no decorators, no wrappers around run, rpc, ctx or retry policies. Users who prefer can construct Resonate(url=…, token=…) themselves.

9. PySpark on Sail (D211)

This section follows the data-processing research of 2026-09-29 (.superpowers/research/rust-data-processing-2026-09.md, not committed; "the research" below) and this document's own check of Sail at 1f6bcde0. They agree on every point used here.

9.1 What Sail covers

Sail 0.7.1 (2026-08-24, Apache-2.0) is a Spark Connect server in Rust on DataFusion 55.1 and arrow 59.2 (source). It accepts PySpark 3.5, 4.0, 4.1 and 4.2 clients (pyspark[connect] or pyspark-client).

AreaSailSource
Spark SQL and the DataFrame APISupported. Spark 3.5.9 Connect tests: 909 passed, 93 failed (90.7 % of non-skipped); Spark 4.2.0: 1,990 passed, 801 failed (71.3 %)Research §2.6 (Sail CI on PR #2681)
Python udf, pandas_udf, UDTFs, mapInPandas, mapInArrow, applyInPandasSupported (run in Sail's embedded CPython)Sail docs, guide/dataframe/features.md
RDD, SparkContextNot available: the Spark Connect protocol does not carry them (Spark's own docs: "APIs such as SparkContext and RDD are unsupported in Spark Connect")Spark 4.2 Connect overview
Java and Scala UDFs, MLlib, pandas-on-Spark, applyInPandasWithState, ORC, CacheTableNot supported or plannedSail docs; research §2.6
Structured StreamingNot ready (Sail's pages disagree between "partial" and "planned"; Iceberg and Delta streaming are unsupported)Sail docs; research §2.6
Icebergv1–v3, merge-on-read DELETE/UPDATE/MERGE with Puffin deletion vectors, REST, Glue, Unity and HMS catalogs; no branch or tag writesguide/sources/iceberg/features.md (2026-09-28)
Object storageS3 (so RustFS), R2, GCS, ADLS, HDFS; Sail's own storage settings, not Hadoop s3aguide/storage
Deploymentsail spark server --port 50051 (Spark's own default is 15002); Kubernetes with SAIL_MODE=kubernetes-cluster: a server Deployment and Service, a driver per session that launches worker pods, object-store shuffle (0.7)guide/deployment/kubernetes.md

Consequence for D55. D55 budgets an upstream contribution of deletion-vector reads to Sail. Both the research and this check find them already on Sail's main branch. The M4 work becomes a verification against a released Sail (the research's recommendation), and this document records that as a note on D55, not a new decision.

9.2 How Loam runs Sail

  • Never linked. Sail's crates are not on crates.io, are on DataFusion 55.1 / arrow 59.2 against Loam's 54 / 58 (§11), and sail-spark-connect pulls an embedded CPython through pyo3. D51 holds: Sail is a separate process.
  • One Sail per namespace. Sail runs tenants' Python UDFs inside its own process, so a shared server would run one tenant's Python beside another's data. Each namespace that uses Spark gets its own Sail server (a Deployment in kubernetes-cluster mode in cloud; a child process in operon dev), scaled to zero after an idle period and started on the first connection (Q96).
  • The endpoint. Clients connect to sc://<ns>.spark.<cloud-domain>:443/;use_ssl=true; the URI carries connection settings only, never the token, so it cannot leak through diagnostics or proxy logs. The token travels as authorization: Bearer <loam token> gRPC metadata: loam.spark_session() builds the session with a PySpark ChannelBuilder subclass that adds that header from LOAM_TOKEN (PySpark's own token= URI parameter produces the same header, but the docs do not show it; verify the metadata path for each PySpark version). A Spark Connect proxy in Loam's gateway role authenticates the token, finds the namespace's Sail Service, starts it if needed, and forwards the gRPC stream, keeping a session on one Sail server (sticky by Spark Connect's session_id). Loam does not implement Spark Connect; it routes to Sail.
  • Data. Sail's Iceberg REST catalog points at Lakekeeper with credential vending (§10 §4), so tables live on RustFS as standard Iceberg (M4). Loam collections are read and written with the format("loam") Python data source (D54, M2), which runs on Sail unchanged (§17 §5.6).
  • Batch jobs. submit_engine_job(EngineJob::SparkBatch { entrypoint, args, conf, python_deps }) runs a PySpark script (uploaded to the object store) in a short-lived driver container against the namespace's Sail endpoint. The run is a Loam durable workflow: start the container, stream its logs into the run's events, record the exit status, retry by policy. schedule can target it, so "run this Spark job nightly" needs no Airflow.

9.3 Resonate at the job level only

A PySpark pipeline with several actions can be written as a Resonate function whose steps are Spark actions:

@resonate.register(name="daily_features")
async def daily_features(ctx, day: str):
    await ctx.run(build_sessions, day)       # spark.sql(...).writeTo("t.sessions").overwritePartitions()
    await ctx.run(build_features, day)       # a crash here resumes at build_features
    await ctx.run(publish, day)

Each step is a checkpoint; a failed step retries; a finished step is not re-run. Steps must write idempotently (an Iceberg partition overwrite or a MERGE keyed on the day, never a blind append), because a step that crashed after its write and before its checkpoint runs again (§6.1). Resonate never enters Spark tasks: task retries, stage recomputation and shuffle recovery are Sail's (or Spark's) own.

9.4 The fallback: Apache Spark on Kubernetes

Jobs that use RDDs, SparkContext, JVM UDFs, MLlib, pandas-on-Spark or Structured Streaming run on Apache Spark 4, unmodified:

JobWhere it runs
Spark Connect-compatible, but uses a feature Sail lacksApache Spark's own Spark Connect server (port 15002) for the namespace; the client only changes the URL
RDD, SparkContext, Scala or Javaspark-submit through the Kubeflow Spark Operator (Apache-2.0) as submit_engine_job(EngineJob::SparkSubmit { … })

Both read the same Iceberg tables through Lakekeeper. This is the only place Loam runs a JVM, and only when a user asks for it; Loam's default deployment has no JVM (research §6). The migration advice is the research's: run on both. loam spark check runs a job on Sail and reports the unsupported calls it hit (from Sail's errors), so a team learns which jobs need the fallback before moving them.

RisingWaveArroyo
LicenceApache-2.0; some features are "Premium" behind a licence key, and the free tier caps Premium features at 4 RWU. The Glue catalog for Iceberg is Premium; the REST catalog is notMIT OR Apache-2.0
Latestv3.1.0, 2026-09-21, very activev0.15.0, 2025-12-01; no release since; main at 0.17.0-dev; Cloudflare owns the roadmap since the 2025 acquisition
SQLPostgres dialect, not Flink SQLIts own dialect on a DataFusion 48 fork, not Flink SQL
Kafka, IcebergKafka, Pulsar, Kinesis sources; Iceberg source, sink and table engine with LakekeeperKafka source and sink; Iceberg sink only (REST catalog, two-phase commit)
StateHummock on S3 (RustFS)Checkpoints to any object store
Already in Loam's planYes: the companion stream processor (D22), returning with the Kafka gateway in M5 (D74)Named as a CI smoke client for the Kafka gateway (research §5)

Decision: Flink SQL jobs move to RisingWave, which Loam already plans as its companion. Arroyo is documented as an alternative client, not a managed engine, because of its release cadence and its forked DataFusion (Q95). Neither speaks Flink SQL, so "little change" means a dialect port: sources and sinks become CREATE SOURCE/CREATE SINK with connector='kafka', time windows use RisingWave's TUMBLE/HOP table functions, and Flink-specific hints and connectors are dropped. The docs carry a porting table with a worked example for each Flink SQL construct in Flink's own examples; an automatic translator is not planned.

submit_engine_job(EngineJob::StreamingSql { engine: RisingWave, statements }) runs the statements on the namespace's RisingWave database, as a Loam durable workflow that applies them in order, idempotently (CREATE … IF NOT EXISTS, and a recorded statement hash per step). Before M5, RisingWave reaches Loam through its Elasticsearch sink (the ES _bulk subset, M1.5), its HTTP sink (the native produce endpoint, M2) and, from M4, its Iceberg sink through Lakekeeper (§02 §7.3, research §5). It reads Loam streams once the Kafka gateway exists (M5).

There is no Rust replacement for Flink's DataStream API (research §2.1, §6). DataStream jobs therefore run on Apache Flink with the Flink Kubernetes Operator (1.16.1, 2026-09-17, Apache-2.0; Flink up to 2.4) (source):

  • Deployment. One FlinkDeployment (application mode) per job, in the namespace's Kubernetes namespace. Checkpoints and savepoints go to RustFS: state.checkpoints.dir: s3://…, s3.endpoint, s3.path-style-access: true, with flink-s3-fs-presto for checkpoints and flink-s3-fs-hadoop where a job's file sink needs a RecoverableWriter (Flink's S3 filesystem docs).
  • Reading Loam. Through Flink's Kafka connector against Loam's Kafka gateway (M5, D74), or Flink's Iceberg connector against Lakekeeper (M4). Flink needs nothing Loam-specific. Before M5, a DataStream job can write to Loam over HTTP or the ES subset but cannot read Loam streams.
  • Lifecycle through control_engine_job:
ActionOperator mechanism
DeployCreate the FlinkDeployment with upgradeMode: savepoint (or last-state; stateless only on request)
SavepointCreate a FlinkStateSnapshot resource (the operator's current mechanism; savepointTriggerNonce is deprecated); periodic savepoints with kubernetes.operator.periodic.savepoint.interval
UpgradePatch the spec (image, jar, parallelism, configuration); the operator suspends with a savepoint and redeploys from it
RollbackThe operator's rollback to lastStableSpec (kubernetes.operator.deployment.rollback.enabled, which defaults to false; Loam turns it on) or, explicitly, a redeploy of a recorded spec from a named savepoint
Suspend, resume, canceljob.state: suspended / running; delete the resource
Statusstatus.lifecycleState (CREATED … STABLE, ROLLING_BACK, ROLLED_BACK, FAILED) and status.jobStatus, mapped onto run events

Each action is a Loam durable workflow: apply the resource, wait for the operator's status, record the outcome. The steps are declarative applies, so a replay after a crash re-applies the same spec and converges. Loam does not checkpoint anything inside Flink: Flink's checkpoints and savepoints are the job's only durable state, and the workflow only records their paths.

  • Flink SQL on real Flink. A team that must keep Flink SQL unported can run it as a FlinkDeployment with the operator's SQL-runner pattern (a small jar that executes a SQL script; verify against the operator's examples at J4). Jobs submitted through the Flink SQL Gateway are not managed by the operator (operator docs), so Loam does not use the Gateway.
  • Later. A Flink VECTOR_SEARCH connector (FLIP-540, Flink 2.2) that calls Loam's hybrid search is a natural adapter for streaming enrichment (research §5, item 6). It is a small Java module outside the Rust workspace and not part of track J.

11. Tenancy and security (D213)

  • Every object is namespaced: queues, jobs, schedules, flows, engine runs, result keys, event streams (_jobs/<queue> inside the namespace) and engine deployments. The namespace comes from the credential (Ctx), never from a request body.
  • Loopback until the unified auth plan. The jobs listener binds 127.0.0.1:7720 and refuses other addresses, like the durable listener (D138) and Live (D121), because unauthenticated enqueue and schedule creation are writes. The auth plan (D111, Q30) lifts it and adds the actions jobs:enqueue, jobs:consume, jobs:admin, jobs:schedule and jobs:engine; an agent's policy (§19 §5.1) can grant jobs:enqueue on one queue only.
  • Workers are principals. A Celery or BullMQ worker holds a token with jobs:consume on its queues. A lease token is useless to another namespace's credential.
  • Payloads are opaque and stay out of ids, tags and logs (§6.7). Client-side encryption is supported by the adapters' serializers (Celery's message signing and custom serializers; BullMQ's job data is user JSON) and by Resonate's encryptor hook for durable mode.
  • Engines run in the namespace's own Kubernetes namespace with network policies that allow only the Loam endpoints, the object store and Lakekeeper. Sail and RisingWave per namespace (§9.2, Q96); Flink per job.
  • Quotas (D65): enqueue rate, payload bytes, stored jobs, active leases, schedules, and concurrent engine runs and their CPU and memory.
  • Erasure (D68): deleting a namespace deletes its job store prefix, result keys, event streams, payload objects and engine deployments. Per-subject erasure cannot look inside payloads; retention bounds how long they stay.

12. Dependencies and licences (D216)

12.1 Licence matrix

The rule: nothing linked into Loam, and nothing in a Loam-published package's required dependencies, may be AGPL, BSL, SSPL or ELv2.

DependencyLicenceHow Loam uses itLinked into the Loam binary?
connect-rust (connectrpc 0.9.1; github.com/anthropics/connect-rust redirects to connectrpc/connect-rust)Apache-2.0The JobsService serverYes (already, for Live)
buffa 0.9.2Apache-2.0Protobuf messagesYes (already)
redb 4MIT OR Apache-2.0LocalJobStoreYes (already a workspace dependency)
tikv-client (fork), operon-tikvApache-2.0TikvJobStoreYes (already, R1)
Resonate server crates and Rust SDK (fork dina-kar/resonate, loam/0.10.1)Apache-2.0Schedules, flows, lifecycle workflowsYes (already, D1)
kube-rsApache-2.0 (verify at J3)Engine runners on KubernetesYes, behind a feature, from J3
Resonate Python SDK 0.8.1, TypeScript SDK 0.11.5Apache-2.0 (the monorepo LICENSE; the Python package declares no licence field)Durable mode, the loam helpersNo: user processes
celery 5.6.3, kombu 5.6.2BSD-3-ClauseRequired by loam-celeryNo
bullmq 6.3.9MITPeer dependency of @loam/bullmqNo
connectrpc (Python) 0.12.1Apache-2.0Generated Python clientNo
@connectrpc/connect 2.2.0, @bufbuild/protobuf 2.15.0Apache-2.0; protobuf-es is Apache-2.0 AND BSD-3-ClauseGenerated TypeScript clientNo
Sail 0.7.1Apache-2.0Separate process per namespaceNo (D51)
RisingWave 3.1Apache-2.0, with licence-key "Premium" featuresSeparate service; Loam uses no Premium feature (REST catalog, not Glue)No
Arroyo 0.15MIT OR Apache-2.0Documented alternative onlyNo
Apache Flink, Flink Kubernetes Operator 1.16.1, Flink Kafka and CDC connectorsApache-2.0Separate servicesNo
Apache Spark 4, Kubeflow Spark OperatorApache-2.0The fallback (§9.4)No

Excluded: BullMQ Pro (commercial), Ververica Platform (proprietary), RisingWave Premium features (licence key), Resonate's resonate-server-scylladb and NATS pieces (BUSL-1.1 lineage, §21 §11.1). None of the linked dependencies is AGPL, BSL, SSPL or ELv2.

12.2 Published packages

PackageRegistryLicence
loam-celeryPyPIApache-2.0
@loam/bullmqnpmApache-2.0 (BullMQ's MIT notice kept for anything adapted from its Postgres backend)
loam helpers (in the existing Python SDK and a new @loam/durable)PyPI, npmApache-2.0
Generated loam.jobs.v1 clientsInside the SDKsApache-2.0

Package names follow the Loam rename (user memory: packages loamdb everywhere). If the Python SDK is loamdb, the helper is loamdb.resonate, and loam-celery becomes loamdb-celery; this document uses the short names.

12.3 Upstream proposals (none required; none posted by this document)

ProjectProposalWhy
kombuAn entry-point group for transports (kombu.transports), as kombu already has for serializers and matchersRemoves the import-order dependency of the loam alias (§7.1.1)
CeleryCapability checks instead of driver_type allow-lists for mingle, gossip and event fanoutLets third-party transports turn them on (§7.1.5)
BullMQNone expected; report any IQueueBackend gaps found by J2The interface is new in v6

13. Risks

#RiskLikelihoodImpactMitigation
J-R1BullMQ's IQueueBackend is new (v6, 2026-07) and changes between minor versionsMediumMediumPin the peer range; run BullMQ's suite per BullMQ release in CI; the backend is small
J-R2Hot keys on busy queues (the ready-index head, the rate bucket, the concurrency counter) cap throughput below RedisMediumMediumBatched leases, queue shards, sharded buckets (§6.5, §6.9); J1 measures and publishes the numbers before claims
J-R3Users expect exactly-once execution after a broker change and are surprised by redeliveryMediumHigh§6.1 in the docs; durable mode as the fix; the Celery and BullMQ semantics are unchanged from Redis
J-R4Celery features gated on driver_type (mingle, gossip) are missedLowLowDocumented; upstream proposal (§12.3)
J-R5Sail coverage (71 % on the Spark 4.2 Connect suite) disappoints teams on Spark 4HighMediumloam spark check, run-on-both, the Spark fallback (§9.4)
J-R6Arroyo's open-source cadence stallsHighLowNot a managed engine (Q95)
J-R7Flink DataStream users need Loam streams before M5MediumMediumIceberg (M4) and write paths before M5; the Kafka gateway is the dependency, not new work
J-R8Per-namespace Sail and RisingWave deployments are expensive for small tenantsMediumMediumScale to zero (Sail); shared RisingWave with a database per namespace for small tenants (Q96)
J-R9The durable store grows without bound under durable-mode jobs (§21 Q40)MediumMediumQueue mode creates no promises (§8.3); retention (Q40) before durable mode is marketed for high-volume queues
J-R10A second event path (outbox relay) lags or duplicates events before D72's idempotent producers existMediumLowSequence numbers in events until M2 (§6.6)

14. Testing

  • Queue core. Property tests over random interleavings of enqueue, lease, extend, complete, crash and clock advance against both stores: no acknowledged enqueue is lost; no job has two outcomes; a fenced write is always refused; priorities and delays are respected per shard. The TiKV store runs under operon-tikv's FaultPlan (region errors, unknown commit outcomes, TSO restarts).
  • Jepsen-style gate (with the M2/M5 stream gates): kill nodes during a busy queue; check the history for lost jobs, double completions and stuck leases.
  • Celery. Celery's own integration suite (t/integration) with broker_url = loam://… and result_backend = loam://…, the canvas tests included; plus Loam's tests for ETA without redelivery, acks_late with a killed worker, chord counters under concurrent completion, and beat twice without duplicate ticks.
  • BullMQ. The backend-neutral suite (§7.2.3).
  • Resonate paths. A schedule driven through several ticks with the durable debug clock enqueues exactly one job per tick; a flow survives kill -9 of the node running it; durable-mode Celery and BullMQ examples resume after a worker is killed mid-step, without repeating finished steps.
  • Engines. Sail: the PySpark examples of §17 and a Sail-vs-Spark comparison on a sample job. Flink: a kind cluster with the operator, a stateful job, savepoint → upgrade → rollback, checking state is kept. RisingWave: a ported Flink SQL example producing the same output as on Flink.

15. Roadmap: track J (D215)

Track J runs beside M, R and D, interleaved on the one-build machine like tracks R and D (D127). It adds new crates and routes and changes no M1 code before M1 exits.

MilestoneScopeDepends onExit gate
J0This document and its decision rows—Owner review
J1: Celeryloam.jobs.v1 and operon-jobs-proto; operon-jobs with LocalJobStore; enqueue, lease, extend, complete, fencing, delays, stalled sweep, DLQ, retention; the jobs listener in operon dev (feature jobs, loopback); schedules and the every loop on the embedded Resonate server; the generated Python client; loam-celery (transport, result backend, beat); TikvJobStoreD1 (embedded Resonate, in-process runtime), R1 (operon-tikv)Celery's integration suite green on operon dev; the property tests on both stores; the Loam Celery tests of §14
J2: BullMQThe generated TypeScript client; @loam/bullmq (IQueueBackend); flows on Resonate promises; event log and watch; @loam/durable helper; the Python BullMQ backend (J2.x)J1BullMQ's backend-neutral suite green at the pinned version, exclusions justified
J3: PySpark on SailEngine-runner framework; Sail in operon dev and on Kubernetes per namespace; the Spark Connect proxy with auth; SparkBatch jobs and schedules; the Spark fallback (Connect server and Spark Operator); loam spark checkJ1; the auth plan (D111) for the proxy beyond loopback; M2 (format("loam")); M4 (Iceberg through Lakekeeper) for tablesPySpark examples pass on Sail through the proxy; a Sail-unsupported job runs on the fallback unchanged; Iceberg tables written by Sail read back in Loam (M4)
J4: FlinkStreamingSql on RisingWave; the Flink operator lifecycle (deploy, savepoint, upgrade, rollback, suspend) as workflows; the porting guide; Arroyo documentedJ1, J3's runner framework; M5 (Kafka gateway) for reading Loam streams; M4 for IcebergA ported Flink SQL example matches Flink's output; a DataStream job keeps its state across upgrade and rollback on a kind cluster, reading Loam through the Kafka gateway

The first PRs of J1, in order, each small: (1) the protos and generated crate; (2) operon-jobs types, the JobStore trait and LocalJobStore with enqueue/lease/complete and fencing; (3) delays, stalled sweep, retention, DLQ; (4) the listener and operon dev wiring; (5) the Python client and the kombu transport; (6) the result backend and chord counter; (7) schedules on Resonate and LoamScheduler; (8) the event outbox and relay; (9) TikvJobStore under the fault plan; (10) the Celery integration suite in CI.

16. Contradictions with earlier decisions, and how they are resolved

EarlierConflictResolution
D42 and §17 §6.1: "Loam serves no Spark Connect endpoint: Sail is the Spark Connect server; Loam is its source and sink"The owner's direction has Loam host Sail as the Spark Connect endpointD211 amends §17 §6.1: Loam still implements no Spark Connect; it runs Sail per namespace and routes sc:// connections to it. Sail stays a separate process (D51 holds)
D51: no embedded enginesLoam now runs Sail, RisingWave and Flink for usersThey are managed services beside Loam, never linked. D51 is about the binary, and it holds
§00 §7: stay out of stateful stream processingLoam manages Flink and RisingWave jobsLoam manages their lifecycle; the stateful processing and its checkpoints are Flink's and RisingWave's (D212)
D55: an upstream contribution of DV reads to Sail in M4Sail's main already has them (§9.1)A note on D55: M4 verifies against a released Sail instead of contributing
D1: object storage is the only source of truthJob state is in TiKV (or redb in dev)Like Live (D130) and durable execution (§21 §15), jobs are a Loam cloud service whose store is TiKV; payloads and results are on object storage
The owner's "drop-in TS package with the same API"BullMQ v6's backend interface makes re-implementing the API unnecessarySame user-facing result, less code: @loam/bullmq re-exports BullMQ's own classes bound to a Loam backend (D209). Changing the import still works
The owner's "Resonate decorators"Queue-mode jobs do not use ResonateDurable mode is opt-in per task with Resonate's own decorators (D210); queue mode stays promise-free for cost (§8.3)

17. Open questions

#QuestionNeeded by
Q95Arroyo: documented alternative only (proposed), or a managed engine beside RisingWave; revisit if Arroyo resumes releasesJ4 plan
Q96Engine tenancy: Sail per namespace with scale-to-zero (proposed); RisingWave per namespace, or a shared cluster with a database per small namespaceJ3 plan
Q97Jobs on self-hosted clusters without TiKV Decided by the owner, 2026-09-29: jobs require TiKV; redb for operon dev only; no Postgres or DynamoDB jobs backend (§6.8, D214)Resolved
Q98A server-streamed LeaseStream for low-latency workers, beside the long-polled unary LeaseJ2 plan
Q99Strict per-key ordering (FIFO groups, like SQS FIFO): out of scope (proposed), or a later queue kindAfter J2
Q100BullMQ v5 apps Decided by the owner, 2026-09-29: ship for BullMQ v6 only; no v5 shim (§7.2.5, D209)Resolved
Q101Celery beat's solar schedules and custom schedule classes: refuse and keep the standard scheduler for them (proposed), or support them in LoamJ1 plan
Q102Loam-hosted workers: Celery and BullMQ workers are user processes in J1–J2; whether Loam hosts them (with the CPU-time runtime of doc 24) is decided thereDoc 24
Q103The jobs listener on its own port (7720, proposed, like Live's 7710) or mounted on the native API listener as a connect-rust serviceJ1 plan
Q104The Spark fallback: Loam-managed (Spark Connect server and the Kubeflow Spark Operator per namespace) or documented onlyJ3 plan

18. Sources

  • Celery f0b1320 (main, 2026-09-28; release 5.6.3) and kombu b1ba4ba (main, 2026-09-28; release 5.6.2): kombu/transport/__init__.py, kombu/connection.py, kombu/transport/virtual/{base,exchange}.py, kombu/transport/{redis,pgmq,gcpubsub}.py, kombu/transport/SQS/__init__.py, kombu/transport/native_delayed_delivery.py; celery/app/{base,backends,amqp,control,task,builtins,defaults,trace}.py, celery/backends/{base,redis}.py, celery/canvas.py, celery/beat.py, celery/result.py, celery/events/event.py, celery/worker/{request,strategy}.py, celery/worker/consumer/{consumer,control,mingle,gossip,tasks}.py, docs/getting-started/backends-and-brokers/{redis,sqs}.rst, docs/userguide/periodic-tasks.rst.
  • BullMQ d1a43ab (master, 2026-09-28; package 6.3.9): src/interfaces/queue-backend.ts, src/utils/{create-backend,with-backend}.ts, src/classes/{queue,queue-base,queue-getters,worker,job,flow-producer,queue-events,queue-keys,redis-connection,redis-queue-backend}.ts, src/commands/*.lua, src/postgres/, src/interfaces/{worker-options,base-job-options,repeat-options,backoff-options}.ts, src/types/{job-options,job-type,deduplication-options,processor}.ts, docs/gitbook/guide/{connections,postgresql}.md, docs/gitbook/changelog.md, docs/gitbook/bullmq-pro/, python/pyproject.toml.
  • Resonate: the fork dina-kar/resonate at e360669 (branch loam/0.10.1 = upstream 28dfd01 plus five commits: RustSec/OpenSSL hygiene, the GCP ID-token feature gate, MySQL errno classification, TiDB support, SDK-rs without default TLS); impl/sdk/py/src/resonate/{resonate,context,retry,schedules,promises,codec}.py (0.8.1), impl/sdk/ts/src/{resonate,context,options,retries}.ts (0.11.5), impl/server/core/crates/resonate-core/src/types.rs, resonate-server-sqlite/src/lib.rs, core/src/serve.rs. Loam: crates/operon-durable on main (the registry, the in-process network and runtime); the TiKV server plugin (crates/operon-durable/src/tikv.rs) is in progress in the working tree, not on main.
  • connect-rust fb5f5aa (github.com/connectrpc/connect-rust; tag v0.9.1); connectrpc 0.12.1 on PyPI; @connectrpc/connect 2.2.0 and @bufbuild/protobuf 2.15.0 on npm.
  • Sail 1f6bcde (0.7.1): docs/introduction/migrating-from-spark, docs/guide/{dataframe/features,sources/iceberg/features,catalog/index,storage,deployment/kubernetes,cli}.md, k8s/sail.yaml, Cargo.toml. Apache Spark 4.2 Spark Connect overview; Spark 4.1.0 release notes.
  • Flink Kubernetes Operator de630f6 (1.16.1): helm/*/crds, docs/content/docs/managing/snapshot-management.md, docs/deployment/overview.md, ResourceLifecycleState.java. Flink 2.3 S3 filesystem and SQL Gateway docs.
  • Arroyo 0630bec (v0.15.0): LICENSE-*, arroyo-api/src/rest.rs, arroyo-rpc/default.toml; the 2025-04-10 Cloudflare announcement. RisingWave 5149d7b (v3.1.0): Iceberg overview and Premium features docs.
  • Data-processing research of 2026-09-29: .superpowers/research/rust-data-processing-2026-09.md (not committed): §1 table, §2.1 Flink, §2.2 Arroyo, §2.4 RisingWave, §2.6 Sail, §4–§6.
  • Loam: §01 §3 (ports), §02 §7 (stream API, idempotent producers), §03 (freshness), §09 §3 (leases), §11 §1, §14, §17 §5.6 and §6.1, §18 §5–§6, §19 §5, §20 §5, §7, §11, §21; crates/operon-meta-tikv/src/leases.rs (check_fence), crates/operon-common/src/meta/types.rs (Fence, Lease), crates/operon-live-proto/build.rs, proto/loam/live/v1/idempotency.proto; crates/operon-stream-grpc (in progress in the working tree, not on main: a tonic StreamService.Produce without producer ids yet); D1, D22, D42, D51, D54, D55, D65, D68, D72–D74, D76, D111, D118, D121, D127, D128, D130, D138–D147.

On this page

1. Summary (D204)2. Goals and non-goals2.1 Goals2.2 Non-goals3. Architecture3.1 Components3.2 Where each piece of state lives (D214)4. Concepts5. The Rust API (D205)5.1 The trait5.2 The main types5.3 Errors5.4 Admin operations5.5 The protobuf surface (D206)5.6 Method → backing primitives6. Semantics (D207)6.1 Delivery6.2 Idempotency and deduplication6.3 Leases and fencing6.4 States6.5 Priorities, ordering and LIFO6.6 The event log and watch6.7 Payloads and results6.8 Stores: TiKV everywhere, redb in dev6.9 Rate limits and concurrency6.10 Dead-letter queues and retention6.11 Observability7. Adapters7.1 Celery: loam-celery (D208)7.1.1 Registration7.1.2 The transport7.1.3 Priorities, ETA and retries7.1.4 The result backend7.1.5 What depends on the broker type, and what does not work7.1.6 Beat7.1.7 Canvas7.1.8 Durable mode for Celery7.2 BullMQ: @loam/bullmq (D209)7.2.1 The backend mapping7.2.2 Semantics to preserve, and how7.2.3 Conformance7.2.4 Why not emulate Redis7.2.5 Versions, Pro features and Python7.2.6 Durable mode for BullMQ8. How Resonate is used (D210)8.1 Schedules8.2 Flows8.3 Queue mode and durable mode8.4 The loam helpers9. PySpark on Sail (D211)9.1 What Sail covers9.2 How Loam runs Sail9.3 Resonate at the job level only9.4 The fallback: Apache Spark on Kubernetes10. Flink (D212)10.1 Flink SQL: RisingWave by default, Arroyo as an option10.2 DataStream jobs: unmodified Flink, managed by the operator11. Tenancy and security (D213)12. Dependencies and licences (D216)12.1 Licence matrix12.2 Published packages12.3 Upstream proposals (none required; none posted by this document)13. Risks14. Testing15. Roadmap: track J (D215)16. Contradictions with earlier decisions, and how they are resolved17. Open questions18. Sources