Examples

Image pipeline

The request validates and stores the original, then gets out of the way. A worker resizes, an event fires, and progress reaches the browser over SSE.

Uploading and processing an image are different problems. Holding the request open while three resize passes run is what makes an upload feel broken. Split them: the request does only what has to happen before the bytes are safe, everything else happens on a worker.

The shape

tsts
POST /images  ├─ validate (magic bytes, size, type)      ~50ms  ├─ store the original                     ~200ms  ├─ enqueue "image:resize"                 instant  └─ 202 { id, statusUrl } worker runs "image:resize"  ├─ thumb.webp / card.webp / hero.jpg  ├─ progress(33) → progress(66) → progress(100)  └─ emit "image.ready" GET /images/:id/progress   (SSE)  └─ pushes every progress tick as it happens

Storage setup

yatta/func/storage.tsts
import { createStorage } from "yatta/storage"; export const storage = createStorage({  driver: process.env.S3_BUCKET    ? {        type: "s3",        bucket: process.env.S3_BUCKET,        region: process.env.AWS_REGION!,        accessKeyId: process.env.AWS_ACCESS_KEY_ID!,        secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY!,      }    : {        type: "local",        path: "./uploads",      },});
Note
The same code runs against local disk in development and S3 in production. Nothing above the driver knows the difference.

Schema

