Transports

A transport is the queue that lets many workers share the load while guaranteeing each execution is advanced by only one worker at a time. This is the hub for the per-backend reference: the invariant they all uphold, the contract, the shared lease/parking mechanics — and a page per backend with its exact queue layout and every operation (the table/keys/documents, how claim leases a group, how ack/nack work, and why). For the concepts (the seam, running workers, scaling) start with distribution; to select one at the worker see STM_TRANSPORT_BACKEND.

The invariant: single active consumer per group

The transport is a FIFO queue partitioned into groups, where group_id = execution_id. The one property the whole design rests on:

at most one message per group is in flight at any moment.

That is what makes concurrency safe — events for one execution are processed in order, by one worker, even with a pool of workers running — and it upholds the store’s single-writer invariant (the store CAS is the backstop if a lease expires). Different groups are claimed by different workers in parallel, which is where the throughput comes from.

The contract every backend implements

Transport is a small Protocol (transport/_base.py) — each backend is a sibling module under harel/engine/transport/, with a twin under harel/engine/aio_transport/:

Method

Purpose

publish(group_id, event, priority=0)

enqueue event in group_id’s FIFO; sets the group’s priority on first publish

claim(worker_id, visibility, min_priority=0)

lease the oldest message of some group with nothing in flight, for visibility seconds, considering only groups with priority >= min_priority; None if nothing deliverable

ack(lease)

the message was handled — remove it, freeing its group

nack(lease, delay=0)

return it to the queue: delay=0 re-claimable now (retry); delay>0 parks it (and blocks its group) until delay passes

close()

release the connection/client

Three mechanisms recur across the backends, which each per-backend page then spells out exactly:

  • The lease / visibility timeout. claim marks a group in-flight until now + visibility. If the worker acks, the message is gone; if the worker dies, the lease expires and the message becomes claimable again — crash recovery with no separate sweeper. (Like any lease-based queue, a lease that expires mid-ack can let two workers briefly touch one group; the store CAS is the backstop.) A Lease carries the group_id, the event, and a handle to identify it on ack/nack — a seq (row id) or a fencing token.

  • Parking (nack with delay). The control plane parks a suspended group’s message rather than spinning a worker on it — see control plane. The _PARKED sentinel marks a message that is non-claimable until its delay elapses.

  • Round-robin fairness. When a group is serviced, its claimability timestamp is set to now (on claim for the SQL / in-memory backends, on ack for postgres / mongo / redis) — not to 0/earliest — so it yields to groups that haven’t been serviced recently. Groups that have never been claimed have timestamp 0 and always go first. This prevents a handful of high-traffic executions from monopolizing worker capacity — the hot-partition problem — while still processing every group fairly.

  • Priority (min_priority). A group’s priority (0-4, fixed at the execution’s create) lets a worker configured with high_ratio>0 prefer high-priority groups (see distribution). Backends that keep their own group index honour it exactly — including finding a high-priority group behind a large backlog of low-priority ones (Redis does this with one ready set per priority level). SQS is the exception: FIFO has no per-group priority, so SqsTransport rejects priority>0 / min_priority>0 (fail-fast) rather than silently ignoring them; its cross-MessageGroupId delivery still gives round-robin fairness.

Backends differ only in how they enforce per-group exclusivity: a database that has a cheap serialization primitive (SQLite’s write-lock, Postgres’s row-lock, SQS’s native MessageGroupId) leans on it; the others (Redis, Mongo) build a per-group lock record by hand, indexed by when the group next becomes claimable so claim reads only the few due groups instead of scanning.

The backends

Each backend has its own page with the full queue layout and every operation broken down with the real SQL/commands and the reasoning:

Backend

Per-group exclusivity mechanism

InMemoryTransport

in-process list + lock (tests, single process)

SqliteTransport

BEGIN IMMEDIATE (global write-lock) + lease

LibsqlTransport

BEGIN IMMEDIATE + lease over libSQL (experimental)

RedisTransport

per-priority ready:{prio} ZSETs + per-group lock; publish/claim/ack are atomic Lua (one round-trip each)

PostgresTransport

per-group row, FOR UPDATE SKIP LOCKED; claim/ack are plpgsql functions (one round-trip each)

RqliteTransport

one serialized UPDATE (Raft orders it)

MongoTransport

per-group ready-index/lock doc; claim is one atomic sorted find_one_and_update

SqsTransport

SQS FIFO MessageGroupId (native) + visibility-timeout lease

Store and transport are independent

They are separate seams — mix freely (e.g. Postgres store + Redis transport) or unify on one backend for a simpler stack: all-Postgres, all-rqlite, all-Mongo, or all-libSQL all need no Redis. The async twins under harel/engine/aio_transport/ are what the async worker uses; the sync classes back the embedded DistributedRunner façade and tests. Each per-backend page has an Async twin section.