Skip to content

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 makeWorker item on a container or the infra/ waitUntil helper 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 restart
yield* outbox.append(GenerateExportFile.make({ … })) // a command, not a fact

✅ Correct — the row and the fact, in one unit:

exports/request-export.ts
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>) plus Outbox.append of 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 CurrentOrg from event.orgId and a System principal 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

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>.ts and exports one subscriber value: the event schema and the handler. An event consumed only inside its slice lives in that slice’s events.ts; a domain-pair event’s schema lives on the subscriber’s side or in @app/domain.
  • infra/outbox/ and infra/jobs/ hold machinery only. main.ts exports subscribers beside jobs, and both entries (the container relay and the Worker’s drainOnce) take the same list.
  • No top-level jobs/, workers/, workflows/ or processes/ folder, and no package for definitions (coupling and cohesion).
  • Accepted cost: a subscriber missing from subscribers compiles, 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.optional or has a decoding default, unknown fields are ignored (the schema default), a growable enum is ForwardCompatibleNullable (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 zero
SELECT count(*) FROM outbox_events
WHERE event_tag = 'ExportRequested' AND processed_at IS NULL;

Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 4

Impact: MEDIUM no infinite loops, and no evidence thrown away

  • A payload SchemaError sets failed_at at once; ten failed attempts set it too. Failed rows are the dead-letter record and are never deleted automatically.
  • An unknown event_tag is 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 *Sync decoder. 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_at
UPDATE outbox_events SET failed_at = NULL, attempts = 0, available_at = <now> WHERE id = …;

Source: notes/08-application-surfaces/workflows-and-jobs.md · Decision 5

  • 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 a due_at.
  • Effect Workflow on ClusterWorkflowEngine as the process engine — trigger: an effect release where effect/workflow and effect/cluster are stable and the cluster’s SQL storage runs on D1.
  • alchemy’s Cloudflare.Workflow as 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.