yatta/func/db.tsts
export const schema = {  images: {    id: col.uuid(),    userId: col.text().references("users.id", { onDelete: "CASCADE" }),    originalKey: col.text(),    status: col.text().default("pending"),   // pending | ready | failed    error: col.text().nullable(),    meta: col.json<Record<string, unknown>>().nullable(),    createdAt: col.createdAt(),    updatedAt: col.updatedAt(),  },}; export const db = createDatabase({ url: process.env.DATABASE_URL, schema });

Typed events

Declare the events once so emit and on are both checked. Getting this wrong used to silently drop the payload — the signature takes a transform function as its fourth argument.

yatta/func/events.tsts
import { createEvents } from "yatta/jobs";import { jobs } from "./jobs"; export interface AppEvents {  "image.uploaded":  { imageId: string; userId: string };  "image.ready":    { imageId: string; variants: string[] };  "image.failed":   { imageId: string; error: string };} declare module "yatta/jobs" {  interface EventRegister extends AppEvents {}} export const events = createEvents(jobs);

The upload route

yatta/backend/images.tsts
import { createAPI, HttpError } from "yatta/api";import { verifyMagicBytes, guessContentType } from "yatta/storage";import { and } from "yatta/db";import { db } from "../func/db";import { storage } from "../func/storage";import { jobs } from "../func/jobs";import { requireAuth } from "../func/middleware"; const images = createAPI("/images"); images.use(requireAuth); images.post("/", async (ctx) => {  const userId = (ctx.state.user as { id: string }).id;  const form = await ctx.req.formData();   const file = form.get("file");  if (!(file instanceof File)) throw new HttpError(400, "file is required");   // Cap the request body before reading it into memory.  const MAX = 10 * 1024 * 1024;  if (file.size > MAX) throw new HttpError(413, "File exceeds 10MB");   // Trust the bytes, not the declared Content-Type: a .exe renamed to .jpg  // arrives here with a perfectly valid image/jpeg header.  const head = new Uint8Array(await file.slice(0, 16).arrayBuffer());  if (!verifyMagicBytes(head)) throw new HttpError(415, "Not an image");   const contentType = guessContentType(head) ?? file.type;  if (!contentType.startsWith("image/")) {    throw new HttpError(415, "Not an image");  }   const id = crypto.randomUUID();  const key = `images/${userId}/${id}/original`;   await storage.put(key).from(file.stream() as any)    .withContentType(contentType)    .asPrivate()    .maxSize("10mb")    .verifyMagic();   const row = await db.images.insert({    id,    userId,    originalKey: key,    status: "pending",  });   // Fire and forget. The response does not wait for a resize.  await jobs.enqueue("image:resize", { imageId: id }, { priority: "high" });   return images.json(    { id: row.id, status: row.status, progressUrl: `/images/${id}/progress` },    { status: 202 },  );});

The worker

yatta/func/jobs.tsts
import { createJobs } from "yatta/jobs";import { db } from "./db";import { storage } from "./storage";import { events } from "./events"; export const jobs = createJobs(); const VARIANTS = [  { suffix: "thumb", width: 200,  format: "webp" },  { suffix: "card",  width: 800,  format: "webp" },  { suffix: "hero",  width: 2000, format: "jpeg" },] as const; jobs.handle("image:resize", async ({ imageId }, ctx) => {  const image = await db.images.findById(imageId);  if (!image) throw new Error(`image ${imageId} vanished`);   try {    const original = await storage.disk().download(image.originalKey);    const buffer = await original.buffer();     const variants: string[] = [];     for (const [i, v] of VARIANTS.entries()) {      // Progress is per-step, so the bar moves instead of sitting at zero.      ctx.progress(        Math.round(((i + 1) / VARIANTS.length) * 100),        `Generating ${v.suffix}`,      );       const resized = await resize(buffer, v);      const key = `images/${image.userId}/${imageId}/${v.suffix}.${v.format}`;       await storage.put(key)        .from(resized)        .withContentType(`image/${v.format}`)        .asPublic()        .maxSize("2mb");       variants.push(key);    }     await db.images      .where((f) => f.id.isEqualTo(imageId))      .update({ status: "ready", updatedAt: new Date() });     // The UI is already listening for this; no polling needed.    await events.emit("image.ready", { imageId, variants });     ctx.progress(100, "Done");  } catch (err) {    await db.images      .where((f) => f.id.isEqualTo(imageId))      .update({ status: "failed", error: String(err) });     await events.emit("image.failed", {      imageId,      error: err instanceof Error ? err.message : "Unknown",    });     // Rethrow so the retry policy applies — a transient disk error should not    // be permanent.    throw err;  }});
Warning
Emitting the failure event and rethrowing is deliberate. The event updates the UI now; the rethrow hands the job to the retry policy. Swallowing the error would leave the row stuck as pending with no retries and no visible error.

Streaming progress over SSE

realtime.job(id) follows a job from the worker. The browser subscribes over SSE and sees each tick as it happens — no polling, and it survives a worker restart because the tracker is keyed by job id.

yatta/backend/images.tsts
import { SSEClient } from "yatta/realtime";import { realtime } from "../func/realtime"; images.get("/:id/progress", async (ctx) => {  const userId = (ctx.state.user as { id: string }).id;   // Ownership first — never let progress be a way to probe for ids.  const image = await db.images    .where((f) => and(      f.id.isEqualTo(ctx.params.id),      f.userId.isEqualTo(userId),    ))    .first();   if (!image) throw new HttpError(404, "Not found");   // Already finished? Send the terminal state and close.  if (image.status !== "pending") {    const client = new SSEClient(ctx.req);    client.send("done", { status: image.status, error: image.error });    client.close();     return new Response(client.readable, {      headers: { "Content-Type": "text/event-stream" },    });  }   const client = new SSEClient(ctx.req);  const tracker = realtime.job(image.id);   // Subscribe to the live tracker.  for await (const update of tracker) {    client.send("progress", {      percent: update.percent,      message: update.message,    });     if (update.percent >= 100) break;  }   client.send("done", { status: "ready" });  client.close();   return new Response(client.readable, {    headers: {      "Content-Type": "text/event-stream",      "Cache-Control": "no-store",      "Connection": "keep-alive",    },  });});

Serving a variant

StorageFile.serve understands range requests, so a large hero image can be partially re-fetched instead of re-downloaded.

yatta/backend/images.tsts
images.get("/:id/:variant", async (ctx) => {  const userId = (ctx.state.user as { id: string }).id;   const image = await db.images    .where((f) => and(      f.id.isEqualTo(ctx.params.id),      f.userId.isEqualTo(userId),    ))    .first();   if (!image || image.status !== "ready") {    throw new HttpError(404, "Not found");  }   const key = `images/${userId}/${image.id}/${ctx.params.variant}`;  const file = await storage.disk().download(key);   // Handles Range, ETag and conditional GET for you.  return file.serve(ctx.req, {    contentType: ctx.params.variant.endsWith(".webp")      ? "image/webp"      : "image/jpeg",  });});

The browser

public/app.jsts
const form = new FormData();form.append("file", input.files[0]); const res = await fetch("/images", { method: "POST", body: form });const { id, progressUrl } = await res.json(); // The server told us where to watch.const events = new EventSource(progressUrl); events.addEventListener("progress", (e) => {  const { percent, message } = JSON.parse(e.data);  bar.style.width = `${percent}%`;  status.textContent = message;}); events.addEventListener("done", (e) => {  const { status, error } = JSON.parse(e.data);  events.close();  status.textContent = status === "ready" ? "Uploaded" : `Failed: ${error}`;});
Tip
EventSource reconnects on its own, so a dropped connection recovers without retry logic on your side. Close it on done or it will keep reconnecting to a finished stream.