DreamLake

Synthesis of the three proposals into one buildable design. Where the three converged (ExtType codec, one wire path, Redis-Streams journal, renamed queue verbs, gather-over-futures) I take the convergent answer. Where they split — the data-UDF file story — I resolve it with the dual-ergonomics filter as the tie-breaker, and say exactly why.


1. Recommended shape + guiding principles

The stance on new runtime types: exactly ONE, and it is a pathlib.Path subclass

The three proposals split precisely here:

  • primitives-only → zero types; a bare Path is the file tag, provenance kept in a module-global _REFS: dict[str, FileRef] keyed by local-path string.
  • annotation-driven → zero runtime types; the return annotation (File/Dir/Tuple(...)) drives encoding; the wire carries bare key strings.
  • one-marker-type → one type, Ref, a Path subclass that carries its own (storage, key, is_dir).

The user's constraint is "default to NO new type; any new runtime type must justify itself against a stdlib primitive." The dual-ergonomics lens is a hard filter, and it eliminates the two zero-type options on footgun grounds:

  • The _REFS global registry is a hidden aliasing contract — provenance lives in a process-global keyed by an ephemeral /tmp/... string, invisible to the reader, unserializable, and wrong the moment two artifacts collide on a scratch path. The lens explicitly bans "hidden ordering/aliasing contracts."
  • The bare-key-string wire is stringly-typed magic — the caller receives "process/poses.json", a plain str an agent will try to open() and get a FileNotFoundError, and a generic queue consumer can only interpret it if a schema fingerprint rides in the envelope. The lens explicitly bans "stringly-typed magic."

Ref is the minimal thing that carries wire identity without a hidden registry and without stringifying semantics. It passes the filter that the stdlib primitive fails:

python
class Ref(pathlib.Path):        # concrete PosixPath/WindowsPath subclass
    key: str                    # prefix-relative, e.g. "process/poses.json"
    storage: str | None         # storage binding; None => purely local
    is_dir: bool = False

Why it justifies itself against Path (the required justification): a plain Path cannot answer the one question the wire must answer — "do these bytes need to travel, and under what prefix-relative key?" A worker-local Path("/tmp/w7/poses.json") is meaningless to the caller; a returned str is indistinguishable from data. Ref answers it as instance state you can read (ref.key, ref.storage) rather than as a side-table. And because it is a Path, it is invisible where it counts: open(ref), np.load(ref), ref / "rgb" / f"{i}.png", os.fspath(ref), and every def f(x: str | Path) signature accept it unchanged.

Everything above one file stays stdlib primitives:

ShapeTypenot
one file / one directoryRef (is_dir=True for a dir)FileBundle
N named outputsdict[str, Ref]FileBundle.rgb
N positional outputstuple[Ref, ...]Tuple(...)
a stream of outputsgenerator / async-gen yielding RefBundles

So: one new type (Ref), everything else is dict / tuple / list / generator. The rejected FileBundle/FileRef/Bundles/Tuple(...) sketch is dropped in full.

The other load-bearing principles (all three agreed; adopted verbatim)

  1. msgpack the library stays; everything above it is fresh. A native ExtType-code codec replaces zaku's {"ztype":...} dict overlay and _greedy marker. zdata.py (the "Vendored from zaku" file) is deleted.
  2. One codec for args, kwargs, results, and stream frames. This deletes the two-@udf split and the msgpack.packb(result) result-path bug (arrays/tuples/Ref now survive as returns, not just as args).
  3. Bytes move at the seam, never in the codec. pack(Ref) emits (storage, key, is_dir) metadata only. Upload/download happens in the runner, the one place a Storage handle exists. This is the cleanest idea across the three and removes the codec's dependency on credentials entirely.
  4. Streaming is real: a per-invocation Redis-Streams journal with an explicit END/ERROR frame and a monotonic seq, replacing zaku's EOF-only termination and lakeshore's fake subscribe_stream (loop + break).

