DreamLake

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:

python
inv = q.submit({"seed": 42})       # 'inv' is a str
inv = train.submit(seed=42)        # same for a queue-bound @udf

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 resolvesq.result(inv, timeout=...)
lose the handle when the process dieswrite the id to a file or DB; any process re-attaches
await itawait q.result(inv) on the async Queue
iterate streamed resultsq.stream(inv, since=cursor)
cancel itq.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:

python
# process A — submit and record
ids = [train.submit(seed=i, _key=f"sweep/{i}") for i in range(100)]
Path("ids.txt").write_text("\n".join(ids))

# process B — hours later, on another machine
q = dls.SyncQueue("gpu")
for inv in Path("ids.txt").read_text().splitlines():
    print(inv, q.result(inv, timeout=600))

_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.

python
@dls.udf(queue="compute")
def square(x: int) -> int:
    return x * x

q = dls.SyncQueue("compute")

ids = [square.submit(i) for i in range(100)]            # fan out
results = [q.result(i, timeout=60.0) for i in ids]      # collect, in order

Async fan-out composes with stdlib asyncio:

python
q = dls.Queue("compute")
ids = [await q.submit({"x": i}) for i in range(100)]
results = await asyncio.gather(*(q.result(i, timeout=60) for i in ids))

As-they-finish is asyncio.as_completed over the same coroutines:

python
for coro in asyncio.as_completed([q.result(i) for i in ids]):
    handle(await coro)
These are asyncio's functions, not the SDK's

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):

python
raw = [q.result(i) for i in [download.submit(u) for u in urls]]
out = [q.result(i) for i in [process.submit(r) for r in raw]]
There is no pipeline run object

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.

python
try:
    val = q.result(inv, timeout=10.0)
except TimeoutError:
    ...            # still running; the id remains valid
except DispatchError as e:
    ...            # the body raised, or the invocation was cancelled

DispatchError lives in dreamlake.lakeshore.dispatch:

python
from dreamlake.lakeshore.dispatch import DispatchError

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:

python
@dls.udf(queue="gpu")
async def traj_to_ego_views(traj: str):
    for i, frame in enumerate(render(dls.run.read(traj))):
        frame.save(dls.run.write(f"images/{i:04d}.png"))
        yield [f"images/{i:04d}.png"]          # a partial manifest, live

inv = traj_to_ego_views.submit("outputs/0007/traj.npz")

q = dls.Queue("gpu")
async for manifest in q.stream(inv):           # frames as they land
    print(manifest)                            # ['images/0000.png'], ...

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:

python
frames = list(q.stream(inv, timeout=5))          # all four frames
tail   = list(q.stream(inv, since=1, timeout=5)) # frames 2 and 3

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.
  • @udf decorator — submit on a queue-bound function, and streaming bodies.
  • Dispatch planes — where invocation rows and journals actually live.