Schedules — Retries, Repeats & Polling
Declarative time policies in Effect v4 — retry/repeat schedules, backoff constructors, min vs max composition, input-typed predicates, jitter, instrumentation, and production-grade retry patterns.
A for loop with await sleep(n) is not a retry policy. It is a side-effectful loop that happens to wait. You cannot test it without waiting, compose it without nesting, log it without interleaving, or bound it without adding branches. In Effect, a schedule is the policy — an immutable value that decides whether to continue, how long to wait, and what it remembered.
// loop-shaped "policy" — not a valuefor (let attempt = 0; attempt < 5; attempt++) { try { return await fetch(url) } catch (e) { if (!isRetryable(e)) throw e await sleep(250 * 2 ** attempt) // jitter? cap? metrics? now what? }}
// value-shaped policy — inspectable, composable, virtual-timedimport { Schedule } from "effect"
const policy = Schedule.exponential("250 millis").pipe( Schedule.jittered, Schedule.while(({ input }) => input.retryable))This chapter is the schedule system end to end: how Schedule is modeled in v4, the constructor catalog, the composition semantics that actually matter (min vs max), typed inputs that gate retryability, delay shaping, instrumentation, and how schedules plug into Effect.retry, Effect.repeat, and streams.
Anatomy of a schedule
Section titled “Anatomy of a schedule”The type
Section titled “The type”import type { Schedule } from "effect"
// output input error services// │ │ │ │type RetryPolicy = Schedule<Duration.Duration, HttpError, never, never>// success value ↑ required input ↓ not yet failed? requirementsSchedule<Output, Input, Error, Env>:
| Parameter | Role | Where it surfaces |
|---|---|---|
Output |
what the schedule produces each step (often a delay or a count) | Effect.repeat returns it; Effect.retry ignores it |
Input |
what each step receives — the error for retry, the success for repeat |
typed by Schedule.setInputType<T>() |
Error |
how a schedule itself can fail (rare; e.g. Schedule.cron parse error) |
propagates into retry/repeat error channel |
Env |
services the schedule needs (opening the policy to logging, random, config) | must be provided alongside the effect |
Two invariants worth internalizing:
- A schedule is not stateful. The value holds no mutable state. Each execution obtains a fresh step function that closes over its own attempt counters and timing.
- Input is the decision signal. Retry schedules inspect the error; repeat schedules inspect the success. The output is not the decision — the decision is whether the schedule continues.
From value to step
Section titled “From value to step”v4 makes stepping explicit and replaces v3’s ScheduleDriver:
import { Duration, Effect, Schedule } from "effect"
const policy = Schedule.recurs(2)
const program = Effect.gen(function* () { // acquire a step function for this run — owns attempt/elapsed counters const step = yield* Schedule.toStep(policy)
// each call: timestamp (ms since epoch) + input → Pull<[Output, Duration]> const [output, delay] = yield* step(Date.now(), undefined) // output = 0 (zero-based count), delay = Duration.zero (recurs has no wait)
return { output, delay }})
// creating a schedule from a raw step function (library / custom policies)import { Cause, Pull } from "effect"
const onceAfter = (d: Duration.DurationInput) => Schedule.fromStep(Effect.sync(() => { let ran = false return (_now: number, _input: unknown) => { if (ran) return Cause.done(Duration.zero) ran = true return Effect.succeed([Duration.zero, Duration.fromInputUnsafe(d)] as const) } }))| Operation | Signature sketch | Use |
|---|---|---|
Schedule.toStep(schedule) |
Effect<(now, input) => Pull<[O, Duration], E>> |
drive manually, advance in tests |
Schedule.toStepWithSleep(schedule) |
Effect<(input) => Pull<O, E>> |
same but sleeps the delay for you |
Schedule.toStepWithMetadata(schedule) |
Effect<(input) => Pull<Metadata<O,I>, E>> |
sleeps and returns full metadata |
Schedule.fromStep(effectOfStep) |
step → schedule | build custom policies |
Schedule.fromStepWithMetadata(effectOfStep) |
metadata-aware step → schedule | predicates on timing without threading now yourself |
Metadata threading is explicit in v4. The per-step view:
import type { Duration, Schedule } from "effect"
type InputMeta<I> = Schedule.InputMetadata<I>// { input: I, attempt: number, start: number, now: number,// elapsed: number, elapsedSincePrevious: number }
type FullMeta<O, I> = Schedule.Metadata<O, I>// InputMeta<I> & { output: O, duration: Duration.Duration }Constructors
Section titled “Constructors”Every constructor returns a schedule that never fails on its own (aside from cron, which validates its expression). Delays that leave the schedule are Duration values — the retry/repeat machinery sleeps them for you.
Counting and spacing
Section titled “Counting and spacing”import { Schedule } from "effect"
// counts steps, zero delayconst threeRetries = Schedule.recurs(3) // outputs: 0, 1, 2, then doneconst forever = Schedule.forever // outputs: 0, 1, 2, ... (never done)
// fixed delay after each step completesconst everySecond = Schedule.spaced("1 seconds")const every30s = Schedule.spaced("30 seconds")
// fixed *cadence* aligned to a wall-clock grid (skips catch-up on slow actions)const onTheSecond = Schedule.fixed("1 seconds")const onTheMinute = Schedule.fixed("1 minutes")
// windowed to interval boundariesconst windowed5s = Schedule.windowed("5 seconds")
// single shot after a duration, then doneconst onceAfter5s = Schedule.duration("5 seconds")
// run for a wall-clock duration, zero delay each stepconst forOneMinute = Schedule.during("1 minutes")| Constructor | Output | Delay | Completes when |
|---|---|---|---|
Schedule.recurs(n) |
number (0-based count) |
Duration.zero |
after n steps |
Schedule.forever |
number |
Duration.zero |
never |
Schedule.spaced(d) |
number |
d |
never |
Schedule.fixed(interval) |
number |
aligned to interval grid |
never |
Schedule.windowed(interval) |
number |
to next window boundary | never |
Schedule.duration(d) |
Duration.Duration |
d (first step only) |
after 1 step |
Schedule.during(d) |
Duration.Duration (elapsed) |
Duration.zero |
when elapsed > d |
Schedule.exponential(base, factor?) |
Duration.Duration |
base * factor^(attempt-1) |
never |
Schedule.fibonacci(one) |
Duration.Duration |
Fib sequence from one |
never |
Schedule.cron(expr, tz?) |
Duration.Duration |
until next cron tick | never (or parse error) |
Backoff family
Section titled “Backoff family”import { Duration, Effect, Schedule } from "effect"
const expo = Schedule.exponential("200 millis") // 200, 400, 800, 1600, ...const expo3 = Schedule.exponential("100 millis", 3) // 100, 300, 900, 2700, ... (factor=3)const fib = Schedule.fibonacci("100 millis") // 100, 100, 200, 300, 500, ...
// cron — calendar scheduling, zone-awareconst everyMinute = Schedule.cron("* * * * *")const nightly2am = Schedule.cron("0 2 * * *", "America/New_York")const weekly = Schedule.cron("0 9 * * MON")There is no Schedule.delayed or Schedule.backoff helper in v4. Tune delays with Schedule.modifyDelay / Schedule.jittered (below) or by composing an exponential schedule with a cap.
Composition — the part that matters
Section titled “Composition — the part that matters”Schedules compose into new schedules. The two combining forms are easy to confuse until you think in terms of continuation and delay selection.
flowchart LR
A["Schedule A<br/>exponential 250ms"]:::a
B["Schedule B<br/>spaced cap 10s"]:::b
C["Schedule C<br/>recurs(5) limit"]:::c
A & B --> MIN{"Schedule.min([A,B])<br/>continue while ANY continues<br/>delay = FASTEST (minimum)"}:::min
A & C --> MAX{"Schedule.max([A,C])<br/>continue while ALL continue<br/>delay = SLOWEST (maximum)"}:::max
MIN --> MOUT["keeps trying:<br/>exponential capped at 10s"]
MAX --> XOUT["bounded: exponential<br/>but stops after 5 attempts"]
classDef a fill:#dbeafe,stroke:#2563eb
classDef b fill:#fef3c7,stroke:#d97706
classDef c fill:#dcfce7,stroke:#16a34a
classDef min fill:#eff6ff,stroke:#2563eb
classDef max fill:#f0fdf4,stroke:#16a34a
| Combinator | Continues while | Delay chosen | Output | Mnemonic / use |
|---|---|---|---|---|
Schedule.min([a, b, ...]) |
ANY continues (OR over termination) |
minimum (fastest) | Duration |
keep trying: cap an exponential, race two intervals picking the sooner |
Schedule.max([a, b, ...]) |
ALL continue (AND over termination) |
maximum (slowest) | Duration |
bound tightly: “exponential but no more than 5 times” |
Schedule.concat(a, b) |
a to completion, then b |
each in turn | O | O2 (merged via Result) |
phase change: fast retries then slow polls |
Schedule.concatResult(a, b) |
same sequencing | same | Result<O2, O> (which phase) |
track which policy produced the current step |
Schedule.upTo(s, { times, duration }) |
bounded by count/time | preserved | O |
shorthand for while(({attempt,elapsed})=>...) on s |
When to reach for which
Section titled “When to reach for which”import { Effect, Schedule } from "effect"
// CAPPED EXPONENTIAL (keep-trying semantics): exponential capped at 10s → min// Continues while EITHER wants to continue → effectively forever (both infinite)// Delay is the faster of the two → min(250*2^n, 10s)const cappedExpo = Schedule.min([ Schedule.exponential("250 millis"), Schedule.spaced("10 seconds")])
// BOUNDED EXPONENTIAL (stop after N): exponential limited to 6 attempts → max// Continues while BOTH want to continue → stops after 6 steps// Delay is the slower (exponential delay, recurs contributes zero)const boundedExpo = Schedule.max([ Schedule.exponential("250 millis"), Schedule.recurs(6)])
// TWO-PHASE: 3 fast retries at 200ms, then hourly pollingconst fastThenSlow = Schedule.concat( Schedule.spaced("200 millis").pipe(Schedule.upTo({ times: 3 })), Schedule.spaced("1 hours"))
// BOUND BY TIME OR COUNT — whichever firstconst atMost = Schedule.forever.pipe( Schedule.upTo({ times: 10, duration: "5 minutes" }))
// CUSTOM stop condition via `while`const untilSettled = Schedule.spaced("500 millis").pipe( Schedule.while(({ attempt, elapsed }) => attempt < 20 && elapsed < 60_000))The worked rule of thumb:
- You want a ceiling on delay (cap) →
min([...exponential, spaced(cap)]). Continues while either would → exponential never stops, cap never stops → still retries indefinitely, but delay never exceeds cap. - You want a limit on attempts or elapsed time →
max([...delayPolicy, recurs(n)])or.pipe(Schedule.upTo(...)). Continues only while both would → the limit’s termination wins.
Typed inputs and predicates
Section titled “Typed inputs and predicates”A schedule that retries everything fails fast on nothing. v4 gates retryability through the Input channel.
import { Schedule, Schema } from "effect"
class HttpError extends Schema.TaggedError<HttpError>()("HttpError", { message: Schema.String, status: Schema.Number, retryable: Schema.Boolean}) {}
// tell the type system what Input looks like — zero runtime costconst retryableOnly = Schedule.exponential("200 millis").pipe( Schedule.setInputType<HttpError>(), Schedule.while(({ input }) => input.retryable))
// equivalent inline via `while` predicate shape:// while receives Metadata<{output,input,duration,attempt,start,now,elapsed,elapsedSincePrevious}>const withStatusGate = Schedule.spaced("500 millis").pipe( Schedule.setInputType<HttpError>(), Schedule.while(({ input }) => input.status >= 500 && input.status < 600))
// effectful predicate — consult config or a serviceconst whileWithEffect = Schedule.spaced("500 millis").pipe( Schedule.setInputType<HttpError>(), Schedule.while(({ input }) => Effect.gen(function* () { if (!input.retryable) return false const budget = yield* RetryBudget // hypothetical service return yield* budget.hasRetriesRemaining }) ))| Step | What it does | Why |
|---|---|---|
Schedule.setInputType<T>() |
retypes the schedule’s Input to T |
without it, downstream while sees unknown and can’t read retryable |
Schedule.while(pred) |
keep recurring while predicate returns true |
failure to pass → Cause.done → retry/repeat stops |
Schedule.tap(f) |
same continuation as input, but runs an effect on each step’s metadata |
instrumentation (see below) |
Predicates are inclusive: true means “continue, wait duration”; false means “done — surface the last error/value”. This is why while(({ input }) => input.retryable) reads as “retry while retryable”.
Delay shaping, jitter, and instrumentation
Section titled “Delay shaping, jitter, and instrumentation”Jitter — not optional in production
Section titled “Jitter — not optional in production”If 500 clients share a spaced("1 seconds") schedule and a backend goes down for exactly 2 seconds, they all wake together and hammer the recovery. Jitter decorrelates wakeups.
import { Schedule } from "effect"
const polite = Schedule.exponential("250 millis").pipe( Schedule.jittered // each delay *= uniform [0.8, 1.2])
// capped exponential with jitter — the production patternconst cappedWithJitter = Schedule.min([ Schedule.exponential("250 millis"), Schedule.spaced("10 seconds")]).pipe(Schedule.jittered)
// per-request decorrelation without cappingconst perRequest = Schedule.spaced("1 seconds").pipe(Schedule.jittered)jittered scales each delay by a uniform factor in [0.8, 1.2] using Random.next. There is no jittered(max, min) overload in v4 — for fully custom distributions, use modifyDelay:
import { Duration, Effect, Random, Schedule } from "effect"
const fullJitter = Schedule.exponential("250 millis").pipe( Schedule.modifyDelay(({ duration }) => Effect.gen(function* () { const r = yield* Random.next // 0..1 const millis = Duration.toMillis(duration) const jittered = millis * r // 0..original return Duration.millis(jittered) }) ))
// bounded exponential via explicit modifyDelay cap (equivalent to min([...]))const cappedExplicit = Schedule.exponential("250 millis").pipe( Schedule.modifyDelay(({ duration }) => Effect.succeed(Duration.min(duration, Duration.millis(10_000))) ))tap — log and meter without changing the policy
Section titled “tap — log and meter without changing the policy”import { Duration, Effect, Schedule, Schema } from "effect"
class HttpError extends Schema.TaggedError<HttpError>()("HttpError", { message: Schema.String, status: Schema.Number, retryable: Schema.Boolean}) {}
const instrumented = Schedule.exponential("250 millis").pipe( Schedule.setInputType<HttpError>(), Schedule.tap((meta) => Effect.logDebug( `retry attempt=${meta.attempt} status=${meta.input.status} ` + `delay=${Duration.toMillis(meta.duration)}ms elapsed=${meta.elapsed}ms` ) ), Schedule.tap((meta) => // any effect — metrics, tracing, updating a budget counter Effect.sync(() => metrics.histogram("retry.delay.ms", Duration.toMillis(meta.duration))) ))
// metadata shape at tap:// { input: HttpError, output, duration, attempt, start, now, elapsed, elapsedSincePrevious }tap runs its effect during the schedule step; the returned schedule has the same Output/Input continuation as before. Multiple taps compose sequentially in pipe order.
Other shapers
Section titled “Other shapers”import { Duration, Effect, Schedule } from "effect"
// map outputs without touching delaysconst labeling = Schedule.recurs(3).pipe( Schedule.map(({ output }) => `attempt #${output + 1}`))
// add a fixed extra delay via a side-computed durationconst withExtraPause = Schedule.spaced("1 seconds").pipe( Schedule.addDelay(() => Effect.succeed("100 millis")))
// cap with flowing due diligence via `modifyDelay`:const decaying = Schedule.exponential("1 seconds").pipe( Schedule.modifyDelay(({ duration, attempt }) => // shrink after 5 attempts Effect.succeed(attempt > 5 ? Duration.millis(500) : duration) ))
// bounded by either count or elapsed timeconst limited = Schedule.forever.pipe(Schedule.upTo({ times: 10 }))const timed = Schedule.forever.pipe(Schedule.upTo({ duration: "30 seconds" }))const either = Schedule.forever.pipe(Schedule.upTo({ times: 20, duration: "2 minutes" }))Using schedules — retry and repeat
Section titled “Using schedules — retry and repeat”One mental model, two combinators
Section titled “One mental model, two combinators”| Combinator | Input fed to schedule | Output of schedule | Success means |
|---|---|---|---|
Effect.retry |
the error E |
ignored (delays only) | original effect eventually succeeded — return its A |
Effect.repeat |
the success A |
returned to the caller | combine repetitions into the schedule’s Output |
A retry schedule that completes means “no more retries — propagate the last error”. A repeat schedule that completes means “done repeating — return the accumulated output”.
Retry — three call shapes
Section titled “Retry — three call shapes”import { Duration, Effect, Schedule, Schema } from "effect"
class TransientError extends Schema.TaggedError<TransientError>()( "TransientError", { message: Schema.String, retryable: Schema.Boolean }) {}
declare const fetchUser: (id: string) => Effect.Effect<User, TransientError>
// 1) policy value — cleanest when the schedule is reusable/namedconst retryPolicy = Schedule.spaced("200 millis").pipe( Schedule.upTo({ times: 3 }), Schedule.setInputType<TransientError>(), Schedule.while(({ input }) => input.retryable))
const viaPolicy = fetchUser("42").pipe(Effect.retry(retryPolicy))
// 2) builder form — assists inference, avoids explicit setInputType// the `$` helper carries the error type throughconst viaBuilder = fetchUser("42").pipe( Effect.retry(($) => $(Schedule.spaced("1 seconds")).pipe( Schedule.while(({ input }) => input.retryable) ) ))
// 3) options object — for the times/while/until/schedule shorthandconst viaOptions = fetchUser("42").pipe( Effect.retry({ times: 3 }))
const viaScheduleField = fetchUser("42").pipe( Effect.retry({ schedule: retryPolicy }))
const viaPredicate = fetchUser("42").pipe( Effect.retry({ schedule: Schedule.exponential("200 millis"), while: (e: TransientError) => e.retryable }))
// `until` is the dual of `while` — stop retrying once the predicate holdsconst untilNotRetryable = fetchUser("42").pipe( Effect.retry({ schedule: Schedule.spaced("200 millis"), until: (e: TransientError) => !e.retryable }))
interface User { id: string }Effect.retry signatures in v4 (from packages/effect/src/Effect.ts):
Effect.retry(options: Retry.Options<E>) // { schedule?, times?, while?, until? }Effect.retry(policy: Schedule<B, NoInfer<E>, Error, Env>)Effect.retry(builder: ($: <O,E,R>(_: Schedule<O,NoInfer<E>,E,R>) => Schedule<O,E,E,R>) => Schedule<B,NoInfer<E>,Error,Env>)The dual two-argument form Effect.retry(self, policy) is also valid. The same three shapes exist on Effect.repeat with Repeat.Options<A> typed over the success channel.
Exhaustion is an error — handle it
Section titled “Exhaustion is an error — handle it”When the schedule says “done” while the effect is still failing, the failure propagates out of retry. If nothing handles it, you will see the last error with extra cause context about the schedule drain. Two patterns:
import { Effect, Schedule, Schema } from "effect"
class HttpError extends Schema.TaggedError<HttpError>()("HttpError", { status: Schema.Number, retryable: Schema.Boolean}) {}
declare const fetchUser: Effect.Effect<User, HttpError>interface User { id: string }
const policy = Schedule.recurs(3).pipe( Schedule.setInputType<HttpError>(), Schedule.while(({ input }) => input.retryable))
// 1) let typed failure propagate — caller pattern-matches the errorconst propagated: Effect.Effect<User, HttpError> = fetchUser.pipe(Effect.retry(policy))
// 2) after exhausting retries, escalate to a defect/fatalconst orDie: Effect.Effect<User, never> = fetchUser.pipe(Effect.retry(policy), Effect.orDie)
// 3) retryOrElse — custom recovery at exhaustion with attempt countconst withFallback: Effect.Effect<User, never> = Effect.retryOrElse(fetchUser, policy, (error, retryCount) => Effect.gen(function* () { yield* Effect.logWarning(`exhausted after ${retryCount} retries: ${error}`) return yield* serveStaleUser }) )
declare const serveStaleUser: Effect.Effect<User, never>Repeat — collecting successes
Section titled “Repeat — collecting successes”import { Effect, Schedule } from "effect"
declare const poll: Effect.Effect<number, string>
// repeat the effect and collect its successes according to the scheduleconst repeated: Effect.Effect<number, string> = poll.pipe(Effect.repeat(Schedule.recurs(3))) // poll runs 4 times total (1 initial + 3 repeats); returns last success (number)
// repeat and transform schedule output — e.g., count repetitionsconst counted = poll.pipe( Effect.repeat(Schedule.recurs(5).pipe(Schedule.map(({ output }) => output))))
// builder form, same as retry — $ carries success typeconst repeatBuilder = poll.pipe( Effect.repeat(($) => $(Schedule.spaced("500 millis")).pipe(Schedule.upTo({ times: 10 }))))
// conditional repeat predicate (over successes)const repeatUntilDone = poll.pipe( Effect.repeat({ until: (n: number) => n >= 100 }))const repeatWhileNotDone = poll.pipe( Effect.repeat({ while: (n: number) => n < 100 }))
// failure inside `repeat` still fails; to fold a failed repeat use repeatOrElseCombining schedules with other combinators
Section titled “Combining schedules with other combinators”Schedules compose with timeout only via explicit sequencing — there is no implicit timeout inside retry. Layer them intentionally:
import { Effect, Schedule } from "effect"
declare const fetchWithTimeout: Effect.Effect<Response, FetchError>class FetchError extends Error {}
const resilient = fetchWithTimeout.pipe( Effect.timeout("3 seconds"), // each attempt has a deadline Effect.retry( // retries wrap the timed effect Schedule.exponential("200 millis").pipe(Schedule.upTo({ times: 4 })) ), Effect.catchTag("TimeoutError", () => Effect.succeed(Response.cached)))interface Response { readonly body: string }declare namespace Response { const cached: Response }Because Effect.timeout interrupts on deadline (see chapter 13), the retried branch releases its resources correctly between attempts.
Schedules and streams — polling without loops
Section titled “Schedules and streams — polling without loops”Two distinct roles for schedules in streaming code:
Schedule as a stream source for polling
Section titled “Schedule as a stream source for polling”import { Effect, Schedule, Stream } from "effect"
declare const fetchStatus: Effect.Effect<Status, FetchError>type Status = { done: boolean }
class FetchError extends Error {}
// Option A — repeat the effect itself on a schedule, streaming each success:const pollingRepeat: Stream.Stream<Status, FetchError> = Stream.fromEffect(fetchStatus).pipe( Stream.repeat(Schedule.spaced("2 seconds")) )
// Option B — produce a tick stream from a schedule and drive fetching off ticks:const tickThenFetch: Stream.Stream<Status, FetchError> = Stream.fromSchedule(Schedule.spaced("2 seconds")).pipe( Stream.flatMap(() => Stream.fromEffect(fetchStatus)) )
// Bounded polling: up to 30 times or until done - predicate on Streamconst boundedPoll = pollingRepeat.pipe( Stream.takeUntil((s) => s.done), Stream.take(30))Reusing a polling schedule as a retry schedule
Section titled “Reusing a polling schedule as a retry schedule”The delay policy is the same value in both contexts — keep it DRY:
import { Effect, Schedule, Stream } from "effect"
const pollingSchedule = Schedule.spaced("1 seconds").pipe(Schedule.upTo({ times: 20 }))
declare const fetchFeed: Effect.Effect<Array<Item>, FetchError>type Item = { id: string }class FetchError extends Error {}
// as retry backoff inside a stream stepconst feedWithRetries: Stream.Stream<Item, FetchError> = Stream.fromIterableEffect(fetchFeed).pipe( Stream.flatMap((items) => Stream.fromIterable(items)), Stream.catchAll((e) => Stream.fail(e)) )// retry at the effect level wraps the *call*, stream governs the *flow*
// as effect retry for the fetch itself, then spread into a streamconst feedStream: Stream.Stream<Item, FetchError> = Stream.fromEffect(fetchFeed.pipe(Effect.retry(pollingSchedule))).pipe( Stream.flatMap((items) => Stream.fromIterable(items)) )Stream-specialized schedule combinators
Section titled “Stream-specialized schedule combinators”Stream exposes schedule-like controls directly on the stream — internally they reuse the same step machinery:
import { Schedule, Stream } from "effect"
const source: Stream.Stream<number> = Stream.range(1, 100)
const shaped = source.pipe( // repeat each element according to a schedule Stream.repeat(Schedule.spaced("100 millis")), // schedule the stream — control when new pulls happen Stream.schedule(Schedule.spaced("500 millis")), // timeout between elements Stream.timeout("3 seconds"), Stream.timeoutOrElse("3 seconds", () => Stream.succeed(-1)))What v4 removed or renamed
Section titled “What v4 removed or renamed”If you are porting v3 schedule code, a short migration checklist:
| v3 | v4 | Note |
|---|---|---|
ScheduleDriver |
Schedule.toStep / fromStep |
no opaque driver type — pure step functions |
Schedule.either(a,b) |
Schedule.min([a,b]) |
continues while any continues, fastest delay |
Schedule.intersect / both |
Schedule.max([a,b]) |
continues while all continue, slowest delay |
Schedule.andThen |
Schedule.concat |
sequential — a to completion, then b |
Schedule.check(pred) |
Schedule.while(pred) |
effectful predicate on Metadata |
Schedule.delayed(f) |
Schedule.modifyDelay(f) |
re-derive delay from Metadata |
Schedule.count |
Schedule.forever |
zero-based output; cap with upTo({ times }) |
Schedule.collectAll / collectWhile / etc |
removed | use Effect.repeat + Stream collection or Schedule.map + Schedule.CurrentMetadata |
Schedule.capped / manual caps |
Schedule.min([expo, spaced(cap)]) |
idiomatic capped exponential |
Schedule.retry(n) on Effect |
Effect.retry({ times: n }) or Schedule.recurs(n) |
no combinator called retry on Schedule itself |
The collect* family removal is intentional: schedules are no longer a collection tool. Returning accumulated arrays from retries couples concurrency, memory, and policy. Collect with Stream.runCollect, or keep the schedule’s Output lean and let the surrounding effect decide accumulation.
The production retry schedule — worked example
Section titled “The production retry schedule — worked example”The canonical production pattern from ai-docs/src/06_schedule/10_schedules.ts combines everything above. Walk it line by line, then adapt it.
import { Duration, Effect, Schedule, Schema, Random } from "effect"
export class HttpError extends Schema.TaggedError<HttpError>()("HttpError", { message: Schema.String, status: Schema.Number, retryable: Schema.Boolean}) {}
// ── the policy ──// 1. exponential starting at 250 ms, doubled each retry (500, 1000, 2000, ...)// 2. capped at 10 s — min([...]) picks the faster of exponential vs the cap// 3. decorrelated with jitter so concurrent clients desynchronize// 4. input-typed to HttpError so `retryable` is checked at the type level// 5. non-retryable errors fail fast without waiting or wasting attemptsexport const productionRetrySchedule = Schedule.min([ Schedule.exponential("250 millis"), Schedule.spaced("10 seconds")]).pipe( Schedule.jittered, Schedule.setInputType<HttpError>(), Schedule.while(({ input }) => input.retryable))
export const fetchUserProfile = Effect.fn("fetchUserProfile")( function* (userId: string) { const r = yield* Random.next const status = r > 0.7 ? 200 : r > 0.3 ? 503 : 401 if (status !== 200) { return yield* new HttpError({ message: `Request for ${userId} failed`, status, retryable: status >= 500 // 5xx → retry, 4xx → fail immediately }) } return { id: userId, name: "Ada Lovelace" } as const })
// ── using it ──export const loadUserWithRetry = fetchUserProfile("user-123").pipe( Effect.retry(productionRetrySchedule), // exhaustion surfaces as HttpError — orDie only at a boundary Effect.orDie)
// same policy, builder inference at callsite (no explicit setInputType needed)export const loadWithBuilder = fetchUserProfile("user-123").pipe( Effect.retry(($) => $(Schedule.spaced("1 seconds")).pipe( Schedule.while(({ input }) => input.retryable) ) ), Effect.orDie)
// observability variant — log each retry's context without altering semanticsexport const instrumentedRetrySchedule = Schedule.min([ Schedule.exponential("250 millis"), Schedule.spaced("10 seconds")]).pipe( Schedule.jittered, Schedule.setInputType<HttpError>(), Schedule.tap((meta) => Effect.logDebug( `retry #${meta.attempt} status=${meta.input.status} ` + `next=${Duration.toMillis(meta.duration)}ms elapsed=${meta.elapsed}ms` ) ), Schedule.while(({ input }) => input.retryable))
// bounded variant: same backoff but at most 6 attempts totalexport const boundedRetrySchedule = Schedule.max([ // capped exponential delays Schedule.min([ Schedule.exponential("250 millis"), Schedule.spaced("10 seconds") ]).pipe(Schedule.jittered), // attempt cap — `max` means "continue only while BOTH continue" Schedule.recurs(6)])Read the composition out loud to check it: “min of exponential and once-per-10s is still an infinite retry policy whose delay is the faster of the two — so it caps. Jitter decorrelates it. While input.retryable gates it so we fail fast on 4xx.” If any clause in that sentence surprises you, translate it back to the table in the composition section — it is a direct rendering of the min/max semantics.
Testing with virtual time — don’t actually wait
Section titled “Testing with virtual time — don’t actually wait”Every schedule ultimately sleeps its computed delay. Under TestClock, sleeps are virtual — TestClock.adjust moves time forward without wall-clock delays:
import { Effect, Schedule, TestClock } from "effect"import { it } from "vitest"
class FlakyError { readonly _tag = "FlakyError" }
it("retries with exponential backoff without waiting", async () => { let attempts = 0 const flaky = Effect.gen(function* () { attempts++ if (attempts < 3) return yield* new FlakyError() return "ok" })
const policy = Schedule.exponential("100 millis").pipe( Schedule.upTo({ times: 5 }), Schedule.setInputType<FlakyError>() )
const program = Effect.gen(function* () { const fiber = yield* Effect.fork(flaky.pipe(Effect.retry(policy))) // advance virtual time past two backoffs: 100ms + 200ms yield* TestClock.adjust("500 millis") const result = yield* fiber.await return result })
const outcome = await Effect.runPromise(Effect.provide(program, TestClock.layer())) // outcome succeeds with "ok" — no real timeouts in the test run})See chapter 28 — Testing for the full TestClock repertoire.
Cheat sheet
Section titled “Cheat sheet”| Goal | Expression |
|---|---|
| Retry up to 3 times | Schedule.recurs(3) |
| Retry every 2 s | Schedule.spaced("2 seconds") |
| Exponential backoff from 200 ms | Schedule.exponential("200 millis") |
| Cap backoff at 10 s | Schedule.min([Schedule.exponential("250 millis"), Schedule.spaced("10 seconds")]) |
| Bound backoff to 6 tries | Schedule.max([Schedule.exponential("250 millis"), Schedule.recurs(6)]) |
| Gate on error field | .pipe(Schedule.setInputType<E>(), Schedule.while(({ input }) => input.retryable)) |
| Jitter all delays | .pipe(Schedule.jittered) |
| Custom jitter / clamp | .pipe(Schedule.modifyDelay(({ duration }) => ...)) |
| Log each decision | .pipe(Schedule.tap((meta) => Effect.logDebug(...))) |
| Bound by time or count | .pipe(Schedule.upTo({ duration: "30 seconds", times: 5 })) |
| Sequence fast→slow | Schedule.concat(fast, slow) |
| Apply to an effect | effect.pipe(Effect.retry(schedule)) or Effect.retry(effect, schedule) |
| Infer input type inline | effect.pipe(Effect.retry(($) => $(schedule).pipe(Schedule.while(...)))) |
| Shorthand times | effect.pipe(Effect.retry({ times: 3 })) |
| Gate retry predicate inline | effect.pipe(Effect.retry({ schedule, while: (e) => e.retryable })) |
| Exhaustion → defect | .pipe(Effect.retry(schedule), Effect.orDie) |
| Exhaustion → custom recovery | Effect.retryOrElse(effect, schedule, (e, n) => ...) |
| Repeat a success | effect.pipe(Effect.repeat(schedule)) |
| Produce ticks as a stream | Stream.fromSchedule(schedule) |
| Drive manually | yield* Schedule.toStep(schedule) then yield* step(Date.now(), input) |