How the design serves dual ergonomics

  • Greppable, guessable verbs: submit / claim / ack, publish / subscribe, call / stream, gather / as_completed / map. An agent's first guess is the real method.
  • Ref is a Path, so there is no new API to learn to read a file — the agent uses open/np.load it already knows.
  • Introspectable returns: -> Ref, -> dict[str, Ref], -> AsyncIterator[Ref] are native hints an agent reads directly; no custom marker vocabulary to construct.
  • No footguns: no consume-once (a Ref is immutable and re-readable), no hidden registry, no stringly-typed keys, no _greedy mode to forget.
  • Errors state the fix: the codec raises wire: no codec for <Type> (register one via wire.register); the runner raises Ref "<key>" was returned but never created via dls.run.write(...).

2. Native TYPE / CODEC layer — dreamlake/lakeshore/wire.py

Replaces zdata.py entirely.

2.1 Leaf tagging — msgpack.ExtType codes via default / ext_hook

python
# wire.py
import msgpack
from dataclasses import dataclass
from typing import Any, Callable, Iterator

@dataclass(frozen=True)
class Codec:
    code: int                                # 0x01..0x7f reserved ext space
    name: str
    match:  Callable[[Any], bool]            # predicate on a live object
    encode: Callable[[Any], bytes]           # obj  -> ext-body bytes
    decode: Callable[[bytes], Any]           # bytes -> obj

_CODECS: list[Codec] = []                    # match-ordered
_BY_CODE: dict[int, Codec] = {}

def register(codec: Codec) -> None:
    _CODECS.append(codec); _BY_CODE[codec.code] = codec

def _default(obj: Any) -> msgpack.ExtType:
    for c in _CODECS:
        if c.match(obj):
            return msgpack.ExtType(c.code, c.encode(obj))
    raise TypeError(
        f"wire: no codec for {type(obj).__name__}; register one via wire.register(...)"
    )

def _ext_hook(code: int, data: bytes) -> Any:
    c = _BY_CODE.get(code)
    return c.decode(data) if c else msgpack.ExtType(code, data)   # forward-compat

def pack(value: Any) -> bytes:
    return msgpack.packb(value, default=_default, use_bin_type=True, strict_types=True)

def unpack(buf: bytes) -> Any:
    return msgpack.unpackb(buf, ext_hook=_ext_hook, raw=False)

strict_types=True is load-bearing: it forces tuple through _default (code 0x05) instead of silently coercing it to a list — so (rgb, depth) comes back a tuple.

2.2 Code table (stable integers; deferred imports preserved)

codePython typeext body (msgpack.packb)decode
0x01numpy.ndarray[dtype_str, shape, raw_bytes]frombuffer(...).reshape(...)
0x02torch.Tensor[dtype_str, shape, raw_bytes] (.detach().cpu().numpy())torch.from_numpy(buf.copy())
0x03PIL.Image.Image[format or "PNG", raw_bytes]Image.open(BytesIO(...))
0x04Refpack({"s": storage, "k": key, "d": is_dir}) — metadata only, no file bytesreconstruct Ref; runner resolves/downloads on first open
0x05tuplepack([*t])tuple(...)
0x06plain pathlib.Pathstr(p)Path(...); runner rejects a non-Ref local Path on a remote hop (see §6)

numpy/torch/PIL predicates try-import and return False on ImportError (exactly as zdata.py does today). Each is one register(Codec(...)) call, so a downstream package can add its own leaf type without editing wire.py.

2.3 Envelope — typed, versioned, no marker-in-data

python
# wire.py
@dataclass
class Envelope:
    v: int                       # wire version — lets the format evolve
    id: str                      # invocation_id (ULID)
    fn: dict                     # {module, qualname, version, source|pickle}
    prefix: str                  # item root -> becomes dls.run.prefix
    args: bytes                  # wire.pack(tuple(args))
    kwargs: bytes                # wire.pack(dict(kwargs))
    cfg: dict                    # {mode, timeout_s, resources, image, retry}
    parent: str | None = None

    def to_bytes(self) -> bytes: ...
    @classmethod
    def from_bytes(cls, b: bytes) -> "Envelope": ...

An explicit v field replaces zaku's _greedy boolean smuggled into user data. args/kwargs are a nested wire.pack so numpy/Ref arguments are codec-tagged; the body/meta split means a generic worker can read prefix/fn/cfg without decoding the payload.

