DreamLake

Lakeshore

Lakeshore is the platform brand. Three pieces ship under it, each with its own deploy story:

PieceWhatWhere it lives
lakeshore CLINode + TypeScript binary; registers providers, manages daemons, runs smoke tests.npm: @dreamlake/lakeshore (lakeshore bin). Source lives in the lakeshore submodule.
Control planeNode + Fastify + Prisma + Mongo HTTP service. Owns canonical state and dispatches Invocations to Workers.Deployed on Heroku from the lakeshore-controlplane submodule.
nymph daemonRust binary that supervises Invocations on each compute host under docker / gvisor / process runners. Outbound-only (long-polls the control plane).The nymph submodule, closed-source for now; build from source.

The rest of this page is about the control plane specifically — the single brain producers submit to, workers long-poll, and the dashboard reads from. One stateless service, one durable store behind it.

It owns:

  • the canonical state of every Namespace, Mode, Queue, Function, Invocation, Worker, and Token
  • the dispatch decisions that match Invocations to Workers
  • the event log a Queue's subscribers tail
  • the scaling decisions that grow and shrink Pools

It speaks the daemon protocol to Workers, the producer API to the Python library, and exposes a small admin surface (HTTP + WebSocket) for the dashboard.

Stack

LayerChoiceWhy
RuntimeNode.js (latest LTS) + TypeScriptFirst-class concurrency for long-poll fan-in; ecosystem fit for the dashboard side.
ORMPrismaGenerated, type-safe client; migrations and schema diffing as a first-class concept.
DatabaseMongoDB 7+Durable record of every Namespace, Mode, Queue, Function, Invocation, Worker. Embedded sub-docs (runConfig, tags, resources), atomic findAndModify. Prisma's mongodb provider.
Hot pathIn-memory queues, in the Lakeshore processPer-queue pending arrays, processing index, and event subscribers live in Node.js memory. No Redis dependency. The Mongo writes are write-through (durable backstop).
Wiremsgpack (binary) + JSON (admin)msgpack for daemon and producer traffic; JSON for the dashboard / CLI.
Object storeS3 / R2Args, results, stdout/stderr that exceed wire.max_payload (default 16 MiB) spill here. Producer and worker handle S3 directly; Lakeshore only stores references.

The split: Mongo for what must survive a restart, in-memory for the hot path, S3 for anything big. No Redis. No bespoke cache layer.

On startup Lakeshore loads all queued Invocations from Mongo into its in-memory pending lists. On graceful shutdown it drains in-flight work and flushes terminal acks. The trade-off is honest: a crash between an in-memory state change and the corresponding Mongo write loses at most one transition for one Invocation, which the daemon's ack flow recovers on the next reconnect.

API surfaces

Lakeshore exposes three surfaces on the same port, distinguished by path prefix.

