DreamLake

F4 — async

The async row of the execution matrix. Write the body with async def and @dls.udf does the rest — the decorator inspects the function once at decoration, so a plain call returns a coroutine and await f(x) is the whole calling convention. No mode flag, no second decorator.

On the queue side, dls.Queue is the async surface. Note the naming is inverted from most SDKs: the unprefixed name is the async one, and dls.SyncQueue is the blocking one. Queue is not a subclass of SyncQueue — it wraps one and exposes it through the .sync property.

L0 — an async body

A plain call returns a coroutine; nothing runs until you await it.

python
import asyncio

import dreamlake.lakeshore as dls


@dls.udf
async def embed(x: float) -> float:
    await asyncio.sleep(0.01)   # real async I/O goes here
    return x * 2.0


async def main() -> None:
    coro = embed(21.0)          # plain call -> coroutine; nothing ran yet
    print(await coro)           # 42.0


asyncio.run(main())

This is the local tier — same process, same event loop, with the run-context fork and the return-your-keys check installed around the await exactly as they are around a sync call.

L1 — concurrent fan-out with asyncio.gather

Local async UDF calls are ordinary awaitables, so the stdlib's asyncio.gather is the fan-out primitive. Five 50 ms calls take 50 ms, not 250.

python
import asyncio
import time

import dreamlake.lakeshore as dls


@dls.udf
async def score(x: int) -> int:
    await asyncio.sleep(0.05)
    return x * x


async def main() -> None:
    t0 = time.perf_counter()
    vals = await asyncio.gather(*(score(x) for x in range(5)))
    print(vals, f"{time.perf_counter() - t0:.2f}s")   # [0, 1, 4, 9, 16] 0.05s


asyncio.run(main())

That is asyncio.gather, from the standard library. There is no dls.gather — the SDK has no gather, no as_completed, and no Future type at any level.

L1 — the async queue surface

dls.Queue mirrors SyncQueue's producer verbs — submit, result, call, stream, cancel, count — plus the low-level worker verbs claim, ack, and nack. Each is await-able and internally offloaded with asyncio.to_thread, so a slow result poll never blocks the loop. topic() stays synchronous and returns a Topic. The blocking-only helpers — working(), post(), and len(q) — live on SyncQueue and are reachable through q.sync.

The handle is still a durable string id, never a Future.

python
import asyncio

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

plane = Dispatch(":memory:")


@dls.udf(queue="gpu")
async def embed(x: float) -> float:
    await asyncio.sleep(0)
    return x * 2.0


async def main() -> None:
    q = dls.Queue("gpu", dispatch=plane)          # async surface, same verbs
    inv = embed.submit(21.0, _dispatch=plane)     # durable string id
    await asyncio.to_thread(dls.run_worker, q.sync, once=True)
    print(await q.result(inv))                    # 42.0

    val = await embed(21.0, _dispatch=plane)      # submit + wait, loop stays live
    print(val)


asyncio.run(main())

A queue-bound async UDF can be awaited directly, as the last two lines show: the plain call routes to .remote(), which submits and then awaits the result on a worker thread. Note that this path resolves with the 30-second default and raises TimeoutError past it — for longer jobs, submit and await q.result(inv, timeout=…) explicitly.

async for frame in q.stream(inv) works on this surface too, but a plain async body posts no journal frames — the stream ends cleanly with zero values once the invocation reaches a terminal state. Frames come from generator and async-generator bodies; streaming is the F5 row.

L2 — an await chain

A fixed pipeline is just sequential awaits; each stage's resolved value feeds the next.

python
import asyncio

import dreamlake.lakeshore as dls


@dls.udf
async def fetch(scene: str) -> dict:
    await asyncio.sleep(0.01)
    return {"scene": scene, "frames": 12}


@dls.udf
async def calibrate(raw: dict) -> dict:
    await asyncio.sleep(0.01)
    return {**raw, "calibrated": True}


@dls.udf
async def summarize(doc: dict) -> str:
    await asyncio.sleep(0.01)
    return f"{doc['scene']}: {doc['frames']} frames, calibrated={doc['calibrated']}"


async def main() -> None:
    raw = await fetch("scene7")
    doc = await calibrate(raw)
    print(await summarize(doc))   # scene7: 12 frames, calibrated=True


asyncio.run(main())

L3 — dynamic spawn off resolved values

The graph's shape is discovered while running: a probe resolves, its value picks the fan-out width, asyncio.create_task spawns the wave, and asyncio.as_completed consumes it in completion order. A second wave spawns off the first wave's values — ordinary Python control flow is the DAG language.

python
import asyncio

import dreamlake.lakeshore as dls


@dls.udf
async def probe(scene: str) -> list[int]:
    await asyncio.sleep(0.01)
    return [3, 1, 4, 1, 5]          # width discovered at run time


@dls.udf
async def refine(seed: int) -> float:
    await asyncio.sleep(0.01 * seed)
    return seed * 1.5


async def main() -> None:
    seeds = await probe("scene7")                # resolved value picks the width
    tasks = [asyncio.create_task(refine(s)) for s in seeds]

    first_wave = []
    for fut in asyncio.as_completed(tasks):      # completion order, not submit order
        first_wave.append(await fut)

    # second wave, spawned off the first wave's values
    rework = [asyncio.create_task(refine(int(v))) for v in first_wave if v > 4.0]
    print(sorted(first_wave), [await t for t in rework])


asyncio.run(main())

Both create_task and as_completed here are stdlib asyncio over local coroutines. Nothing crosses a queue, so nothing is durable; the moment you want the wave to survive the process, the handles become ids and the row is F2.

What is awkward today

  • There is no async worker verb. dls.run_worker takes the sync queue, so draining from async code means reaching through the async surface for its blocking sibling and offloading by hand: await asyncio.to_thread(dls.run_worker, q.sync, once=True). The .sync property makes it a one-liner, but the composition is the caller's job. An arun_worker(q) — or run_worker accepting either sibling — would close the seam.
  • Waiting for a worker to pick up a job is a poll loop. once=True returns 0 when the queue is empty, so "drain exactly one job, whenever it arrives" from async code is a retry loop around to_thread. There is no awaitable "a job was executed" signal.
  • The async offload is thread-per-verb. Every Queue verb rides asyncio.to_thread over the blocking core — fine at notebook scale, but a native-async transport is what would make it honest at high concurrency. The Queue class docstring reserves exactly that: an httpx-native async transport slotting in behind this same surface when the remote control plane wire lands. It is not built.

All four body kinds under the one decorator — including the async pair — run side by side in examples/03_body_kinds.py in the dreamlake-lakeshore repo (python examples/03_body_kinds.py).