2.4 Frame — one shape for scalar results, item streams, and chunked binary

python
# wire.py — Frame = packed [kind, seq, body, meta]
VALUE, CHUNK, BEGIN, END, ERROR = 0, 1, 2, 3, 4

@dataclass
class Frame:
    kind: int
    seq: int = 0                 # monotonic per invocation; the journal cursor
    body: Any = None             # codec value (VALUE) | raw bytes (CHUNK) | None
    meta: dict | None = None     # {size, content_type, name} on BEGIN; {state, error} on END

    def to_bytes(self) -> bytes:            # single pack of the whole array
        return pack([self.kind, self.seq, self.body, self.meta])
    @classmethod
    def from_obj(cls, o: list) -> "Frame":
        return cls(o[0], o[1], o[2], o[3])

class FrameReader:
    """Feed raw chunks; yield decoded Frames. msgpack length-prefixes make each
    frame self-delimiting over ONE Unpacker — no separator byte."""
    def __init__(self) -> None:
        self._u = msgpack.Unpacker(raw=False, ext_hook=_ext_hook)
    def feed(self, chunk: bytes) -> Iterator[Frame]:
        self._u.feed(chunk)
        for o in self._u:
            yield Frame.from_obj(o)

class FrameWriter:
    def value(self, obj, seq=0)     -> bytes: return Frame(VALUE, seq, obj).to_bytes()
    def chunk(self, data, seq)      -> bytes: return Frame(CHUNK, seq, data).to_bytes()
    def begin(self, meta, seq=0)    -> bytes: return Frame(BEGIN, seq, None, meta).to_bytes()
    def end(self, meta=None)        -> bytes: return Frame(END, -1, None, meta).to_bytes()
    def error(self, exc)            -> bytes: return Frame(ERROR, -1, None, {"error": repr(exc)}).to_bytes()

body is packed once as part of the frame array, so a VALUE frame carries a Ref/ndarray directly and ext_hook decodes it in place — no double-pack.


3. Native @udf — dreamlake/lakeshore/decorator.py

One decorator, four body kinds, detected once at decoration; args and results share wire.pack.

python
class UDF(Generic[T]):
    def __init__(self, fn, defaults, *, remote_by_default: bool):
        self.is_async     = inspect.iscoroutinefunction(fn)
        self.is_gen       = inspect.isgeneratorfunction(fn)
        self.is_async_gen = inspect.isasyncgenfunction(fn)
        self.ret          = typing.get_type_hints(fn).get("return")  # introspectable schema

    def __call__(self, *a, **k) -> Any:  ...   # dispatch per the matrix below
    def local(self, *a, **k) -> T: ...          # always in-process; still commits under a local prefix
    def submit(self, *a, **k) -> "Future[T]": ...        # non-blocking remote handle
    def stream(self, *a, **k): ...              # force Stream/AsyncStream for a generator body

Decorator surface (spellings kept — they are the right native words):

python
@udf                                     # local; sync | async | gen | async-gen
@udf("gpu")                              # remote-by-default
@udf(mode="local" if DEBUG else "gpu")
@udf("gpu", timeout_s=..., resources={...}, image=..., retry=2)
@udf.pure / @udf.pure("gpu")             # re-invocation-safe (memoizable in-runtime)

Dispatch matrix:

bodyfn(...) localfn(...) remote
def (sync)resultsubmit(...).result() (blocks)
async defcoroutineawait submit(...)
def … yield (gen)the generatorStream (sync-iterable)
async def … yield (async-gen)the async-genAsyncStream (async-iterable)

Explicit entrypoints: fn.local(...), fn.submit(...) -> Future, fn.stream(...) -> Stream | AsyncStream. Stream/AsyncStream are thin iterators over a FrameReader (§5) — handles, never data types; a consumer never constructs one.

The one result path (fixes the bug):

python
# runner (callee)  — entry.py / _runners.run_in_process
bound  = udf.resolve_inputs(args, kwargs, ctx)     # Ref args -> local paths (download)
raw    = fn(*bound.a, **bound.k)                    # or asyncio.run(...) for async
result = ctx.commit(raw)                            # Ref outputs -> upload; returns metadata Refs
client.settle(id, "succeeded", wire.pack(result))  # was: msgpack.packb(result)  <-- the bug

