DreamLake

Living dev note. Iterate freely.

Where this fits

The tiered pass-by-handle design below is the SDK-side surface of a broader caching strategy. The same machinery powers:

  • Avoiding re-serialization across UDF hops (this doc's primary motivation — handing a 500 MB tensor from tokenize to embed without round-tripping bytes).
  • Caching expensive computations by content-hash. A @udf(cache=True) function that returns a payloadRef whose key is a hash of its args + fn bytes; a subsequent call with the same args resolves from cache without re-running. Tier 1 (S3) gives durable cache; tier 2 (worker-resident) gives hot cache.
  • Cross-pipeline data reuse. When two pipelines share an intermediate (loaded dataset, tokenizer state, embedding table), the handle is the dedup primitive.

So the tiers below are the substrate — pass-by-handle is the first consumer; explicit caching is the second; cross-pipeline reuse is the third. The wire shape and lifecycle are the same.

See also: /dev/payload-store for the storage substrate this tiered system rides on.

Motivation

Today every @udf call serializes args + return on every hop. That's correct and decoupled (producer/worker can run different Python runtimes — see /dev/python-sdk for why we're OK with it) but it's expensive when the result of one @udf is the argument to the next:

python
big = fetch_dataset()          # → 500 MB numpy
toks = tokenize(big)           # serialize 500 MB across the wire
emb = embed(toks)              # serialize again

We pay the serialize + S3-upload + S3-download cost twice when the data never actually leaves the worker fleet's data plane.

Goal

A first-class "by reference" mode where a value produced by one @udf is addressable by handle on the next, without re-serializing the bytes.

Three plausible backings, in increasing order of ambition.

Tier 1 — S3 reference passthrough (cheapest)

fn returns a value larger than INLINE_THRESHOLD → wire shape is resultRef = { bucket, key, size } (already exists via Phase B). The next @udf taking that as input can declare pass_by_ref=True on its parameter, and the SDK forwards the payloadRef to the worker without producer-side download. Net: zero serialize/download on the producer, one S3 fetch on the next worker (which is cheap if the worker is in the same region as the bucket).

This is the easy first cut. It composes cleanly with the existing payload-store wire — no new primitives.

Tier 2 — Worker-resident handles

A value lives on a specific worker's RAM (or local SSD). The result handle from @udf carries a (worker_id, local_key) reference. The next @udf requesting that handle is scheduled on the same worker if possible, then resolves the handle from local state. Falls back to S3 round-trip if the worker has aged out or the next call lands elsewhere.

This is Ray's ObjectRef in spirit. Requires:

  • Scheduler-aware dispatch ("prefer worker X if alive, fall back to S3 ref")
  • Worker-side TTL + memory-pressure eviction
  • A blob server endpoint per worker for cross-worker fetches

Tier 3 — Inter-worker zero-copy (most work)

For the in-region multi-worker case: workers can fetch from each other's blob endpoint without an S3 hop. Adds:

  • Cross-worker network bandwidth (often free or cheap intra-VPC)
  • Authentication between workers
  • A small distributed cache protocol

Non-goals (for v1)

  • True shared-memory across hosts (Ray Plasma). Out of scope — network is fine.
  • Avoiding cloudpickle entirely. The fn definition still cloudpickles; only the arg/return values are by-reference.
  • ObjectRef-style implicit chaining. Caller still has to pass the handle explicitly.

Open questions

  • Producer awareness — does the producer side need to know it's getting a handle vs. a value? Probably yes for tier 2/3, no for tier 1 (the handle is just a thin wrapper around payloadRef).
  • Serialization protocol negotiation — pairs with cloudpickle protocol-negotiation on /hello. Worker advertises which handle tiers it supports; producer picks the highest mutual.
  • GC + TTL — handles are mortal. Need explicit handle.persist() for things that should outlive the producing worker's session.
  • Failure semantics — if the worker holding tier-2 state dies, the handle becomes a S3 ref (fall through) OR an error (HandleLost)? Probably the former, with a metric for the fallback rate.

Phasing

  1. Tier 1 first — small change on top of Phase B. ~3 days.
  2. Tier 2 — Phase G after scheduler-aware dispatch lands.
  3. Tier 3 — Phase H. Probably never if Tier 1+2 cover real workloads.

Bookmark — no code lands until Tier 1 is scoped against a real workload.