Workflows & background jobs
Durable work is a row in our own database, on every runtime profile. This page answers which
mechanism carries which kind of work, how a multi-step process is modelled, where its
definitions live, and what happens to a row that never succeeds. Vocabulary: a
subscriber handles one outbox event; a process is a multi-step state machine on a slice’s
row; a job is recurring only; workflow is reserved for an engine-run workflow, and none
exists today. Running example: an org asks for a data export, the server generates a file,
emails a link, and expires the link after 7 days, all in a slice called exports/.
Make durable work an outbox row; let loseable work go to the helper
Section titled “Make durable work an outbox row; let loseable work go to the helper”Impact: HIGH nothing promised to a user disappears on a deploy
- The test: “if this is lost on a deploy, is any state wrong, or any promise to a user broken?”
If yes, it is an outbox event with a subscriber
(idempotency and outbox). If no, it goes to a
makeWorkeritem on a container or theinfra/waitUntilhelper on a Worker; their loss is counted, not prevented. - The row records a fact the producer committed (
ExportRequested), never a command (GenerateExportFile), even when producer and subscriber are the same slice. - No Effect cluster, no workflow engine and no platform queue by default. The container profile has no queue, and a second mechanism would be a second retry system for the same kind of row.
- Accepted cost: throughput is the relay’s (10 rows per claim, 30s each). On a Worker the relay runs after a response and on a 1-minute cron, so a burst waits its turn.
| the work | Node / Bun | Cloudflare Worker |
|---|---|---|
| may be lost | a makeWorker item |
the infra/ waitUntil helper |
| must not be lost, one step | an outbox event and its subscriber | the same rows, drained in waitUntil(drainOnce) and a 1-minute cron |
| several steps or a wait | a process (below) | a process (below) |
| recurring | Jobs.layer |
one Cron Trigger per job |
❌ Incorrect — a forked fiber dies with the deploy; the event names the consumer’s action:
yield* exports.insert(org.orgId, row)yield* Effect.forkDetach(generateFile(row.id)) // lost on restartyield* outbox.append(GenerateExportFile.make({ … })) // a command, not a fact✅ Correct — the row and the fact, in one unit:
yield* transactions.withTransaction(Effect.gen(function* () { yield* exports.insert(org.orgId, row) // status: requested yield* outbox.append(ExportRequested.make({ id: yield* newId(OutboxEventId), exportId: row.id, orgId: org.orgId, }))}))Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 1
Model a process as a status column plus chained events
Section titled “Model a process as a status column plus chained events”Impact: HIGH one SELECT answers “where is export X?”, on every profile
- The state is a status column on the owning slice’s row. Each transition is one atomic unit: a
guarded update (
WHERE org_id = … AND id = … AND status = <from>) plusOutbox.appendof the next event. If the guard matches no row, the step already ran and the subscriber returns success. - A subscriber is one atomic unit plus at most one outbound call, idempotent, and fits the relay’s
30s timeout. Long work is split into chunks with a cursor column. The relay provides
CurrentOrgfromevent.orgIdand aSystemprincipal per row. - A step that calls out does so before its transition commits, under an idempotency key derived from our id. A step that fails for good moves the row to its failure state and appends a failure event; undoing earlier steps is that event’s subscribers.
- Accepted cost: transitions are written by hand, with no step memoization. On D1 a guard does not stop the next statement in a batch, so a replayed step can append a duplicate next event; the next step’s guard absorbs it.
❌ Incorrect — no guard, no org, and the outbound write after the commit:
(e: ExportRequested) => Effect.gen(function* () { yield* transactions.withTransaction(Effect.gen(function* () { yield* sql`UPDATE exports_exports SET status = 'generated' WHERE id = ${e.exportId}` yield* outbox.append(next) // appended again on every replay })) yield* files.put(`exports/${e.exportId}`, yield* render(e)) // a crash here loses the file})✅ Correct — the call first, keyed by our id; then a guarded transition and the next event:
(e: ExportRequested) => Effect.gen(function* () { const current = yield* exports.findById(e.orgId, e.exportId) if (Option.isNone(current) || current.value.status !== "requested") return const fileKey = yield* files.put(`exports/${e.exportId}`, yield* render(e)) // idempotent by key const next = ExportGenerated.make({ id: yield* newId(OutboxEventId), exportId: e.exportId, orgId: e.orgId }) yield* transactions.withTransaction(Effect.gen(function* () { if (yield* exports.markGenerated(e.orgId, e.exportId, fileKey)) yield* outbox.append(next) }))})Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 2, amended
Make a wait a column read by a job
Section titled “Make a wait a column read by a job”Impact: MEDIUM a pending wait is visible in a SELECT and cancelled by changing the row
- A delay that belongs to the domain is a
due_at-style column, read by a job as of its slot (scheduling and retry). Never a sleeping fiber or a delayed message. - An external callback or a human decision is an HTTP or webhook handler that performs a guarded
transition (
WHERE status = 'awaiting_…') in its own unit. - A process whose stalling a user would notice has a deadline column and a job that moves overdue rows to the failure state. The job enumerates orgs first, then works per org.
- Accepted cost: a timer has the granularity of its job: minutes, not seconds.
❌ Incorrect — the wait lives in memory, invisible and uncancellable:
yield* Effect.sleep(Duration.days(7))yield* exports.expire(e.orgId, e.exportId)✅ Correct — the column and a job:
export const expireExports = Jobs.make({ name: "exports.expire-links", cron: "*/15 * * * *", exclusive: false, run: (slot) => Effect.gen(function* () { const orgIds = yield* exports.dueOrgIdsAcrossOrgs(slot, "expire export links that are due") yield* Effect.forEach(orgIds, (orgId) => exports.expireDueAt(orgId, slot)) }),})Source: notes/08-application-surfaces/workflows-and-jobs.md · Decisions 1, 2, amended
Define the work in the slice that owns the state
Section titled “Define the work in the slice that owns the state”Impact: MEDIUM each piece of work sits next to the state it changes
- A subscriber file is
<slice>/on-<event-in-kebab-case>.tsand exports one subscriber value: the event schema and the handler. An event consumed only inside its slice lives in that slice’sevents.ts; a domain-pair event’s schema lives on the subscriber’s side or in@app/domain. infra/outbox/andinfra/jobs/hold machinery only.main.tsexportssubscribersbesidejobs, and both entries (the container relay and the Worker’sdrainOnce) take the same list.- No top-level
jobs/,workers/,workflows/orprocesses/folder, and no package for definitions (coupling and cohesion). - Accepted cost: a subscriber missing from
subscriberscompiles, and its rows retry until they fail. The relay’s unknown-tag warning is how that shows up.
❌ Incorrect — background work pulled away from its state:
apps/server/src/ jobs/expire-exports.ts workers/generate-export.ts infra/outbox/handlers.ts # one switch over every kind of work✅ Correct — the process is read by reading the slice folder:
apps/server/src/ exports/ export.ts # row schema, ExportStatus events.ts # ExportRequested, ExportGenerated exports-repo.ts # the guarded transitions request-export.ts # use-case: insert + append, one unit on-export-requested.ts # subscriber, step 1 on-export-generated.ts # subscriber, step 2 expire-exports.ts # job infra/outbox/ infra/jobs/ # machinery only main.ts # exports `subscribers` and `jobs`Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 3
Read stored payloads tolerantly; break with a new tag
Section titled “Read stored payloads tolerantly; break with a new tag”Impact: HIGH a row written by the previous release still decodes
- A row written by the previous release may still be pending, so changes are additive: a new
field is
Schema.optionalor has a decoding default, unknown fields are ignored (the schema default), a growable enum isForwardCompatibleNullable(API evolution). No field is removed, renamed or narrowed in place. - A breaking change is a new tag (
ExportRequestedV2). The old subscriber stays until a count of unprocessed old-tag rows returns zero, failed rows requeued or deleted first, then goes in a later PR. - A subscriber ships in the same release as its first producer, or earlier.
- Accepted cost: during a breaking change, two subscribers live side by side.
❌ Incorrect — a version field, and a field renamed in place:
{ version: Schema.Literal(2), fileFormat: ExportFormat } // was `format`: every pending row breaks✅ Correct — an optional new field; a new tag when it must break:
format: Schema.optional(ExportFormat), // old rows still decode-- the old subscriber can go when this returns zeroSELECT count(*) FROM outbox_eventsWHERE event_tag = 'ExportRequested' AND processed_at IS NULL;Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 4
Fail at once only what can never succeed
Section titled “Fail at once only what can never succeed”Impact: MEDIUM no infinite loops, and no evidence thrown away
- A payload
SchemaErrorsetsfailed_atat once; ten failed attempts set it too. Failed rows are the dead-letter record and are never deleted automatically. - An unknown
event_tagis a retryable failure with the normal backoff and a warning naming the tag: the replica that knows it may be one deploy step away. Ten attempts span about 32 minutes, longer than a rolling deploy. - The relay decodes with
Schema.decodeUnknownEffect, never a*Syncdecoder. A defect is caught by the per-row boundary and counts as a failed attempt, counted at claim. A domain outcome (the org was deleted) is not poison: the subscriber returns success with an Info line. - Requeue is one statement an operator runs, written down next to the failed-row alert, not built as a command.
❌ Incorrect — a sync decoder throws outside the per-row boundary, so the row loops forever:
const payload = Schema.decodeUnknownSync(subscriber.schema)(row.payload)✅ Correct — the Effect decoder, and a written-down requeue:
const payload = yield* Schema.decodeUnknownEffect(subscriber.schema)(row.payload) // SchemaError → failed_atUPDATE outbox_events SET failed_at = NULL, attempts = 0, available_at = <now> WHERE id = …;Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 5
Deferred
Section titled “Deferred”- A Cloudflare Queue between the relay and its subscribers on the Worker — trigger: the oldest unprocessed outbox row is older than 5 minutes for more than 1 hour in total within a week.
Outbox.append(event, { availableAt })for one-off delayed delivery — trigger: the first delay that has no domain row to carry adue_at.- Effect
WorkflowonClusterWorkflowEngineas the process engine — trigger: an effect release whereeffect/workflowandeffect/clusterare stable and the cluster’s SQL storage runs on D1. - alchemy’s
Cloudflare.Workflowas the process engine — trigger: the container profile is retired, leaving only the Worker profile. - Running a long step outside the relay — trigger: a subscriber’s p95 run time exceeds 20s over a week.
- A requeue command or admin endpoint — trigger: more than 5 manual requeues in one month.