Skip to content

Deterministic replay workflows, ContinueAsNew, history paging, payload limits

Moduledur.06 · build · Go · Pass 8 · 6 to 8 h
You buildgo/durable/workflow/workflow.go: ExecuteActivity, Sleep, Now, SideEffect, GetVersion, ContinueAsNew, ContinueAsNewSuggested; go/durable/workflow/replay.go: the replay engine, Registry, ReplayHistory, LoadHistory, ReplayHistoryFromFile; go/durable/server/continue.go: the record_marker and continue_as_new commands
Contractthe Go SDK surface of DESIGN 2.7, the Command and history messages of durable.proto; the package API is section 4
Testscourse/tests/go/dur_06/ (what they check: section 4); recorded runs of the course test workflows in course/fixtures/dur/histories/
Needsdur.05 activities and their failures (and through it the whole server)
Used bydur.07 timers and dur.08 signal/cancel handling extend this workflow SDK
MilestoneMS-durable
Optional depththe Temporal docs on deterministic constraints and versioning; Fateev and Abbas, Temporal: durable execution talks; Hellerstein et al. on deterministic replay
  • A workflow never runs once from start to finish. Every workflow task re-runs it from the start against the history: past results come from events, and each command it issues must match the event recorded for it (TestReplayHandExample).
  • Determinism is the price: no wall clock (workflow.Now), no randomness (SideEffect), no goroutines or I/O. A code change that reorders commands is caught as ErrNondeterminism, never silently continued (TestNondeterminismDetected, TestReorderFailpointNondeterministic).
  • GetVersion lets a deploy change code under runs already in flight: old histories replay the old path, new runs take the new one (TestGetVersionMigration).
  • A long workflow bounds its history with ContinueAsNew when the server suggests it, at 10,000 events or 32 MiB (TestContinueAsNewKeepsHistoryBounded).
  • Kill the server in the middle of a run: the worker replays and the result equals an uninterrupted run’s (TestKillServerMidWorkflowSameResult).
Terminal window
ol start dur.06
ol tests dur.06
ol check dur.06 # needs dur.05 (or --ref-deps)
ol diff dur.06

Build the coroutine first (start, block, step) with a workflow that only calls ExecuteActivity; pass TestReplayHandExample; then the matcher and ErrNondeterminism; then markers; then ContinueAsNew.


Workers run activities, but the logic that decides what to run next (train, then evaluate, then if the loss is good release, else retry with a smaller learning rate) still lives nowhere durable. If you write it as ordinary Go in a worker, a crash in the middle loses its local variables and its place in the code. Durable execution keeps writing ordinary Go and gets crash-proof state anyway, by recording the outcome of every step in history and replaying the code against it. This module is the SDK that does it. After it, CorpusBuild, TrainRun, ModelRelease, and AgentRun are plain functions that survive any crash.

SymbolMeaning
WWthe workflow: deterministic code from its input and the results of its commands to commands
hhthe run’s history (dur.02)
T1,…,TmT_1, \dots, T_mits completed workflow tasks: each WorkflowTaskCompleted names the WorkflowTaskStarted it answers
CjC_jthe command events recorded right after TjT_j‘s WorkflowTaskCompleted, in one batch (dur.04)
sjs_jthe event id of TjT_j‘s WorkflowTaskStarted
s∗s_*the id of the current task’s WorkflowTaskStarted, the last one in hh

To handle the current task:

start W as a coroutine
for j in 1..m: # past tasks, in order
deliver every event with id < s_j # results resolve futures
run W until it blocks; each command it issues must match the next event of C_j
if a command differs, or C_j has events left over: ErrNondeterminism
deliver every event with id < s_*
run W until it blocks; the commands it issues now are new: send them

The workflow cannot tell a replay from a first run: Get on an activity future returns the recorded result if its ActivityTaskCompleted was delivered, and otherwise blocks. Because commands are matched in order, the engine detects any change in what the code does at any past point.

The coroutine. The workflow runs in its own goroutine, but only one of the two (engine, workflow) runs at a time: the engine resumes it through a channel and waits; the workflow, when it must wait on a future, signals back and waits. So a replay is a deterministic function of the code and the history. When the task is done the engine kills the coroutine (it exits with runtime.Goexit from inside its wait).

Matching. A command matches an event when:

CommandEventCompared
ScheduleActivityActivityTaskScheduledactivity id and activity type
StartTimerTimerStartedtimer id
RecordMarkerMarkerRecordedmarker name
CompleteWorkflow, FailWorkflow, ContinueAsNewcompleted, failed, continuedkind

Activity and timer ids are the command’s sequence number in the run (“1”, “2”, …), so the same code issues the same ids in the same order on every replay.

Anything that can differ between two runs of the same code breaks replay:

Instead ofUseWhy it is safe
time.Now()workflow.Now(ctx)the time of the current task’s WorkflowTaskStarted, recorded in history
rand, uuidworkflow.SideEffect(ctx, f)f runs once; its JSON result is recorded as a SideEffect marker; replays return the recorded value and never call f
time.Sleepworkflow.Sleep(ctx, d)a durable timer (dur.07)
goroutines, channels, I/Oactivitiesactivities run in workers and their results are recorded
iterating a mapa sorted sliceGo randomizes map order