# Future (caller)
def _unpack(self):
    return wire.unpack(self._blob)                  # was: msgpack.unpackb(...)

The function reference is content-addressed by source (sha256(source)[:16], already computed for versioning) with cloudpickle kept only as an opt-in fallback for closures — because cloudpickle across a Python/deps mismatch is the classic remote-exec footgun and our args are already fully codec-able.


4. Native queues — dreamlake/lakeshore/queue.py

One Queue(name) carrying four faces; verbs renamed to name the actual job lifecycle. submit/poll/ack already exist in the control-plane client.py, so these align with the CP.

python
class Queue:
    def __init__(self, name: str, *, server=None, namespace="default", **pool_opts): ...

    # ── task queue (durable) ─────────────────────────────
    def submit(self, value: dict, *, key: str | None = None,
               priority: float = 0.0, meta: dict | None = None) -> str          # -> job_id
    def claim(self, *, timeout_s: float = 30.0) -> "Job | None"
    def ack(self, job: "Job | str", *, result: Any = None) -> None              # settle success
    def nack(self, job: "Job | str", *, error: str, retry: bool = True) -> None
    def count(self, status: str = "pending") -> int
    def __len__(self) -> int
    @contextmanager
    def working(self):                     # claim() + auto-ack / auto-nack on exception
        j = self.claim()
        try: yield j; self.ack(j)
        except BaseException as e: self.nack(j, error=repr(e)); raise

    # ── pubsub ───────────────────────────────────────────
    def publish(self, value: dict, *, topic: str) -> int                        # fan-out count
    def subscribe_one(self, topic: str, *, timeout: float = 5.0) -> dict | None
    def subscribe(self, topic: str, *, timeout: float | None = None) -> Iterator[dict]   # REAL stream
    async def subscribe_async(self, topic, *, timeout=None) -> AsyncIterator[dict]

    # ── RPC (submit + reply on the invocation_id topic) ──
    def call(self, *args, timeout: float = 30.0, **kwargs) -> Any               # request -> one reply
    def stream(self, *args, timeout: float | None = None, **kwargs) -> Iterator[Any]
    def call_future(self, *args, **kwargs) -> "Future"
    async def call_async(self, *args, timeout=30.0, **kwargs) -> Any
python
@dataclass(frozen=True)
class Job:
    id: str          # invocation_id == reply topic
    value: dict      # wire.unpack'd inline (S3 round-trip only for Ref leaves)
    attempt: int
    meta: dict

Gather + functional constructs (module-level; operate on Future/Stream):

python
dls.gather(*futs, timeout=None)          -> list                  # ordered join (submit order)
dls.as_completed(futs, timeout=None)     -> Iterator[Future]      # ready order; RE-RAISES failures
dls.map(fn, iterable, *, unordered=False)-> list                  # ordered parallel map
dls.starmap(fn, arg_tuples)              -> list
dls.scatter(fn, iterable)                -> list[Future]          # submit-all, no wait
# async siblings: dls.gather_async / dls.as_completed_async / dls.map_async

Two concrete fixes over today's code, carried from the proposals: submit(priority=) actually attaches priority to the envelope (today's PriorityQueue has no submit-side field); as_completed uses asyncio.wait semantics with no hardcoded 60 s fallback and no except BaseException: pass — a failed future re-raises.


5. Streaming + chunked streaming

Substrate — a per-invocation Redis-Streams journal in the Node control plane

Today publish is a terminal ack into a single result slot, so a "stream" physically cannot hold >1 message, and subscribe_stream fakes it with a break. The fix:

  • Producer: each frame → XADD ls:{inv}:out * f <frame_bytes>.
  • Terminal END (or ERROR) frame is the last entry.
  • Consumer: XREAD BLOCK <ms> STREAMS ls:{inv}:out <cursor>, feed bytes to a FrameReader, advance the cursor by entry id.

Replayable, survives consumer reconnect, and distinguishes producer finished from connection dropped — the two things zaku's write_eof() cannot.

