DreamLake

Living dev note. Iterate freely.

Problem

Redis OOMs and goes slow when payload bytes live in the queue row. zaku hit this; the workaround in zaku-service was to move bytes to S3 and keep only {s3_key, s3_bucket, size} in Redis. We're doing the same — with one upgrade: clients upload/download S3 directly via server-issued presigned URLs, so the controlplane is never on the bandwidth path.

This is the cache for job arguments and return values. It has nothing to do with the log-sink wire (which streams stdout/stderr). Those two stores have different lifetimes, different access patterns, and different keys.

Two tiers

TierWhereWhenCost
InlineRedis stream entrysize ≤ INLINE_THRESHOLD (default 256 KB)Redis memory + a single round-trip
S3object store, client-direct via presigned URLsize > thresholdone extra HTTP for URL, one to S3

The threshold is a CP env var, not a queue-level knob — same value across the fleet so the wire shape is predictable.

Job row schema

ts
interface ExecJob {
  // existing
  id: string;
  namespaceId: string;
  queueId: string;
  // payload reference (mutually exclusive with `payloadInline`)
  payloadRef?: {
    kind: "s3";
    bucket: string;
    key: string;
    size: number;
    contentType?: string;
    sha256?: string;            // optional integrity check
  };
  payloadInline?: Uint8Array;   // only when size ≤ INLINE_THRESHOLD
  // return-value reference (filled by daemon after exec)
  resultRef?: { kind: "s3"; bucket; key; size; ... } | { kind: "inline"; bytes };
  // ...
}

Constraint: exactly one of payloadRef / payloadInline is set. Same for resultRef once the job completes.

Wire — enqueue (large payload)

client                                          CP                            S3
  │                                              │                             │
  │ POST /v1/.../get-started/payloads/presign?op=put&size=N  │                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { url, bucket, key, expires_at, fields } │                             │
  │                                                                            │
  │ PUT <url> (binary body, msgpack-encoded args)─────────────────────────────▶│
  │ ◀── 200 ─────────────────────────────────────────────────────────────────  │
  │                                              │                             │
  │ POST /v1/.../get-started/queues/:q/get-started/jobs                  │                             │
  │   { payloadRef: { bucket, key, size, sha256 } }                            │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { job_id }                               │                             │

CP validates that the key it's being told about matches one it issued (prevents clients from stuffing the queue with arbitrary S3 keys).

Wire — daemon pulls a job

daemon                                          CP                            S3
  │ POST /v1/daemon/poll                         │                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { job_id, payloadRef | payloadInline }   │                             │
  │                                                                            │
  │ (if payloadRef:)                                                           │
  │ POST /v1/.../get-started/payloads/presign?op=get&key=…   │                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { url, expires_at }                      │                             │
  │ GET <url> ────────────────────────────────────────────────────────────────▶│
  │ ◀── 200 (binary) ────────────────────────────────────────────────────────  │
  │                                                                            │
  │ ...runs the @udf...                                                        │

Wire — return value

Symmetric: daemon presigns a PUT for the result, uploads, then POSTs the resultRef to a per-job return channel. Client subscribes by job_id (zaku-style return topic) and downloads via a presigned GET.

daemon                                          CP                            S3
  │ POST /v1/.../get-started/payloads/presign?op=put         │                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { url, key }                             │                             │
  │ PUT <url> (binary result) ────────────────────────────────────────────────▶│
  │ POST /v1/.../get-started/jobs/:id/result                 │                             │
  │   { resultRef: { bucket, key, size, sha256 } }                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── 204                                                                    │

client                                          CP                            S3
  │ GET /v1/.../get-started/jobs/:id/result (long-poll or SSE)                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { resultRef: { ... } }                   │                             │
  │ POST /v1/.../get-started/payloads/presign?op=get&key=…   │                             │
  │ ────────────────────────────────────────────▶│                             │
  │ ◀── { url }                                  │                             │
  │ GET <url> ────────────────────────────────────────────────────────────────▶│
  │ ◀── 200 (binary)                                                           │

Encoding

Same as zaku — msgpack with the ZData extension for numpy/torch/PIL. Bytes only; no base64 anywhere. contentType: application/msgpack on the S3 object.

Bucket layout

Shared bucket with prefix:

<bucket>/
├── payloads/<YYYY>/<MM>/<DD>/<job_id>/args          # request payload
├── payloads/<YYYY>/<MM>/<DD>/<job_id>/result        # return value

Lifecycle rules:

  • payloads/... — expire after 48 hours (job-result TTL). Tunable per namespace later.
  • The log-sink bucket prefix is separate (exec/...) and has its own retention. Don't conflate.

Credentials

CP holds the master AWS creds (Heroku config var or IAM role on EC2 in the future). Presigned URLs are issued from those.

  • Default TTL on a presigned URL: 15 minutes (covers slow uploaders
    • clock skew).
  • Per-URL size cap: server-enforced via Content-Length header on the PUT. CP refuses to presign anything > MAX_PAYLOAD_BYTES (default 100 MB; raise as needed).
  • Auth: the payloads/presign endpoints sit behind the same namespace bearer-token gate as the rest of /v1/.

The client never sees AWS creds. The daemon never sees AWS creds. Only the URL.

Inline cutover

For payloads ≤ INLINE_THRESHOLD (256 KB default):

  • Client serializes with msgpack.
  • POSTs { payloadInline: <base64> } directly to the enqueue route.
  • No S3 round-trip.
  • CP stores the bytes inline in the Redis stream entry for that job.

This keeps the common "hello world" flow snappy — no S3 hop for tiny RPCs.

Failure modes

FailureBehaviour
S3 PUT times out client-sideclient retries the PUT (key is idempotent); CP doesn't see the job until the final enqueue POST
Daemon fetches but S3 returns 404job moves to payload_lost state; CP requeues with --retry-once or fails the job depending on policy
sha256 mismatch on downloaddaemon refuses to run, marks job as payload_corrupt; alerts
Presigned URL expires mid-uploadclient requests a new one (CP allows N presigns per job_id)
Bucket misconfig / region drifthealth-check endpoint surfaces it before job dispatch

What we're NOT building (yet)

  • Cross-region replication
  • Client-side encryption (S3 SSE-S3 only for now; SSE-KMS later)
  • Streaming-multipart upload for > 5 GB payloads (use single-PUT cap)
  • Deduplication by content-hash (each job_id gets a fresh key)
  • Web upload via the dashboard (CLI/SDK only)

Reference

zaku-service/zaku/s3_helper.py + zaku/interfaces.py:Job.add — same data shape (s3_key, s3_bucket, size), different wire: zaku-service proxies through the server; we go client-direct via presigned URLs.

Rollout

  1. Schema additions (Job row gets payloadRef / payloadInline / resultRef).
  2. POST /v1/.../get-started/payloads/presign route. Tests for size cap + auth + signing.
  3. Client SDK encode + presigned-upload helper. Wire to Queue.add.
  4. Daemon-side presigned-download helper. Wire to job execution.
  5. INLINE_THRESHOLD switch on the enqueue route — small payloads bypass S3.
  6. Return-value path (resultRef + per-job result subscribe).
  7. Cleanup: drop any code path that puts payload bytes into Redis directly.

Estimated ~5 days end-to-end. Lands after the queue table (phase A) and before the elasticity controller.