Reliability
Three failure modes that normally end in a hung request or a dead thread: a worker that dies, a task that never returns, and a process that is asked to stop.
Worker supervision
When a worker emits error or close, the scheduler does three things:
- Rejects every in-flight request and every queued task with
WorkerCrashError— so callers fail fast instead of hanging forever - Spawns a replacement worker with an identical configuration
- Re-mounts every subsystem graph the dead worker was hosting, then re-points routing at the replacement
import { WorkerCrashError } from "yatta/runtime"; try { await runtime.execute(graphId, "computeReceipt", payload);} catch (err) { if (err instanceof WorkerCrashError) { console.log(`${err.workerId} (${err.role}) died: ${err.message}`); // The work is lost, but nothing is left dangling. }}error is always followed by close, and a crashed flag ensures replacement happens exactly once. A replacement that itself dies is replaced again — tested across cascading crashes.Task deadlines
Every dispatched task carries a deadline. It starts at dispatch, so time spent queued counts against it. When it expires the task is removed from the scheduler's bookkeeping and rejects with TaskTimeoutError.
// Scheduler-wide defaultconst runtime = createRuntime({ taskTimeoutMs: 30_000 }); // Per subsystemdefineSubsystem({ name: "exports", entrypoint: "./func/exports.ts", workload: "cpu", timeoutMs: 120_000, // a long report gets more time}); // Opt out entirelydefineSubsystem({ name: "worker", ..., timeoutMs: 0 });import { TaskTimeoutError } from "yatta/runtime"; try { await runtime.execute(graphId, "export", payload);} catch (err) { if (err instanceof TaskTimeoutError) { console.log(`${err.handler} exceeded ${err.timeoutMs}ms`); }}Graceful shutdown
On SIGINT or SIGTERM the scaffolded server:
- Stops accepting new connections
- Waits up to ten seconds for in flight tasks to finish
- Only then terminates the worker threads
const server = Bun.serve({ /* … */ }); let shuttingDown = false;const shutdown = async () => { if (shuttingDown) return; shuttingDown = true; server.stop(true); console.log("[shutdown] draining…"); const drained = await runtime.drain(10_000); if (!drained) { console.warn(`${runtime.getActiveTaskCount()} task(s) still in flight — forcing.`); } await runtime.shutdown(); process.exit(0);}; process.on("SIGINT", () => void shutdown());process.on("SIGTERM", () => void shutdown());The re-entrancy guard matters: a second signal arriving while draining must not start a second teardown.
Draining manually
// Resolves true if the fleet went idle, false on timeoutconst clean = await runtime.drain(5_000); console.log(runtime.getActiveTaskCount()); // in flight + queuedError types
class WorkerCrashError extends Error { workerId: string; role: "cpu-pool" | "io-pool" | "dedicated";} class TaskTimeoutError extends Error { graphId: string; handler: string; timeoutMs: number;}Multi-process notes
Under SO_REUSEPORT several OS processes share the port and the database. SQLite is opened in WAL mode with a busy timeout, and schema creation runs under BEGIN IMMEDIATE with retries, so concurrent boots are safe.