python
# client.py additions (thin CP routes over the Redis Stream)
def journal_append(self, inv_id, frame: bytes) -> None                                  # XADD
def journal_read(self, inv_id, cursor="0", *, block_ms=5000) -> tuple[list[bytes], str, bool]
def journal_close(self, inv_id, state, error=None) -> None

Author side — a UDF emits N frames by being a generator

python
@dls.udf("gpu")
async def traj_to_ego_views(splats: str | Path, traj: str | Path) -> AsyncIterator[Ref]:
    steps = load_traj(dls.run.read(traj))
    for i, wp in enumerate(steps):
        p = dls.run.write(f"images/ind_{i:04d}.png")     # local scratch Ref, registered
        render(splats, wp).save(p)
        dls.run.progress(i + 1, len(steps))
        yield p                                          # -> one VALUE frame (a metadata Ref)

Each yield → FrameWriter().value(commit(p), seq=i) → journal_append. Exhaustion → END{state:"succeeded"}; an exception → ERROR.

Raw chunked binary (a single large value split rather than one giant msgpack bin) — yield bytes:

python
@dls.udf("gpu")
def stream_tar(path: str | Path) -> Iterator[bytes]:
    yield dls.begin(content_type="application/x-tar")     # optional BEGIN metadata frame
    with open(dls.run.read(path), "rb") as f:
        while (b := f.read(1 << 20)):
            yield b                                        # a yielded `bytes` -> a CHUNK frame

A yielded bytes is inherently a raw chunk; anything else is a VALUE. dls.begin(...)/dls.end(...) are helper functions that emit BEGIN/END metadata frames — not types, so the "one new type" budget holds. All frames ride the same single Unpacker, so "raw chunked binary over one msgpack.Unpacker" is just consecutive CHUNK frames the reader concatenates until an END.

Consumer side — iterate the call / the future

python
async for p in traj_to_ego_views.stream(splats, traj):   # p is a Ref per rendered frame
    ...
data = b"".join(stream_tar(path))                         # CHUNK run rejoined by the reader
fut = traj_to_ego_views.submit(splats, traj)
for p in fut:                                             # Future.__iter__ drains the journal
    ...

Future.__iter__/__aiter__ feed journal_read bytes into one FrameReader: VALUE → yield frame.body; a BEGIN…CHUNK…END run → yield joined bytes; END → stop; ERROR → raise the reconstructed exception.


6. Data-UDF returns with simple primitives — worked on the real functions

The seam — RunContext gets prefix + write/read; bytes move here, not in the codec

python
# context.py  (RunContext today: id, attempt, event_sink, parent_id, root_id)
@dataclass
class RunContext:
    ...
    prefix: str = ""                 # item root, e.g. "scenes/0007" (storage-relative)
    scratch: Path = field(default_factory=_mk_scratch)   # local working dir
    storage: str | None = None       # None => local mode, plain fs, no S3
    _writes: dict[str, Ref] = field(default_factory=dict)

class _RunProxy:                     # dls.run — the ENTIRE author-facing I/O surface
    @property
    def prefix(self) -> str: ...
    def write(self, key: str, *, dir: bool = False, eager: bool = False) -> Ref:
        """Declare an output at prefix/key. Returns a Ref (a Path) the body writes
        to with plain fs ops. Lazy(default): uploaded atomically on UDF success.
        eager=True: each file streams up as it lands (sequence UDFs)."""
    def read(self, ref_or_key: "Ref | str | Path") -> Ref:
        """Resolve an input to a local Ref, downloading prefix/key on first open
        if remote; identity in local mode."""
    def progress(self, i: int, total: int) -> None: ...

def scope(subpath: str):             # context manager pushing a sub-prefix (per-episode fan-out)
    ...
def ref(prefix: str, key: str, *, storage: str | None = None) -> Ref:   # reconstruct from (prefix,key)
    ...

dls.run.write/dls.run.read are the only place a Storage handle is reached — which is exactly why byte movement lives here and not in the codec (pack(Ref) is metadata-only). The Storage STS/presign brokering already in storage.py stays behind this seam; the body never sees a bucket, a credential, or an absolute path.

Local mode (@dls.udf with no mode, or mode="local"): storage=None, so write/read are plain filesystem ops under scratch — generate_scene runs end-to-end with no S3.

