Invocation ids
There is no Future type in the SDK. A pending job is its invocation id — a plain string (a ULID), returned by every submit:
The id is the whole handle. It is durable state on the dispatch plane, not an object tied to your process — so everything a Future does, a string does better:
| With a Future you'd… | With an id you… |
|---|---|
| hold the object until it resolves | q.result(inv, timeout=...) |
| lose the handle when the process dies | write the id to a file or DB; any process re-attaches |
| await it | await q.result(inv) on the async Queue |
| iterate streamed results | q.stream(inv, since=cursor) |
| cancel it | q.cancel(inv) |
Reconnect-by-id
Because the id is just a string, crash recovery is trivial — persist ids at submit time, re-attach later from anywhere:
_key (the idempotency key) closes the last gap: if process A crashes
between submitting and recording, rerunning it resubmits the same keys
and gets the same ids back — no duplicate jobs.
Fan-out + collect
The Python for-loop is the spawn; the result loop is the barrier. There is no scheduling DSL — the call graph is the program.
Async fan-out composes with stdlib asyncio:
As-they-finish is asyncio.as_completed over the same coroutines:
dls.gather and dls.as_completed do not exist. They belonged to the
pre-rewrite surface and were deleted. Use the stdlib versions over
q.result coroutines, or a plain loop on SyncQueue.
Pipelines
Chain stages by feeding one collect into the next comprehension — that
is the whole DAG mechanism (examples/02_pipeline.py):
The SDK has no PipelineRun, no run handle, and no DAG object — a
pipeline is the ids you keep. This is an acknowledged gap (the F7 row in
the SDK repo's docs/execution-patterns.md); see
F7 · The run object.
Errors
q.result raises DispatchError when the invocation failed or was
cancelled (the message carries the worker's error) and TimeoutError
when the wait window elapses. A timeout is not a failure — the id is
still live; call result again with a longer window.
DispatchError lives in dreamlake.lakeshore.dispatch:
The default timeout on result is 30 seconds — including for the
implicit fetch behind a plain call on a queue-bound UDF. Long jobs want
an explicit submit + result(timeout=...) pair.
Streaming
A generator or async-generator body appends one frame to the invocation's journal per yield; consumers iterate live while the job runs:
The journal has an explicit terminal frame, so producer finished and
connection dropped are distinguishable: END stops iteration cleanly,
ERROR raises DispatchError carrying the producer's exception type
and message, and silence past timeout raises TimeoutError.
Cursor reconnect
Every frame carries a monotonic seq, zero-based per invocation. Pass
the last seq you saw as since to resume without loss or duplication:
The journal is derived data — the terminal result always lives on the
invocation row, so even a trimmed journal degrades to q.result(inv),
never to wrong state.
Read next
- Queue API — the verbs these ids flow through.
@udfdecorator —submiton a queue-bound function, and streaming bodies.- Dispatch planes — where invocation rows and journals actually live.