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
Pathis 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, aPathsubclass 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
_REFSglobal 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 plainstran agent will try toopen()and get aFileNotFoundError, 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:
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:
| Shape | Type | not |
|---|---|---|
| one file / one directory | Ref (is_dir=True for a dir) | |
| N named outputs | dict[str, Ref] | |
| N positional outputs | tuple[Ref, ...] | |
| a stream of outputs | generator / async-gen yielding Ref |
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)
- msgpack the library stays; everything above it is fresh. A native
ExtType-code codec replaces zaku's{"ztype":...}dict overlay and_greedymarker.zdata.py(the "Vendored from zaku" file) is deleted. - One codec for args, kwargs, results, and stream frames. This deletes the two-
@udfsplit and themsgpack.packb(result)result-path bug (arrays/tuples/Refnow survive as returns, not just as args). - 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 aStoragehandle exists. This is the cleanest idea across the three and removes the codec's dependency on credentials entirely. - Streaming is real: a per-invocation Redis-Streams journal with an explicit
END/ERRORframe and a monotonicseq, replacing zaku's EOF-only termination and lakeshore's fakesubscribe_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. Refis aPath, so there is no new API to learn to read a file — the agent usesopen/np.loadit 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
Refis immutable and re-readable), no hidden registry, no stringly-typed keys, no_greedymode to forget. - Errors state the fix: the codec raises
wire: no codec for <Type> (register one via wire.register); the runner raisesRef "<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
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)
| code | Python type | ext body (msgpack.packb) | decode |
|---|---|---|---|
0x01 | numpy.ndarray | [dtype_str, shape, raw_bytes] | frombuffer(...).reshape(...) |
0x02 | torch.Tensor | [dtype_str, shape, raw_bytes] (.detach().cpu().numpy()) | torch.from_numpy(buf.copy()) |
0x03 | PIL.Image.Image | [format or "PNG", raw_bytes] | Image.open(BytesIO(...)) |
0x04 | Ref | pack({"s": storage, "k": key, "d": is_dir}) — metadata only, no file bytes | reconstruct Ref; runner resolves/downloads on first open |
0x05 | tuple | pack([*t]) | tuple(...) |
0x06 | plain pathlib.Path | str(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
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
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.
Decorator surface (spellings kept — they are the right native words):
Dispatch matrix:
| body | fn(...) local | fn(...) remote |
|---|---|---|
def (sync) | result | submit(...).result() (blocks) |
async def | coroutine | await submit(...) |
def … yield (gen) | the generator | Stream (sync-iterable) |
async def … yield (async-gen) | the async-gen | AsyncStream (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):
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.
Gather + functional constructs (module-level; operate on Future/Stream):
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(orERROR) frame is the last entry. - Consumer:
XREAD BLOCK <ms> STREAMS ls:{inv}:out <cursor>, feed bytes to aFrameReader, 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.
Author side — a UDF emits N frames by being a generator
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:
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
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
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.
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/.
Named multi-output — poses_to_views → rgb/ + depth/. The name-carrying primitive is a plain dict[str, Ref]; downstream selects with views["rgb"].
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:
Per-episode fan-out with a sub-prefix:
Properties that fall out, each mapped to a requirement:
- Cross-pipeline resolution by
(prefix, key)— aRefis reconstructable fromprefix+keyalone (dls.ref(...)), soposes_to_viewsopensprocess/poses.jsonwritten by an earlierextract_freespace_keypointsrun not in the current object graph. Also what makes migration incremental. - Fan-out / aliasing is free — a
Refis immutable and re-readable;source/splats.plyfeedssplats_to_mesh,poses_to_views, andtraj_to_ego_viewswith no consume-once semantics. source/*pre-existing inputs are keys no UDF produced — addressed asdls.run.read("source/splats.ply")like any other key; the three-rootsource/·process/·outputs/layout is just prefix-relative keys.- Idempotent re-runs —
(prefix, key)is deterministic acrossattempt; a retry overwritesoutputs/{i:04d}/traj.npzrather 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)
| Concept | zaku's approach | Our native approach | Verdict + one-line why |
|---|---|---|---|
| Envelope / codec pass | Payload(SimpleNamespace), greedy ZData walk, _greedy marker injected into user data | typed Envelope{v,id,fn,prefix,args,kwargs,cfg}, codec always-on, default/ext_hook recursion | RENAME + 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≡str | msgpack.ExtType integer codes 0x01–0x06; tuple + Path preserved | RENAME — ext codes can't collide with a user "ztype" dict; loud on unknown; keeps tuple/Path identity |
| File / data returns | none — S3 is an opaque queue transport; take returns one S3 ref for the whole payload | one Ref (Path subclass) + plain dict/tuple; bytes move at the dls.run seam, codec is metadata-only | NEW (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/Ref | same wire.pack as args/kwargs/frames | DROP the bug — a result is just a value; one codec everywhere |
| Task-queue verbs | add / take(S3-ref) / mark_done | submit / claim(inline Job) / ack / nack | RENAME — names the real lifecycle; payload decoded inline, S3 only for Ref leaves |
| PubSub | publish = terminal ack (1 message max); subscribe_stream fakes it with break | publish → Redis-Streams journal; subscribe/subscribe_async read by cursor | KEEP-LESSON + NEW — real N-message fan-out, not a single-emit stub |
| RPC | mints rpc-{uuid4()} topic carried in _request_id | call / stream / call_future, reply topic is the invocation_id | RENAME — CP already indexes by invocation_id; no side-channel id field |
| Gather | dedicated {name}.return-queue.{uuid} + gather-{uuid} token set drained by is_done | gather / as_completed / scatter / map over Futures reading the journal cursor | RENAME + DROP — futures already carry completion; the second queue + token bookkeeping is redundant state |
| Streaming framing | self-delimited msgpack over chunked HTTP body, terminated by write_eof(); no seq | Frame[kind,seq,body,meta] over one Unpacker, Redis-Streams journal, explicit END/ERROR + seq | KEEP-LESSON + NEW — keep self-delimited framing; add real close signal + resumable cursor + chunked-binary frames |
@udf / client | two client classes (TaskQ sync, TaskQAsync async); function is an opaque payload | one @udf; detects sync/async/gen/async-gen; local/remote; content-addressed fn | NEW — typed defs already declare their nature; no client class to pick |
| Data-UDF I/O seam | none — worker owns raw S3 keys | dls.run.{prefix, write, read} + dls.scope + dls.ref | NEW — lets a body return a plain Ref/dict and never see a credential |
8. Dual-ergonomics audit of the final public surface
| Surface | Greppable / guessable? | Precise hints? | Introspectable / self-doc? | Footguns removed | Error quality |
|---|---|---|---|---|---|
@udf / @udf("gpu") | ✅ first guess for "remote function" | ✅ preserves the wrapped signature | ✅ .ret, .is_async/.is_gen exposed | one decorator, not two client classes | mode "<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 readable | no hidden _REFS global; no bare-string key; immutable → no consume-once | Ref "<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 mode | body never touches S3/creds/abs paths | read("<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 wired | nack requires error= (keyword-only) |
publish / subscribe | ✅ | ✅ subscribe() -> Iterator[dict] | ✅ real stream, documented | not a 1-message stub; explicit END | — |
call / stream (RPC) | ✅ | ✅ call(...) -> Any, stream(...) -> Iterator | ✅ reply keyed by invocation_id | no rpc-{uuid} side-channel to leak | call timeout raises TimeoutError, not None |
gather / as_completed / map / scatter | ✅ mirrors concurrent.futures | ✅ typed over Future[T] | ✅ ordered-by-default documented | as_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 schema | unknown type raises at pack, not deep in msgpack | wire: no codec for <Type>; register one via wire.register(...) |
Things an agent would likely get wrong, and how the design prevents them:
- Returning
str(path)out of habit. A barestris data, not a file → it would silently not ship. Prevented: theout=migration bridge wraps a returned bare path when anout=key is present, and the primary path (dls.run.write(...) -> Ref) makes the file-ness a type the agent sees in the signature. ARefnever created viawriteraises with the fix in the message. - Trying to
open()a value received from a streaming/named output. Because the value is aRef(aPath),open/np.loadjust work — there is no bare key string to fumble. - Assuming a stream is consume-once or that iterating twice re-reads. The journal has a cursor and an explicit
END; aFuturere-iterated replays from cursor 0.Refs are immutable and re-readable, so fan-out (source/splats.ply→ three UDFs) is well-defined. - Forgetting a serialization "mode" (zaku's
_greedy). There is no mode — the codec is always on and free for scalar payloads (_defaultfires only on non-native types). - Picking the wrong queue verb.
submit/claim/ackare the CF-standard lifecycle words;add/take/mark_doneare gone.
9. Open questions for the user
Ref— accept the one-type exception? The synthesis introduces exactly one runtime type,Ref(aPathsubclass), on the grounds that the two zero-type alternatives each trip the no-footgun filter (hidden global registry / stringly-typed wire). Do you acceptRef, or do you want to hold the line at literally zero and live with a bare-Path+ resolve-registry (accepting the aliasing footgun)?RefsubclassingPathon Python ≤3.11. Direct subclassing ofpathlib.Pathis clean only on 3.12+. On ≤3.11 we need the_flavour = type(Path())._flavourtrick. What is the floor Python version for lakeshore workers?- Directory artifacts — upload granularity. For
is_dir=TrueRefs, do we upload the subtree as many keys under a prefix (listable, resumable) or as one packed object (atomic, one round-trip)? Affectseager=Truesequence UDFs and cross-pipeline directory reads. - 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? - 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?
out=bridge lifespan. How long must the legacyout=-passing call sites keep working — one release, or indefinitely? Determines whether the bridge is a deprecation shim or a permanent second path.- Journal retention. Redis-Streams result journals need a trim/TTL policy (
XADD MAXLENor time-based). What retention do streamed results need for replay after a consumer reconnect?