/v1/daemon/* — worker protocol (msgpack)

The set documented in daemon protocol: /hello, /poll (long-poll), /ack, /event. Daemons connect here and stay connected.

/v1/producer/* — library protocol (msgpack)

Used by the Python library. Mostly thin wrappers over Prisma writes.

MethodPathPurpose
PUT/v1/producer/get-started/queuesUpsert a queue by name. Body carries mode + elasticity bounds.
POST/v1/producer/submitSubmit one Invocation envelope. Auto-creates the queue if missing.
POST/v1/producer/awaitLong-poll for an Invocation's completion. Backs q.result(inv).

The producer-side API is what a queue-bound fn(...) / fn.submit(...) and q.submit(...) / q.result(...) ultimately call.

/v1/admin/* — dashboard + CLI (JSON)

REST-shaped CRUD over the same entities, but in JSON for human-tool ergonomics. Read-mostly; writes are gated by admin tokens.

Data model

Full Prisma schema lives at prisma/schema.prisma in the lakeshore-controlplane repo (currently private). Excerpt:

schema.prisma (excerpt)prisma
datasource db {
  provider = "mongodb"
  url      = env("DATABASE_URL")
}

model Namespace {
  id        String  @id @default(auto()) @map("_id") @db.ObjectId
  name      String  @unique
  modes     Mode[]
  queues    Queue[]
  workers   Worker[]
}

model Mode {
  id          String    @id @default(auto()) @map("_id") @db.ObjectId
  namespaceId String    @db.ObjectId
  name        String
  backend     String    // "local" | "fabric"
  runner      String    // "docker" | "gvisor" | "process"
  image       String?
  tags        String[]
  resources   Json      // embedded sub-doc
  env         Json
  timeoutS    Int?
  @@unique([namespaceId, name])
}

model Queue {
  id            String   @id @default(auto()) @map("_id") @db.ObjectId
  namespaceId   String   @db.ObjectId
  name          String
  modeId        String   @db.ObjectId
  state         String
  minHosts      Int
  maxHosts      Int
  policy        String
  @@unique([namespaceId, name])
}

model Invocation {
  // ULID from the producer — we override _id so the doc is keyed by
  // the same id that travels on the wire.
  id            String   @id @map("_id")
  queueId       String   @db.ObjectId
  functionId    String   @db.ObjectId
  state         String
  workerId      String?  @db.ObjectId
  argsBlob      Bytes?       // inline msgpack
  argsRef       String?      // S3 key if spilled
  resultBlob    Bytes?
  resultRef     String?
  error         Json?
  runConfig     Json
  priority      Int      @default(0)
  attempt       Int      @default(0)
  submittedAt   DateTime
  startedAt     DateTime?
  finishedAt    DateTime?
  // Compound index drives the dispatch hot path:
  @@index([queueId, state, priority, submittedAt])
}

Mongo-native choices:

  • runConfig, resources, env, tags, capabilities are embedded sub-documents / arrays. No join cost on read.
  • Invocation._id is the producer's ULID, not an ObjectId — the id you see in the daemon protocol, the invocation id you hold in the SDK, and Mongo are the same string.
  • TTL on Event.expiresAt is configured outside Prisma — db.Event.createIndex({ expiresAt: 1 }, { expireAfterSeconds: 0 }).

Hot paths

Two interactions dominate throughput.

Submit

POST /v1/producer/submit is the call your fn(...), fn.submit(...), or q.submit(...) ultimately makes. Steps (see src/server.ts in the lakeshore-controlplane repo):

  1. Resolve the Queue id by (namespaceId, name) — ensureQueueByName upserts so a never-before-seen queue auto-creates with defaults.
  2. prisma.invocation.create — durable Mongo insert with state: "queued".
  3. bumpAdmin() — wakes any admin long-poll watching the change feed.

Dispatch (the poll reply)

When a daemon's long-poll arrives, the dispatcher runs one findAndModify against Mongo to claim a queued Invocation atomically:

ts
// lakeshore-controlplane/src/server.ts — claimOne, paraphrased
const result = await prisma.$runCommandRaw({
  findAndModify: "Invocation",
  query: { state: "queued", queueId: { $in: queueIds } },
  sort: { priority: -1, submittedAt: 1 },  // priority desc, FIFO ties
  update: { $set: { state: "running", workerId, startedAt: <now> } },
  new: true,
});

For prefix-subscribed daemons (--queue 'train/*'), resolveQueueIds does a name: { startsWith: prefix } lookup and feeds the matching set into claimOne's $in filter — fairness is priority-then-FIFO across the whole prefix, not round-robin between queues.

In-memory hot path — planned, not built

The architectural direction is to move pending queue state into a per-queue in-memory pending array inside the Lakeshore process so poll becomes an O(1) shift instead of a Mongo findAndModify, with Mongo demoted to a write-through durable record. Today both submit and dispatch hit Mongo directly. The scaling discussion below already assumes the in-memory plan; until that lands, throughput is bounded by Mongo write IOPS rather than Node event-loop turnover.

Scaling

The in-memory queue state means a single Lakeshore process owns its queues. Two consequences:

  • One process per queue family. Multiple Lakeshore instances cannot share an in-memory queue. To scale horizontally, partition by queue name (e.g. one Lakeshore for prod/*, another for experiments/*) and have the load balancer route by queue name.
  • Vertical first. Node.js on a single modern core comfortably handles tens of thousands of submits/sec when the hot path is in-memory; the bottleneck moves to Mongo write throughput for the durable transitions, which a replica set absorbs well.

Scale-out plan, applied in order as needed:

  1. Bigger box. More CPU / memory for the Lakeshore process. The in-memory queues use bounded memory proportional to pending depth, which the queue's max_invocations cap controls.
  2. Partition by namespace or queue prefix. Run one Lakeshore per partition; load balancer routes producers and daemons by queue name to the owning instance.
  3. Shard Mongo by namespaceId. When storage or durable-write IOPS saturate. The schema's compound indexes are namespace-prefixed, so this is a clean cut.

No tiered queue-server hierarchy, no Redis. If we ever need to share hot-path state across processes, that's the moment to revisit — but for ML-shape workloads the steps above are plenty.

What lives where

ConcernWhere
Identity, tags, Mode catalogMongo (Namespace, Mode, Token)
Per-Queue config and boundsMongo (Queue)
Invocation durable recordMongo (Invocation) — every state transition written through
Invocation hot pathIn-memory pending arrays + processing index inside the Lakeshore process
Args / result blobs > 16 KiBS3 (referenced from argsRef / resultRef)
stdout / stderr beyond headS3
Live event fan-out to subscribersIn-memory subscriber map; long-poll resolution wakes them
Durable event log for replayMongo (Event, TTL'd)
Worker livenessIn-memory heartbeat timestamps; Mongo (Worker.lastSeenAt) on transition
Dispatch claimIn-memory pop, then Mongo write-through

Read next