The commit rule (in the runner, after the body returns)

ctx.commit(value) walks the return tree: a Ref leaf uploads to storage[prefix/key] and re-emits as a metadata Ref; a Ref that arrived from upstream is re-shipped as the same (storage, key) (no re-download/re-upload); dict/tuple/list recurse; anything else is inline data. A Ref returned but never created via dls.run.write raises Ref "<key>" was returned but never created via dls.run.write(...).

Worked examples on the real UDFs

Single file — splats_to_mesh (also maps_to_keypoints, contact_sheets_to_segments, segment_sheets_to_labels). Lazy/atomic: a body that raises publishes nothing.

python
@dls.udf("gpu")
async def splats_to_mesh(splats: str | Path) -> Ref:
    mesh = dls.run.write("process/mesh.ply")        # was: out=f"{scene_root}/process/mesh.ply"
    extract_mesh(dls.run.read(splats)).save(mesh)
    return mesh

Directory as one artifact — mesh_to_maps (also video_to_contact_sheets, segments_to_segment_sheets). The is_dir=True Ref denotes the completed directory, so maps_to_keypoints(maps) never sees a half-populated maps/.

python
@dls.udf
async def mesh_to_maps(mesh: str | Path) -> Ref:
    d = dls.run.write("process/maps", dir=True)
    build_maps(dls.run.read(mesh), out_dir=d)
    return d

Named multi-output — poses_to_views → rgb/ + depth/. The name-carrying primitive is a plain dict[str, Ref]; downstream selects with views["rgb"].

python
@dls.udf("gpu")
async def poses_to_views(splats: str | Path, poses: str | Path) -> dict[str, Ref]:
    rgb   = dls.run.write("process/images/rgb",   dir=True, eager=True)
    depth = dls.run.write("process/images/depth", dir=True, eager=True)
    for pid, cam in enumerate(load(dls.run.read(poses))):
        render_rgb(splats, cam).save(rgb / f"{pid}.png")
        render_depth(splats, cam).save(depth / f"{pid}.png")
        dls.run.progress(pid, total=n)
    return {"rgb": rgb, "depth": depth}             # two stdlib primitives already say "N named files"

Sequence (streaming) — traj_to_ego_views returns AsyncIterator[Ref] (§5). Ordering lives in the key (ind_{i:04d}.png), recoverable from ref.key — so the row-alignment invariant lives in the key, not re-derived in two UDFs.

Upstream → downstream threading via _write/_prefix

Pipelines thread the returned Refs directly; the shared prefix makes refs resolvable across separately-decorated pipelines. build_scene_graph (today string-templates f"{scene_root}/process/images") becomes:

python
@dls.pipeline
async def build_scene_graph(scene_root: str) -> None:
    with dls.scope(scene_root):                          # prefix = scene_root
        views = await poses_to_views("source/splats.ply", "process/poses.json")
        await views_to_scene_graph(views["rgb"], "process/poses.json")   # select one named stream

Per-episode fan-out with a sub-prefix:

python
@dls.pipeline
async def build_traj_tuples(scene_root: str, *, num_episodes: int = 100) -> None:
    with dls.scope(scene_root):
        for i in range(num_episodes):
            with dls.scope(f"outputs/{i:04d}"):
                traj = await goal_to_traj("goal.json", "../../process/maps")
                await traj_to_ego_views("../../source/splats.ply", traj)   # eats the traj Ref

