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

| Shape | Type | not |
|---|---|---|
| one file / one directory | `Ref` (`is_dir=True` for a dir) | ~~FileBundle~~ |
| N **named** outputs | `dict[str, Ref]` | ~~FileBundle.rgb~~ |
| N **positional** outputs | `tuple[Ref, ...]` | ~~Tuple(...)~~ |
| a **stream** of outputs | generator / async-gen yielding `Ref` | ~~Bundles~~ |

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

| 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

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

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

```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 `Ref`s 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)

| 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 `Future`s 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 ; 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. `Ref`s 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?