A deploy changes workflow code while runs are in flight. Their histories were made by the old code; replaying them with new code that issues different commands is nondeterminism. GetVersion(ctx, changeID, min, max) makes a change safe:

if workflow.GetVersion(ctx, "use-b", 0, 1) == 0 {
// old path, as the code was before the change
} else {
// new path
}

On a first execution it returns max and records a Version marker {change_id, version}. On replay, if the history has that marker at this point, it returns the recorded version (and the marker is consumed); if not, the history predates the change and it returns min without consuming anything. Later calls with the same change id return the same value.

A workflow that loops (train for 10 epochs, evaluate every 500 steps for a week) grows its history without bound, and every workflow task replays all of it. The server suggests continuing once the history passes 10,000 events or 32 MiB (continue_as_new_suggested); workflow.ContinueAsNewSuggested(ctx) recomputes the same rule from the history as of the current task, so it is the same on every replay. The workflow returns workflow.ContinueAsNew(ctx, state): the SDK issues a ContinueAsNew command; the server closes the run with WorkflowExecutionContinuedAsNew{new_run_id, input} and, once that is durable, starts the new run under the same workflow id with continued_from_run_id. The server refuses to grow a history past 20,000 events at all.

History paging. A history over 1 MiB reaches the worker in pages (dur.04); the engine needs every event, and the first page alone is not enough.

Payload limits. Payloads over 2 MiB are refused by the server: activity inputs and results, workflow results (dur.02, dur.04), and here marker details and continue-as-new inputs. Large data goes to /artifacts and travels by path.

The course test workflows (Sequence, Parallel, SideEffects, Versioned, Reorder, Clock, Counter) were run once against the reference and their histories recorded in course/fixtures/dur/histories/. Your engine must replay them all with the same results. For your own workflows the recipe is the same in two steps: run once through your server and worker, fetch the history, and replay it with ReplayHistory; any change that breaks determinism then fails a test instead of a production run.

Sequence(2) calls Upper("s1"), then Upper("s2"), and returns the list. Its history after both activities, with the current task started at event 15:

idevent
1started{Sequence, input 2}
2, 3wt_scheduled, wt_started
4wt_completed{started 3}
5act_scheduled{id "1", Upper, "s1"}
6, 7act_started, act_completed{scheduled 5, "S1"}
8, 9wt_scheduled, wt_started
10wt_completed{started 9}
11act_scheduled{id "2", Upper, "s2"}
12, 13act_started, act_completed{scheduled 11, "S2"}
14, 15wt_scheduled, wt_started (current)

Handling the task whose history ends at 15:

  1. Past task 1 (s1=3s_1 = 3): deliver events 1 and 2 (nothing to resolve). Run: the code calls ExecuteActivity("Upper", "s1"), sequence 1, and blocks on Get. Command ScheduleActivity{"1", Upper} matches event 5. C1C_1 is consumed.
  2. Past task 2 (s2=9s_2 = 9): deliver events 3 to 8; event 7 resolves activity “1” with "S1". Run: Get returns "S1"; the code calls ExecuteActivity("Upper", "s2"), sequence 2, and blocks. Command "2" matches event 11.
  3. Current task (s∗=15s_* = 15): deliver events 9 to 14; event 13 resolves “2” with "S2". Run: the code returns ["S1","S2"]. The new command is CompleteWorkflow{["S1","S2"]}.

Handling the earlier tasks (history up to 3, then up to 9) gives ScheduleActivity{"1"} and then ScheduleActivity{"2"}: the same function, more history. This is TestReplayHandExample.

package workflow // import "tinyllm/durable/workflow"
type Workflow func(ctx Context, input []byte) ([]byte, error)
type Context struct{ /* unexported */ }
var ErrNondeterminism = errors.New("durable: command does not match history")
type ActivityOptions struct {
TaskQueue string; StartToClose, ScheduleToClose, HeartbeatTimeout time.Duration; Retry *durablev1.RetryPolicy
}
type ActivityError struct { ActivityID, ActivityType string; Failure *durablev1.Failure }
type Future[O any] interface { Get(ctx Context) (O, error); Ready() bool }
func ExecuteActivity[O any](ctx Context, name string, in any, o ActivityOptions) Future[O] // in: []byte as is, else JSON
func Sleep(ctx Context, d time.Duration) error
func Now(ctx Context) time.Time
func SideEffect[T any](ctx Context, f func() T) T
func GetVersion(ctx Context, changeID string, min, max int) int
func ContinueAsNew(ctx Context, input []byte) error // return it from the workflow
func ContinueAsNewSuggested(ctx Context) bool
func WorkflowID(ctx Context) string
func RunID(ctx Context) string
type Registry struct{ /* unexported */ }
func NewRegistry() *Registry
func (r *Registry) Register(name string, wf Workflow)
func (r *Registry) HandleWorkflowTask(ctx context.Context, task *durablev1.WorkflowTask, history []*durablev1.HistoryEvent) ([]*durablev1.Command, error)
type HistoryFile struct { WorkflowID, RunID, WorkflowType string; Events []json.RawMessage }
func LoadHistory(path string) (*HistoryFile, []*durablev1.HistoryEvent, error)
func ReplayHistory(r *Registry, workflowType string, history []*durablev1.HistoryEvent) ([]byte, error)
func ReplayHistoryFromFile(r *Registry, path string) ([]byte, error)
// go/durable/server/continue.go: registers "record_marker" and "continue_as_new" with the dur.04 command table

