DreamLake

Task queue

Strip the launch policy away and a Lakeshore queue is an ordinary task queue: producers submit work, workers claim it, results come back by id. This page is that half on its own — everything here holds whether or not a provider is attached.

Set providerRef: null and the queue is only this. Nothing autoscales, nothing is provisioned, and the queue runs on whatever workers you attached yourself — a lab machine, an allocation you already hold, your laptop. Adding a provider is covered in Elastic queues; it changes who supplies the workers, not the behaviors below.

Ids, not futures

submit returns a durable string id — a ULID — not a handle bound to your process:

python
inv = f.submit(x=1)          # -> "01JD2K7Q...": a string, persist it anywhere
result = q.result(inv)       # poll from this process, or any other, later

That is the load-bearing choice. There is no Future type in the SDK, no gather, no as_completed. A submitting process can exit; another can pick the id up and read the result. Reconnection is the default case rather than a recovery path.

The lifecycle

An invocation moves through a small state machine:

              submit()
                 │
                 ▼
            [ queued ]  ──── a worker claims it
                 │
                 ▼
            [ running ]
           ┌─────┼─────┬──────────┐
           ▼     ▼     ▼          ▼
     [succeeded] [failed] [killed] [timeout]

queued and running are the live states; the other four are terminal. The terminal state is what q.result(inv) resolves against — a failed invocation carries its error rather than raising at submit time.

If you know zaku

The shape is the same one zaku uses — created → in_progress → completed | failed — with two extra terminal states for the cases a remote worker adds: killed when the invocation is cancelled out from under the worker, and timeout when it exceeds timeout_s. The practical difference is where the payload lives: zaku puts it in S3 under zaku/{queue}/{job_id} and the client deletes it on completion, while Lakeshore keeps small payloads inline as msgpack on the invocation and spills larger ones to a storage reference.

Idempotent submits

Pass _key= and the submit becomes idempotent:

python
inv = f.submit(x=1, _key="nightly-2026-08-06")

Resubmitting the same key returns the existing invocation id rather than creating a second one. The key is unique per namespace + queue, so the same key against a different queue is a different job. This is the retry-safe path for anything a scheduler might fire twice.

Streaming results

A generator body does not have to finish before you can read from it. Yields become VALUE frames in a journal, terminated by an END frame, and the stream replays from a cursor:

python
for frame in q.stream(inv, since=cursor):
    ...

since is a 0-based sequence cursor, so a consumer that drops off can resume from where it stopped instead of re-reading the whole stream. Topic is the same journal spelled as pub/sub.

Back-pressure

The admission block on the queue is where depth limits are declared — max_depth, on_full: reject | block, rate_limit, deadline_cutoff_s.

Admission is not enforced yet

These fields are validated and persisted, but nothing applies them at submit time. A queue with max_depth: 100 will accept the 101st invocation today. Treat the block as declared intent until the reference says otherwise.

Running a worker

The whole plane fits in one process when you want to test:

python
import dreamlake.lakeshore as dls

q = dls.Dispatch(":memory:")      # a full plane, in-process
dls.run_worker(q, once=True)      # drain exactly one job and return

once=True drains a single job, which is what makes a test deterministic. Drop it for a worker that keeps pulling.

Bodies and I/O

One decorator covers four body kinds — sync, generator, async, and async-generator. inspect picks the shape once at decoration; there is no mode flag and no second decorator. See Simple functions for the decorator itself and UDF patterns for the shapes side by side.

File bytes never ride the wire. Bodies read and write through dls.run.read / dls.run.write under a dls.scope(prefix) and return the keys they wrote, so the payload on the invocation stays small and the data stays in storage.

Elastic queues →

Attach a provider and the same queue starts supplying its own workers.

Queue model →

The record behind all of this, field by field.