Skip to content

Activities: retries, backoff + jitter, timeouts, heartbeats, idempotency keys

Moduledur.05 · build · Go · Pass 8 · 4 to 5 h
You buildgo/durable/activity/activity.go: IdempotencyKey, Attempt, HeartbeatDetails, RecordHeartbeat, NewError, NewNonRetryable, IsNonRetryable, Backoff, Jitter; go/durable/server/activities.go: RecordHeartbeat, FailActivityTask, ListDeadLetters, RedriveDeadLetter
ContractRetryPolicy, ActivityOptions, Failure, and the heartbeat and fail rpcs in durable.proto; the Go API is section 4
Testscourse/tests/go/dur_05/ (what they check: section 4)
Needsdur.04 the protocol and the worker; reading: S-M07a (the expected value of a jittered backoff)
Used bydur.06 workflows and dur.08 cancellation use these activity failure and heartbeat handlers
MilestoneMS-durable
Optional depthBrooker, Exponential Backoff and Jitter (AWS Architecture Blog, 2015); Helland, Idempotence Is Not a Medical Condition (ACM Queue, 2012)
  • The delay before attempt n+1n+1 is min⁡(b0⋅rn−1,bmax⁡)\min(b_0 \cdot r^{n-1}, b_{\max}) (TestBackoffHandExample), and the server applies full jitter, a uniform draw in [0,d)[0, d), so a burst of failures does not retry in lockstep (TestJitterBoundsAndMean).
  • A failure that cannot succeed on retry (bad input) is non-retryable: marked by the activity or named by the policy’s non_retryable types, it ends the activity at once (TestNonRetryableStops).
  • Exhausted retries park the task in the DLQ and record ActivityDeadLettered; the workflow waits for a redrive instead of failing (TestRetryScheduleUnderFakeClock).
  • The idempotency key "<workflow_id>/<activity_id>" is the same on every attempt; with it, a flaky activity under SIGKILLs applies each effect exactly once (TestFlakyActivityKillLoopExactlyOnce).
Terminal window
ol start dur.05
ol tests dur.05
ol check dur.05 # needs dur.04 (or --ref-deps)
ol diff dur.05

Backoff and Jitter are pure functions: write and test them first, then the server’s FailActivityTask, then heartbeats.


With dur.04 a failing activity is retried at once, forever: a downstream service that is down gets hammered by every worker, a bad input burns every attempt, and an activity that was killed after charging a card charges it again on the retry. This module makes failure a first-class part of the engine: retries with exponential backoff and jitter, a dead-letter queue when they run out, non-retryable failures that stop at once, heartbeats that let a long activity keep its lease and resume from its last checkpoint, and the idempotency key that makes “at least once” safe. <system> wf start FlakyDemo --fail-rate 0.5 completes, with every effect applied once.

SymbolMeaningDefault
nnthe attempt that just failed, n≥1n \ge 1
b0b_0initial_ms1 s
rrbackoff, the growth factor (below 1 counts as unset)2.0
bmax⁡b_{\max}max_interval_ms, the cap60 s
dnd_nthe backoff before attempt n+1n+1
UUa uniform random number in [0,1)[0, 1)
δn\delta_nthe delay actually used

dn=min⁡(b0⋅r n−1,  bmax⁡)d_n = \min\left(b_0 \cdot r^{\,n-1},\; b_{\max}\right)

The first retry waits b0b_0, each later one rr times longer, never more than bmax⁡b_{\max}. A zero field means the default: a policy left empty must not retry in a hot loop.

If a dependency fails for 1,000 activities at once, they all retry after exactly d1d_1, fail together again, and so on: synchronized waves. Jitter spreads them:

δn=U⋅dn,U∼Uniform[0,1).\delta_n = U \cdot d_n, \qquad U \sim \mathrm{Uniform}[0, 1).

Then 0≤δn<dn0 \le \delta_n < d_n and E[δn]=dn/2\mathbb{E}[\delta_n] = d_n / 2 (the computation of S-M07a): the expected wait halves, and retries land spread uniformly over the window. The server draws UU from its own seeded generator (Options.Seed); tests can replace the whole step with Options.Jitter.

2.3 Retryable, non-retryable, dead-lettered

Section titled “2.3 Retryable, non-retryable, dead-lettered”

FailActivityTask(failure) on a live lease:

ConditionHistoryQueue
failure.non_retryable, or failure.type listed in RetryPolicy.non_retryableActivityTaskStarted, ActivityTaskFailed (the workflow is woken and sees an error)Dropped
the attempt reaches max_attempts (0 = [durable].dlq_after_attempts, 5)ActivityDeadLettered (informational: the workflow keeps waiting)DeadLettered
otherwisenothing (the retry lives in the queue; describe shows the last failure)Retrying after δn\delta_n

Why does an exhausted activity not fail the workflow? Because exhausting retries usually means a bug or an outage, which an operator fixes and then redrives (RedriveDeadLetter, attempt 1 again); failing the workflow would throw away hours of finished steps. ListDeadLetters shows each dead letter with its workflow, activity, attempts, and last failure (<system> wf dlq list data).

Activities choose their failure type: activity.NewNonRetryable("BadSpec", msg) fails for good; activity.NewError("Timeout", msg) is retryable with a type the policy can name. The subprocess runner of dur.09 maps exit code 65 to type ExitCode65, which policies list as non-retryable.

An attempt’s lease is its start-to-close timeout, or its heartbeat timeout hh when it has one: a heartbeating activity keeps its lease alive by calling activity.RecordHeartbeat(ctx, details) at least every hh. Each heartbeat moves the deadline to now +h+ h and stores the details (a checkpoint path, a byte offset). If the worker dies, the next attempt starts with activity.HeartbeatDetails(ctx) = the last details: it resumes instead of starting over. A heartbeat with a stale lease is FAILED_PRECONDITION, which the worker turns into a canceled context (dur.04).

Delivery is at least once: an attempt can apply its effect and die before reporting, so the next attempt applies it again. The fix is not in the engine but in the effect: give it a key that is the same for every attempt of the activity, activity.IdempotencyKey(ctx) = "<workflow_id>/<activity_id>", and make the receiver ignore a key it has seen (a unique constraint, an output directory named by the key, the effects sink of the course tests). Attempt number, worker id, or run id in the key would make every attempt a different effect.

Policy: b0=1b_0 = 1 s, r=2r = 2, bmax⁡=10b_{\max} = 10 s, max_attempts 4.

Failed attempt nnb0rn−1b_0 r^{n-1}dnd_nfull jitter at U=0.5U = 0.5
11 s1 s0.5 s
22 s2 s1 s
34 s4 s2 s
48 s8 s4 s
516 s10 s5 s
632 s10 s5 s

With jitter switched off (Options.Jitter = identity) and the clock at 0: attempt 1 fails at 0 and attempt 2 is pollable at exactly 1 s (not at 999 ms); attempt 2 fails and attempt 3 is pollable 2 s later; attempt 3 fails, 4 s later attempt 4; attempt 4 fails and, 4≥4 \ge max_attempts, the task is dead-lettered with ActivityDeadLettered{attempts 4} and no workflow task. A redrive delivers it as attempt 1. These are TestBackoffHandExample and TestRetryScheduleUnderFakeClock.

