Scheduling & retry
Retry only what is safe to replay, inside a budget the caller can afford. This page answers where retry policies live, which failures are retried, and how a recurring job is defined and kept from running twice when several replicas are up.
Name every retry policy in infra/retry.ts and wrap exactly the replayed unit
Section titled “Name every retry policy in infra/retry.ts and wrap exactly the replayed unit”Impact: HIGH copies drift; a retry in the wrong place replays committed work
- Every retry schedule and retry predicate is a named export of
infra/retry.ts. A call site names one. A feature that needs a new budget adds an export there, in the PR that needs it. - The retry sits in the pipeable of the
Effect.fnthat is the atomic unit: the use-case’swithTransaction, or a repository method that owns a D1batch. A use-case that does more than its unit gives the unit its ownEffect.fn, so a retry never re-runs steps after the commit. - A repository method never retries. It cannot know whether it runs inside a transaction, and a statement retried inside an aborted Postgres transaction fails again.
- An HTTP adapter builds its retrying client once, with
HttpClient.retryTransient({ schedule: Retry.transient }), and uses it only forGET,PUT,DELETEand aPOSTthat carries an idempotency key. Everything else uses the base client.
❌ Incorrect — an inline policy, inside the repository:
activate: (orgId: OrgId, ruleId: RuleId) => sql`UPDATE labeling_rules SET active = 1 WHERE org_id = ${orgId} AND id = ${ruleId}`.pipe( Effect.retry({ schedule: Schedule.exponential("100 millis"), times: 3 }), )✅ Correct — the atomic unit is its own Effect.fn, retried with a named policy:
const persist = Effect.fn("ActivateRule.persist")( function* (ruleId: RuleId) { yield* transactions.withTransaction(Effect.gen(function* () { yield* rules.activate(orgId, ruleId) yield* outbox.append(RuleActivated.make({ id: yield* newId(OutboxEventId), ruleId, orgId })) })) }, Effect.retry({ schedule: Retry.contention, while: Retry.isContention }),)Source: notes/08-application-surfaces/scheduling-and-retry.md · Decision 1, amended
Retry contention always, connection loss only on an idempotent unit
Section titled “Retry contention always, connection loss only on an idempotent unit”Impact: HIGH a lost COMMIT acknowledgement can mean “already applied”
- Contention (
DeadlockError,SerializationError,LockTimeoutError, including SQLite’sSQLITE_BUSY) is retried on any atomic unit. The database rolled it back, so a replay cannot apply anything twice. ConnectionErroris retried withRetry.transientonly on an idempotent unit: a read, a set-state write, a create by client id withON CONFLICT, or a write guarded by a receipt (see idempotency and the outbox).- Never retried:
StatementTimeoutError, any D1 failure (D1 classifies nothing), a domain error,SchemaError, auth failures and other 4xx, defects and interruptions. isRetryableanswers “could a later attempt succeed?”, not “is a replay safe?”. Within a request, aRetry-Afterlonger than the remaining budget is not waited for; the adapter fails with its typed error.
❌ Incorrect — any “retryable” failure replayed on a unit that appends:
const persist = Effect.fn("ActivateRule.persist")( function* (ruleId: RuleId) { /* activate + outbox.append in one transaction */ }, Effect.retry({ schedule: Retry.transient }), // lost COMMIT ack → event appended twice)✅ Correct — the predicate says what is safe:
Effect.retry({ schedule: Retry.contention, while: Retry.isContention }) // any atomic unitEffect.retry({ schedule: Retry.transient, while: Retry.isConnectionLost }) // idempotent unit onlySource: notes/08-application-surfaces/scheduling-and-retry.md · Decision 2
Cap every policy, jitter it, and keep one retry layer per failure
Section titled “Cap every policy, jitter it, and keep one retry layer per failure”Impact: MEDIUM unbounded or stacked retries multiply load
- Every in-process policy is
capped(base, cap, retries): exponential, a per-gap cap, an attempt cap and jitter. The attempt cap is mandatory, and jitter is mandatory because two or more replicas hit the same contention at once. - Two policies:
Retry.contention(50ms base, 200ms cap, 3 retries, about 0.35s) andRetry.transient(200ms base, 2s cap, 3 retries, about 1.4s). Anything longer is the relay’s job, or the next job slot’s. - Work that already runs under a durable retrier (an outbox subscriber, a Queue consumer) retries
in process on
Retry.contentiononly. The durable retrier owns everything else. - A durable delay is computed from the attempt count stored in the row (
Retry.durableDelay), never from aSchedule, whose state resets on restart. A retry attempt logs nothing; the final failure is logged once.
❌ Incorrect — uncapped gaps, no jitter, stacked under a queue that retries again:
// inside a Queue consumer configured with maxRetries: 8Effect.retry({ schedule: Schedule.exponential("100 millis"), times: 2 })✅ Correct — one shape, every policy ends:
const capped = (base: Duration.Input, cap: Duration.Input, retries: number) => Schedule.max([ Schedule.min([Schedule.exponential(base), Schedule.spaced(cap)]), Schedule.recurs(retries), ]).pipe(Schedule.jittered)
export const contention = capped("50 millis", "200 millis", 3)export const transient = capped("200 millis", "2 seconds", 3)Source: notes/08-application-surfaces/scheduling-and-retry.md · Decision 3
Define a recurring job once in its slice, keyed by its slot
Section titled “Define a recurring job once in its slice, keyed by its slot”Impact: HIGH one schedule on both runtime profiles, and missed slots heal themselves
- A job is a value exported by its slice, built with
Jobs.make: name, cron,exclusive, a timeout andrun(slot). On the container profilemain.tspasses the list toJobs.layer; on the Worker profileisolate.tsregisters the same list withCloudflare.Workers.cron. Jobs.run(job, slot)is the boundary on both profiles: timeout, lease if exclusive, a span, catch and log once, and interruption let through. A failed run is retried by the next slot.- A job processes what is due as of its slot (
WHERE due_at <= ${slot}), never “since the last run”. Cron is parsed in UTC,nowcomes fromClock, and the granularity is a minute at most. - Jobs run in the server process. There is no external cron and no job bin.
❌ Incorrect — schedule buried in a layer, with an in-memory “since the last run” cursor:
Effect.gen(function* () { const since = yield* Ref.getAndSet(lastRun, yield* Clock.currentTimeMillis) yield* trials.closeExpiredSince(since)}).pipe(Effect.repeat(Schedule.cron("0 3 * * *")), Effect.forkScoped)✅ Correct — the slice owns what, when, and whether it is exclusive:
export const closeExpiredTrials = Jobs.make({ name: "billing.close-expired-trials", cron: "0 3 * * *", // UTC on both profiles exclusive: true, timeout: "5 minutes", // the lease is timeout + 1 minute run: (slot) => trials.closeDueAt(slot), // what is due as of slot})Source: notes/08-application-surfaces/scheduling-and-retry.md · Decision 4
Claim an exclusive job’s slot with one lease upsert
Section titled “Claim an exclusive job’s slot with one lease upsert”Impact: HIGH two replicas must not both send the email
- Ask in review: “if two replicas run this slot at once, is anything wrong or paid for twice?” No
means
exclusive: falseand no lease (the outbox relay, bulk deletes). Yes meansexclusive: true(email, unkeyed provider calls, expensive scans, billing). infra/jobs/ownsjobs_leases, one row per job. The claim is one upsert that runs unchanged on Postgres, SQLite and D1.last_slot < slotstops a retried or late fire;lease_until < nowstops overlap; the far-future clause reclaims a lease written by a clock that jumped.- The claim takes
nowas a parameter. The lease is the job’s timeout plus 1 minute. The release is guarded byholder, a run id. A replica that loses the claim does nothing and logs nothing. - The lease cuts duplicates; it does not make them impossible. An exclusive job is still idempotent per item.
❌ Incorrect — a lease with no slot: a retried fire runs the slot again:
UPDATE jobs_leases SET lease_until = ${now + leaseMs}, holder = ${runId}WHERE name = ${job.name} AND lease_until < ${now}✅ Correct — one row back means this replica runs the slot:
INSERT INTO jobs_leases (name, last_slot, lease_until, holder)VALUES (${job.name}, ${slot}, ${now + leaseMs}, ${runId})ON CONFLICT (name) DO UPDATE SET last_slot = excluded.last_slot, lease_until = excluded.lease_until, holder = excluded.holder WHERE jobs_leases.last_slot < excluded.last_slot AND (jobs_leases.lease_until IS NULL OR jobs_leases.lease_until < ${now} OR jobs_leases.lease_until > ${now + 2 * leaseMs})RETURNING nameSource: notes/08-application-surfaces/scheduling-and-retry.md · Decision 5, amended
Back off a poll loop whose tick fails
Section titled “Back off a poll loop whose tick fails”Impact: LOW a database outage logs once a minute, not once a second
- A sub-minute poll loop whose tick fails waits before the next tick. The wait doubles from its interval up to 1 minute and resets on the first successful tick. Today that is only the container relay’s 1s idle poll.
- The tick boundary is unchanged: catch, log once, continue (see error boundaries).
❌ Incorrect — a fixed interval through an outage:
tick fails → log → 1s → fails → log → 1s → … (one error line per second per replica)✅ Correct — doubling on failure, reset on success:
tick fails → log → 2s → fails → log → 4s → … → 60s (cap) → succeeds → back to 1sSource: notes/08-application-surfaces/scheduling-and-retry.md · Decision 6
Deferred
Section titled “Deferred”- In-process retry of D1 failures — trigger: an
@effect/sql-d1release whoseclassifyErrorreturns a reason other thanUnknownError. infra/jobs/(Jobs.make,Jobs.run, the loop, the Worker registration) — trigger: the first recurring job.- The
jobs_leasestable, claim and release — trigger: the first job declaredexclusive: true. - Lease renewal (a heartbeat) — trigger: the first exclusive job that needs a timeout above 15 minutes.
- Full or decorrelated jitter — trigger: an incident traced to synchronized retries from several replicas.
- Honoring an upstream
Retry-Afterin the relay — trigger: the first adapter whose upstream error carries a retry-after delay.