Guides
Send email from a job
A request that sends email, calls an API, or generates a PDF makes the user wait for all of it — and fails them if any of it fails. Jobs fix both.
The problem
api.post(async (ctx) => { const user = await createUser(ctx.json()); // fast await mailer.send({ // slow, and can fail to: user.email, template: "welcome", data: { name: user.name, verifyUrl }, }); return API.json({ ok: true }); // user waited for all of it});Two failure modes: the user waits several seconds for an SMTP round trip, and if the mail server is down they get a 500 for something that actually succeeded.
Register the handler
Handlers live in yatta/func/workers.ts, which yatta/main.ts imports at boot.
import { jobs } from "./jobs";import { mailer } from "./mail"; jobs.handle("send-welcome", async ({ data, log }) => { log(`sending welcome to ${data.to}`); await mailer.send({ to: data.to, subject: "Welcome to Yatta!", text: "Thanks for creating an account.", }); return { sent: true };}); export const defaultWorker = jobs.worker("default", { concurrency: 5, pollInterval: "1s", lockDuration: "60s",});Note
concurrency is per worker process. Raising it above the database's comfortable write concurrency causes lock contention — start at 5 and raise deliberately.Enqueue from a route
import { API, createAPI } from "yatta/api";import { db } from "../func/db";import { jobs } from "../func/jobs"; const api = createAPI(); api.post(async (ctx) => { const { email, name } = await ctx.json(); const user = db.users.insert({ email, name }); // Returns immediately; the worker does the rest. await jobs.enqueue("send-welcome", { to: user.email, name: user.name, }); return API.json({ ok: true, user: user.id }, { status: 201 });}); export default api;Event-driven instead
If the same thing should happen in several places, pipe an event into the queue rather than enqueueing from each one.
events.on("user.registered", (data) => { console.log(`new user: ${data.email}`);}); // Event payload → job payload.events.pipe( "user.registered", "send-welcome", undefined, (data) => ({ to: data.email, name: data.name }),); // Emit from anywhere.await events.emit("user.registered", { userId: user.id, email: user.email, name: user.name,});Failures and retries
A throwing handler is retried with exponential backoff and jitter. After the attempt limit the job moves to the dead-letter queue.
await jobs.enqueue("send-welcome", data, { attempts: 5, timeout: "30s", retry: { type: "exponential", delay: 2000, factor: 2, jitter: true, // spreads retries out maxDelay: "10m", },});const dead = await jobs.dlq.list("default"); // Send one back through the queue after fixing the cause.await jobs.dlq.retry(dead[0].id); // Or clear them out.await jobs.dlq.purge("default"); const stats = await jobs.metrics("default");// { queued, delayed, running, completed, dead }Warning
Do not retry a permanent failure. If the address is invalid, retrying three times just delays the inevitable — validate before enqueueing.
Progress reporting
Long jobs can report progress, which is also how you broadcast state to the browser over the realtime connection.
jobs.handle("generate-report", async (ctx) => { await ctx.progress(10, "Fetching rows"); const rows = await fetchAllRows(); await ctx.progress(50, "Rendering PDF"); const pdf = await render(rows); await ctx.progress(100, "Done"); // Bail out if the job was cancelled or timed out. if (ctx.signal.aborted) return; return { url: upload(pdf) };});import { realtime } from "./func/realtime"; // Inside the handlerrealtime.job(ctx.id).progress(50, "Rendering"); // In the browserconst tracker = realtime.job(jobId);tracker.on("progress", (p) => { bar.style.width = `${p.percent}%`;});Scheduled work
cron.schedule("nightly-cleanup", "0 0 * * *", async () => { await jobs.enqueue("cleanup-stale-tokens");}); cron.every("10m", () => { console.log("[cron] heartbeat");});