DreamLake

F2 — deferred ids

Row F2 of the execution matrix: deferred calls. submit() returns a durable string invocation id (a ULID) — not a Future, not a handle object. The id is the whole contract: persist it, poll it, resolve it, from this process or any other.

LevelShapeSpelling
L0one callf.submit(x) → id
L1fan-out[f.submit(x) for x in xs]
L2static pipelineresolved values feed stage 2
L3dynamic graphwave 2 shaped by wave 1's results

All snippets below share one setup — an in-memory plane, one queue, two UDFs. Draining is dls.run_worker(q, once=True), one call per job.

python
import dreamlake.lakeshore as dls
from dreamlake.lakeshore.dispatch import Dispatch

plane = Dispatch(":memory:")
q = dls.SyncQueue("ids", dispatch=plane)

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

@dls.udf(queue="ids")
def add(a: int, b: int) -> int:
    return a + b

L0 — one deferred call

Submit, drain, resolve. The id is a plain string from the moment submit returns.

python
inv = square.submit(7, _dispatch=plane)   # "01KX…" — a plain string
assert isinstance(inv, str)

dls.run_worker(q, once=True)              # drain one job
assert q.result(inv, timeout=5) == 49
Control kwargs on UDF.submit are underscore-prefixed

UDF.submit forwards everything else straight to your function, so its own knobs are namespaced out of the way: _dispatch=, _key=, _priority=. The queue-level SyncQueue.submit has no user function to collide with and uses the plain names envelope=, key=, priority=. Mixing them up does not fail at the call site: f.submit(x, key="…") packs key into the envelope as one of your function's keyword arguments, and the job blows up inside the worker instead.

L1 — fan-out

A list comprehension over inputs is a list of ids. There is no gather and no as_completed; resolve in submission order and the results line up with the inputs by position.

python
invs = [square.submit(x, _dispatch=plane) for x in range(5)]
for _ in invs:
    dls.run_worker(q, once=True)

results = [q.result(i, timeout=5) for i in invs]
assert results == [0, 1, 4, 9, 16]        # submission order

L2 — static pipeline

Stage boundaries are resolution points. Resolve stage 1, feed the values to stage 2 — the DAG is ordinary Python control flow, not a graph API.

python
stage2 = add.submit(results[3], results[4], _dispatch=plane)
dls.run_worker(q, once=True)
assert q.result(stage2, timeout=5) == 25  # add(9, 16)

L3 — dynamic graph

The second wave's width depends on the first wave's results — structure discovered at run time. Ids are the edges; a hand-kept edge list is the run record, and it serializes for free because every node is a string.

python
edges = []                                # (parent_id, child_id) pairs
wave2 = []
for parent, r in zip(invs, results):
    if r > 5:                             # width decided by the data
        child = square.submit(r, _dispatch=plane)
        edges.append((parent, child))
        wave2.append(child)

for _ in wave2:
    dls.run_worker(q, once=True)
assert [q.result(c, timeout=5) for c in wave2] == [81, 256]

Two of five first-wave results cleared the bar, so the second wave has two nodes. Dump edges to a journal line and the whole run replays. Keeping that list by hand is exactly the work the missing run object would take over.

Idempotent submits

The same key on the same queue returns the same id — resubmit after a crash and you get the invocation you already had, not a duplicate job. Idempotency is scoped per namespace, queue, and key.

python
a = square.submit(7, _key="scene7/x", _dispatch=plane)
b = square.submit(7, _key="scene7/x", _dispatch=plane)
assert a == b                             # same key → same id

The queue-level spelling drops the underscore. It enqueues a raw wire payload rather than a UDF envelope, so it is claimed by a hand-rolled worker (q.claim() / q.working()), not by dls.run_worker — that loop decodes every job as an Envelope:

python
c = q.submit({"x": 1}, key="scene7/y")
d = q.submit({"x": 1}, key="scene7/y")
assert c == d

An id is a fact, not state

A Future is process-bound state: it dies with the event loop that made it, cannot be written to disk, and means nothing to another process. An id is a fact about the world — "invocation 01KX… exists on queue ids" — true regardless of who holds it. That is why the whole Future family was deleted: every capability it had is a string plus q.result, and the string survives restarts.

Persist the id, resolve it later — the reader needs nothing but the id and a queue on the same plane.

python
with open("/tmp/run-id.txt", "w") as f:
    f.write(inv)                          # the checkpoint is one line

later = open("/tmp/run-id.txt").read()    # any process, any time
assert q.result(later, timeout=5) == 49

Note that q.result defaults to timeout=30.0 seconds and raises TimeoutError past the deadline; pass an explicit timeout= for anything long-running.

Full example

examples/04_ids_not_futures.py in the dreamlake-lakeshore repo walks exactly this row, L0 through L3 plus idempotency keys, runnable with zero infrastructure:

bash
python examples/04_ids_not_futures.py