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

the slow versionts
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.

yatta/func/workers.tsts
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

yatta/backend/signup.tsts
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.

yatta/func/events.tsts
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.

control the policyts
await jobs.enqueue("send-welcome", data, {  attempts: 5,  timeout: "30s",  retry: {    type: "exponential",    delay: 2000,    factor: 2,    jitter: true,     // spreads retries out    maxDelay: "10m",  },});
inspect failurests
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.

progressts
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) };});
broadcast it livets
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

yatta/func/cron.tsts
cron.schedule("nightly-cleanup", "0 0 * * *", async () => {  await jobs.enqueue("cleanup-stale-tokens");}); cron.every("10m", () => {  console.log("[cron] heartbeat");});