package activity // import "tinyllm/durable/activity"
const DefaultInitial, DefaultBackoff, DefaultInterval = time.Second, 2.0, time.Minute
func IdempotencyKey(ctx context.Context) string // "<workflow_id>/<activity_id>"
func Attempt(ctx context.Context) int
func HeartbeatDetails(ctx context.Context) []byte
func RecordHeartbeat(ctx context.Context, details []byte)
type Error struct { Type, Message string; NoRetry bool; Cause error } // FailureType(), NonRetryable()
func NewError(typ, msg string) error
func NewNonRetryable(typ, msg string) error
func IsNonRetryable(err error) bool
func Backoff(p *durablev1.RetryPolicy, n int) time.Duration
func Jitter(d time.Duration, u float64) time.Duration // u in [0, 1)
// go/durable/server/activities.go (methods on *Server)
func (s *Server) RecordHeartbeat(ctx context.Context, req *durablev1.HeartbeatRequest) (*durablev1.HeartbeatResponse, error)
func (s *Server) FailActivityTask(ctx context.Context, req *durablev1.FailActivityRequest) (*durablev1.Empty, error)
func (s *Server) ListDeadLetters(ctx context.Context, req *durablev1.ListDeadLettersRequest) (*durablev1.ListDeadLettersResponse, error)
func (s *Server) RedriveDeadLetter(ctx context.Context, req *durablev1.RedriveDeadLetterRequest) (*durablev1.RedriveDeadLetterResponse, error)
TestKINDChecksWhy it matters downstream
TestBackoffHandExampleunitsection 3’s schedule, the defaults, a non-integer factor, jitter at 0.5the formula the server applies
TestJitterBoundsAndMeanstatistical20,000 draws in [0,d)[0, d) with mean d/2d/2 within 4 sdretries spread, as S-M07a predicts
TestErrorTypesunitnon-retryable and typed errors survive fmt.Errorf("%w")the worker inspects them with errors.As
TestRetryScheduleUnderFakeClockunit, faultretries due at exactly 1, 2, 4 s; the 4th failure dead-letters; describe, ListDeadLetters, redrivewf dlq list, drill ops.03
TestNonRetryableStopsunita marked failure and a policy-listed type both end in ActivityTaskFailed after one attemptbad input is reported, not retried
TestWorkerMapsActivityErrorsunitan activity returning NewNonRetryable reaches history as such, after one attemptthe SDK end to end
TestHeartbeatExtendsLeaseAndResumesunit, faultthe deadline moves to now + h; describe shows the heartbeat; attempt 2 sees ckpt-3, attempt 2, and the same keyresume from checkpoint (dur.09)
TestStaleLeaseHeartbeatAndFailRefusedfaulta stale lease cannot heartbeat or faila presumed-dead worker cannot steer a retry
TestFlakyActivityKillLoopExactlyOncefault20 flaky activities, 8 worker SIGKILLs, server restarts: completed, every key applied once, some delivered twiceMS-durable’s exactly-once effects
PitfallSymptomCaught by
1. the exponent off by one, no cap, or the backoff of the next attempt instead of the failed onethe first retry waits too long, or retries wait for daysTestBackoffHandExample, TestRetryScheduleUnderFakeClock (mutants s01, s02, s08)
2. zero policy fields used as they area default policy retries in a hot loopTestBackoffHandExample (mutant s03)
3. jitter added on top of the backoffdelays up to twice the backoff, still in lockstep at the floorTestJitterBoundsAndMean (mutant s04)
4. non-retryable failures retried (the flag lost, the type matched against the message, the type dropped)bad input burns every attempt and lands in the DLQTestErrorTypes, TestNonRetryableStops, TestWorkerMapsActivityErrors (mutants s05, s09, s10, s17)
5. dead letters not recorded, without counts or names, or a redrive of an unknown task reported as an outageoperators cannot see what is stuck or whyTestRetryScheduleUnderFakeClock (mutants s06, s07, s14, s15)
6. heartbeat details not kept or not passed on; the deadline reported in the wrong unita resumed attempt starts from scratch; the SDK misjudges its leaseTestHeartbeatExtendsLeaseAndResumes (mutants s11, s13, s16)
7. an idempotency key that is not stable per activitythe same effect is applied once per attempt, or keys collide across workflowsTestFlakyActivityKillLoopExactlyOnce, TestHeartbeatExtendsLeaseAndResumes (mutant s12)
8. failing an attempt without checking its leasea stale worker fails the live attempt and triggers a needless retryTestStaleLeaseHeartbeatAndFailRefused (mutant s18)
DirectionModuleHow it uses this
Backdur.04the worker reports failures and heartbeats with these rpcs; Info carries the key and the details
BackS-M07athe expected value of a uniform draw, d/2d/2
Forwarddur.06ExecuteActivity futures resolve with ActivityTaskFailed; options carry the RetryPolicy
Forwarddur.08cancellation uses these activity failure and heartbeat handlers
Forwarddur.09the subprocess runner heartbeats checkpoint paths and maps exit codes to failure types
Your pieceProduction equivalentWhat it addsWhere to look
heartbeat detailsTemporal heartbeat detailsresumable activities with typed progress, throttled heartbeatsTemporal “Activity Heartbeats”
a worker reports each attemptasynchronous activity completionan external system completes the activity later by its task tokenTemporal “Asynchronous Activity Completion”
retry state in the queueretry state in visibilitythe current attempt, last failure, and next retry time are searchableTemporal PendingActivityInfo
start-to-close or heartbeat leaseschedule-to-close and schedule-to-start timeoutsbounds on the total time including retries, and on waiting in the queueTemporal “Activity Timeouts”