Registry is the worker.WorkflowHandler of dur.04: worker.New(conn, worker.Options{Workflows: reg, ...}). Workflows are Go only; Python implements activities (DESIGN 2.8).

TestKINDChecksWhy it matters downstream
TestReplayHandExampleunitsection 3: the three tasks issue "1", "2", then complete with ["S1","S2"]the definition of replay
TestRecordedHistoriesReplayconformance, goldenevery recorded course test run replays with its recorded result (continued runs end in ContinueAsNew)the SDK conformance suite
TestNondeterminismDetectedunit, faultanother activity type, fewer commands, an extra command: ErrNondeterminism, replayed and livecode changes cannot corrupt runs
TestReorderFailpointNondeterministicfaultdur/workflow/reorder flips the order: the same ErrNondeterminism three times; passes without itthe failpoint of the catalog
TestSideEffectNotReexecutedunitreplay returns recorded values and never calls the function; a first run calls it oncerandom ids stay stable
TestGetVersionMigrationunit, regressionold history on the old path, new history on the new path, a new run records {use-b, 1}safe deploys of dur.11 and dur.12
TestNowComesFromHistoryunitNow equals the WorkflowTaskStarted times, on replay toono wall clock in workflows
TestRecordThenReplayOwnRununita run through your server and worker replays to the same resulthow you test your own workflows
TestKillServerMidWorkflowSameResultfaultthe server dies mid-run; the result equals an uninterrupted run and the history replaysMS-durable’s replay determinism
TestContinueAsNewKeepsHistoryBoundedfault, boundary12,000 steps continue as new past 10,000 events; each run under 20,000; pages over 1 MiB fetchedTrainRun, AgentRun
TestServerCommandLimitsboundarymarker and continue-as-new payloads over 2 MiB refused; the next run keeps type and queuepayload limits (2.7)
PitfallSymptomCaught by
1. replay plumbing: dropping a delivered result, collecting a task’s commands from the wrong place, completing with the wrong value, taking the first WorkflowTaskStarted as currentthe workflow sees empty results, every past task looks nondeterministic, or old commands are issued againTestReplayHandExample, TestRecordThenReplayOwnRun, TestKillServerMidWorkflowSameResult (mutants s02, s04, s07, s14)
2. matching too loosely (id without type, any completion)a changed workflow continues silently on a history it did not makeTestReorderFailpointNondeterministic, TestNondeterminismDetected (mutants s03, s09)
3. ids that do not advance with each commandtwo activities share an id; results go to the wrong futureTestReplayHandExample (mutant s01)
4. calling a side effect again on replaya new random id on every replay forks the workflow’s stateTestSideEffectNotReexecuted (mutant s05)
5. versioning: the new path for old histories, the wrong version recorded, the marker not consumedevery in-flight run breaks at the next deployTestGetVersionMigration (mutants s06, s11, s12)
6. ContinueAsNew treated as a failure, suggested too late, or the next run without its type and queuelong workflows fail, hit the history cap, or start a run no worker pollsTestContinueAsNewKeepsHistoryBounded, TestServerCommandLimits (mutants s10, s15, s16)
7. reading the wall clockeach replay sees a different time and takes a different branchTestNowComesFromHistory (mutant s13)
8. no payload limits on markers or continue-as-new inputsmessages over 4 MiB, histories that cannot be servedTestServerCommandLimits (mutants s17, s18)
DirectionModuleHow it uses this
Backdur.05activity futures resolve with results, ActivityErrors, and redriven attempts
Forwarddur.07Sleep issues StartTimer; the server’s timer commands and the timer service fire it
Forwarddur.08signals and cancellation add event hooks and commands to the same engine
Forwarddur.11, dur.12, data.09, ag.05TrainRun, EvalSuite, ModelRelease, CorpusBuild, and AgentRun are workflows on this SDK
Your pieceProduction equivalentWhat it addsWhere to look
replay from event 1 on every tasksticky cachesthe worker keeps a run’s coroutine between tasks and replays nothing; the server sends only new eventsTemporal “Sticky Execution”
one coroutine per workflowfull SDK command state machinesworkflow.Go, selectors, child workflows, cancellation scopes, each with its own state machinetemporalio/sdk-go: internal/internal_decision_state_machine.go
GetVersion by handworker deployment versioningold runs stay pinned to the build that started themTemporal “Worker Versioning”
replay on fixturesdeterminism checks on upgradereplay production histories against new code in CI before deployingTemporal “Replay Testing”
a crash mid-activity retries itsaige agent/durable/local: ErrIndeterminate and Reconcilean effect whose outcome is unknown is reconciled, not blindly retriedsaige agent/durable/local