Skip to content

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.fn that is the atomic unit: the use-case’s withTransaction, or a repository method that owns a D1 batch. A use-case that does more than its unit gives the unit its own Effect.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 for GET, PUT, DELETE and a POST that carries an idempotency key. Everything else uses the base client.

❌ Incorrect — an inline policy, inside the repository:

labeling/labeling-rules-repo.ts
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:

labeling/activate-rule.ts
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’s SQLITE_BUSY) is retried on any atomic unit. The database rolled it back, so a replay cannot apply anything twice.
  • ConnectionError is retried with Retry.transient only on an idempotent unit: a read, a set-state write, a create by client id with ON 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.
  • isRetryable answers “could a later attempt succeed?”, not “is a replay safe?”. Within a request, a Retry-After longer 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 unit
Effect.retry({ schedule: Retry.transient, while: Retry.isConnectionLost }) // idempotent unit only

Source: 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) and Retry.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.contention only. The durable retrier owns everything else.
  • A durable delay is computed from the attempt count stored in the row (Retry.durableDelay), never from a Schedule, 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: 8
Effect.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 and run(slot). On the container profile main.ts passes the list to Jobs.layer; on the Worker profile isolate.ts registers the same list with Cloudflare.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, now comes from Clock, 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:

billing/close-expired-trials.ts
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: false and no lease (the outbox relay, bulk deletes). Yes means exclusive: true (email, unkeyed provider calls, expensive scans, billing).
  • infra/jobs/ owns jobs_leases, one row per job. The claim is one upsert that runs unchanged on Postgres, SQLite and D1. last_slot < slot stops a retried or late fire; lease_until < now stops overlap; the far-future clause reclaims a lease written by a clock that jumped.
  • The claim takes now as a parameter. The lease is the job’s timeout plus 1 minute. The release is guarded by holder, 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 name

Source: notes/08-application-surfaces/scheduling-and-retry.md · Decision 5, amended

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 1s

Source: notes/08-application-surfaces/scheduling-and-retry.md · Decision 6

  • In-process retry of D1 failures — trigger: an @effect/sql-d1 release whose classifyError returns a reason other than UnknownError.
  • infra/jobs/ (Jobs.make, Jobs.run, the loop, the Worker registration) — trigger: the first recurring job.
  • The jobs_leases table, claim and release — trigger: the first job declared exclusive: 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-After in the relay — trigger: the first adapter whose upstream error carries a retry-after delay.