Properties that fall out, each mapped to a requirement:

  • Cross-pipeline resolution by (prefix, key) — a Ref is reconstructable from prefix+key alone (dls.ref(...)), so poses_to_views opens process/poses.json written by an earlier extract_freespace_keypoints run not in the current object graph. Also what makes migration incremental.
  • Fan-out / aliasing is free — a Ref is immutable and re-readable; source/splats.ply feeds splats_to_mesh, poses_to_views, and traj_to_ego_views with no consume-once semantics.
  • source/* pre-existing inputs are keys no UDF produced — addressed as dls.run.read("source/splats.ply") like any other key; the three-root source/·process/·outputs/ layout is just prefix-relative keys.
  • Idempotent re-runs — (prefix, key) is deterministic across attempt; a retry overwrites outputs/{i:04d}/traj.npz rather than forking.

Migration — no flag day

The commit walk checks, in order: (1) an explicit out= the pipeline still passes (today's convention) → the runner infers key = out.relative_to(prefix) and wraps the returned bare Path/str as a Ref; else (2) the Ref the body returned from dls.run.write. So the eight pipelines keep running while UDFs migrate one at a time (splats_to_mesh first, poses_to_views when its two named streams matter) — each migrates by changing out: str|Path → dls.run.write(...) and its annotation → -> Ref, independently.


7. Concept-by-concept comparison vs zaku (primary deliverable)

Conceptzaku's approachOur native approachVerdict + one-line why
Envelope / codec passPayload(SimpleNamespace), greedy ZData walk, _greedy marker injected into user datatyped Envelope{v,id,fn,prefix,args,kwargs,cfg}, codec always-on, default/ext_hook recursionRENAME + NEW — versionable header, no marker-in-data, no hand-walk (msgpack recurses for us)
Type tags{"ztype":"numpy.ndarray", b, dtype, shape} dict overlay; tuple→list; Path≡strmsgpack.ExtType integer codes 0x01–0x06; tuple + Path preservedRENAME — ext codes can't collide with a user "ztype" dict; loud on unknown; keeps tuple/Path identity
File / data returnsnone — S3 is an opaque queue transport; take returns one S3 ref for the whole payloadone Ref (Path subclass) + plain dict/tuple; bytes move at the dls.run seam, codec is metadata-onlyNEW (exactly one type) — files are the data-UDF's actual results, threaded per-leaf; stdlib already gives named/positional access
Result serialization(lakeshore) bare msgpack.packb(result) — drops ndarray/tuple/Refsame wire.pack as args/kwargs/framesDROP the bug — a result is just a value; one codec everywhere
Task-queue verbsadd / take(S3-ref) / mark_donesubmit / claim(inline Job) / ack / nackRENAME — names the real lifecycle; payload decoded inline, S3 only for Ref leaves
PubSubpublish = terminal ack (1 message max); subscribe_stream fakes it with breakpublish → Redis-Streams journal; subscribe/subscribe_async read by cursorKEEP-LESSON + NEW — real N-message fan-out, not a single-emit stub
RPCmints rpc-{uuid4()} topic carried in _request_idcall / stream / call_future, reply topic is the invocation_idRENAME — CP already indexes by invocation_id; no side-channel id field
Gatherdedicated {name}.return-queue.{uuid} + gather-{uuid} token set drained by is_donegather / as_completed / scatter / map over Futures reading the journal cursorRENAME + DROP — futures already carry completion; the second queue + token bookkeeping is redundant state
Streaming framingself-delimited msgpack over chunked HTTP body, terminated by write_eof(); no seqFrame[kind,seq,body,meta] over one Unpacker, Redis-Streams journal, explicit END/ERROR + seqKEEP-LESSON + NEW — keep self-delimited framing; add real close signal + resumable cursor + chunked-binary frames
@udf / clienttwo client classes (TaskQ sync, TaskQAsync async); function is an opaque payloadone @udf; detects sync/async/gen/async-gen; local/remote; content-addressed fnNEW — typed defs already declare their nature; no client class to pick
Data-UDF I/O seamnone — worker owns raw S3 keysdls.run.{prefix, write, read} + dls.scope + dls.refNEW — lets a body return a plain Ref/dict and never see a credential

8. Dual-ergonomics audit of the final public surface

SurfaceGreppable / guessable?Precise hints?Introspectable / self-doc?Footguns removedError quality
@udf / @udf("gpu")✅ first guess for "remote function"✅ preserves the wrapped signature✅ .ret, .is_async/.is_gen exposedone decorator, not two client classesmode "<x>" unknown; expected local|<pool>
fn.submit / .local / .stream✅ verbs an agent tries first✅ Future[T] / Stream[T]✅ handle types name themselves.local always in-process (no accidental remote)—
Ref✅ short, one concept✅ is a Path; every str|Path sig accepts it✅ .key, .storage, .is_dir readableno hidden _REFS global; no bare-string key; immutable → no consume-onceRef "<key>" returned but never created via dls.run.write(...)
dls.run.write / read / prefix✅ write/read are the obvious pair✅ write(key) -> Ref, read(x) -> Ref✅ docstring states lazy-vs-eager + local modebody never touches S3/creds/abs pathsread("<key>"): not found under prefix "<p>"
Queue.submit / claim / ack / nack✅ standard lifecycle verbs✅ claim() -> Job | None✅ Job.{id,value,attempt,meta}working() ctx-mgr auto-acks; priority= actually wirednack requires error= (keyword-only)
publish / subscribe✅✅ subscribe() -> Iterator[dict]✅ real stream, documentednot a 1-message stub; explicit END—
call / stream (RPC)✅✅ call(...) -> Any, stream(...) -> Iterator✅ reply keyed by invocation_idno rpc-{uuid} side-channel to leakcall timeout raises TimeoutError, not None
gather / as_completed / map / scatter✅ mirrors concurrent.futures✅ typed over Future[T]✅ ordered-by-default documentedas_completed re-raises (no silent except BaseException)failed future re-raises through the iterator
wire.pack / unpack / register✅✅ pack(Any) -> bytes✅ Codec dataclass is the schemaunknown type raises at pack, not deep in msgpackwire: no codec for <Type>; register one via wire.register(...)

Things an agent would likely get wrong, and how the design prevents them:

  1. Returning str(path) out of habit. A bare str is data, not a file → it would silently not ship. Prevented: the out= migration bridge wraps a returned bare path when an out= key is present, and the primary path (dls.run.write(...) -> Ref) makes the file-ness a type the agent sees in the signature. A Ref never created via write raises with the fix in the message.
  2. Trying to open() a value received from a streaming/named output. Because the value is a Ref (a Path), open/np.load just work — there is no bare key string to fumble.
  3. Assuming a stream is consume-once or that iterating twice re-reads. The journal has a cursor and an explicit END; a Future re-iterated replays from cursor 0. Refs are immutable and re-readable, so fan-out (source/splats.ply → three UDFs) is well-defined.
  4. Forgetting a serialization "mode" (zaku's _greedy). There is no mode — the codec is always on and free for scalar payloads (_default fires only on non-native types).
  5. Picking the wrong queue verb. submit/claim/ack are the CF-standard lifecycle words; add/take/mark_done are gone.

9. Open questions for the user

  1. Ref — accept the one-type exception? The synthesis introduces exactly one runtime type, Ref (a Path subclass), on the grounds that the two zero-type alternatives each trip the no-footgun filter (hidden global registry / stringly-typed wire). Do you accept Ref, or do you want to hold the line at literally zero and live with a bare-Path + resolve-registry (accepting the aliasing footgun)?
  2. Ref subclassing Path on Python ≤3.11. Direct subclassing of pathlib.Path is clean only on 3.12+. On ≤3.11 we need the _flavour = type(Path())._flavour trick. What is the floor Python version for lakeshore workers?
  3. Directory artifacts — upload granularity. For is_dir=True Refs, do we upload the subtree as many keys under a prefix (listable, resumable) or as one packed object (atomic, one round-trip)? Affects eager=True sequence UDFs and cross-pipeline directory reads.
  4. Return annotations — plain hints only, or optional marker aliases? The recommendation uses native hints (-> Ref, -> dict[str, Ref], -> AsyncIterator[Ref]). The annotation-driven proposal offered richer markers (File/Dir/Files) that also encode the key template for the dashboard/canvas schema. Do you want those as optional introspection-only aliases, or keep the surface strictly primitive?
  5. Content-addressed fn vs cloudpickle default. We propose shipping the fn by content-addressed source with cloudpickle as an opt-in closure fallback. Any UDFs today that rely on closures/partials that would force cloudpickle-by-default?
  6. out= bridge lifespan. How long must the legacy out=-passing call sites keep working — one release, or indefinitely? Determines whether the bridge is a deprecation shim or a permanent second path.
  7. Journal retention. Redis-Streams result journals need a trim/TTL policy (XADD MAXLEN or time-based). What retention do streamed results need for replay after a consumer reconnect?