Skip to content

Delivered stream records on Workflow Tasks and on replay. - #7

Open
moedash wants to merge 17 commits into
moe/AI-198-srv-4-workflow-commandsfrom
moe/AI-198-srv-5-delivery-replay
Open

moedash wants to merge 17 commits into
moe/AI-198-srv-4-workflow-commandsfrom
moe/AI-198-srv-5-delivery-replay

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR delivers stream records on the workflow task and re-supplies them on replay.

A subscribed workflow gets its records as slices on RecordWorkflowTaskStarted and on the poll response. The task records only the offset range it consumed, on WorkflowTaskCompleted. Replay re-reads those ranges from the stream and hands them back, so a cold worker sees what the original execution saw without History ever having held a payload.

What changed?

  • service/history/api/recordworkflowtaskstarted/stream_slices.go: staging a range per subscription, reading it from an owned stream in this execution or from a standalone one through the routed client, and attaching the slices. A duplicate of the same request hands back the same slice, since the staged range is redelivered rather than recomputed.
  • Replay walks every page of the history being sent, not just the first, and re-reads each recorded range. A range that is truncated, deleted or over the replay budget fails the workflow task with its own cause rather than handing matching a retry.
  • chasm/lib/workflow/workflow.go and service/history/workflow/: the cursor advance is folded in as the completed event is built, so the advance and the event carrying the range are one transaction. A subscription with offsets left to deliver schedules a workflow task on transaction close, which also covers the transaction that registers a subscription against a stream that already has data and completes no task of its own.
  • GetStreamReplaySlices on the history service, called by matching after it has fetched the history for a non-sticky query task. A query dispatched straight through matching is built without RecordWorkflowTaskStarted, so nothing else re-supplies its ranges and a cold worker could not replay a consuming workflow to answer.
  • tests/testcore/onebox.go can start more than one history host, which the cross-host routing test needs.
  • Functional tests for delivery, continue-as-new, replay past the first history page, a deleted stream failing the task loudly, the query path and cross-host routing.
  • tests/sdk_server_host_test.go is an opt-in host that serves a real SDK against this build. tests/stream_publish_cost_test.go is an opt-in measurement of what publishing costs in History; it holds one cluster per arm for the whole run, which outlasts the dedicated pool and would park the functional suite until its timeout.

Why?

This is the half of the design that has to hold under replay. Reviewing it next to the publish path would mix "what does a workflow ask for" with "what does the server owe it afterwards".

How did you test it?

go build ./...
make proto            # tree stays clean, generated artifacts match the protos
make fmt              # tree stays clean
go test ./chasm/lib/... ./service/history/... ./service/matching/ ./common/... -count=1
go test -tags test_dep ./tests/ -count=1 -run 'TestStream|TestReplay|TestSubscrib|TestWorkflowConsumes|TestQueryTask|TestMessagesAppended|TestConsumerOf|TestACompletedConsumer|TestANewRun|TestTruncation|TestRetentionWaits|TestDeleteStreamRefuses|TestOutsideAppends|TestOwnedStream|TestOverLimit|TestPublishOnAFailedTask'

All green.

Series

Server 5 of 6. Replaces #2.

Previous: #6 moe/AI-198-srv-4-workflow-commands. Next: #8 moe/AI-198-srv-6-reset-failover.

Review round

A speculative workflow task is given no stream range. It can be thrown away, and a discarded one writes no completed event, so nothing commits the cursor: the range came back on the next task while the worker's in-memory workflow had already consumed it, and live execution saw the batch twice where replay saw it once. Without the fix, TestSpeculativeTaskConsumesNoStreamRange also trips the server's dirty-mutable-state assertion, because the delivery dirties a transaction the discard throws away. A workflow with records waiting gets a normal task of its own, which is where they belong.

The three cold-replay bounds are namespace dynamic config now: stream.replayMaxRecords, stream.replayMaxBytes and stream.replayMaxPages. A consumer past one of them still cannot be replayed by this path, but it no longer fails the same task forever: the second consecutive STREAM_RANGE_UNAVAILABLE terminates the workflow with that cause. One retry is kept because a shard that was moving can serve the range on the next attempt; a second failure is not a race, and the old behaviour grew History until the size limit terminated the workflow for an unrelated reason.

A cursor naming a stream the workflow neither owns nor marked external now reports an unavailable range rather than a bare FailedPrecondition, so the task fails once with a cause an operator can read instead of coming back from matching forever. That is the state an event-based-replication standby rebuilds into, and #8 asserts it end to end.

Every routed read one task start makes now shares a single deadline, the same stream.RoutedSetBudget the completion path uses, because a workflow may consume as many streams as the subscription limit allows and the execution's lock is held across all of them.

Matching only asks History for the ranges a query task needs when the workflow actually holds a cursor. GetMutableStateResponse carries consumes_streams, read from the mutable state matching already fetched, so a query on any other workflow no longer takes a workflow lease on the history side or inherits a new way to fail.

Known limits, argued rather than fixed:

  • There is no bounded paged re-supply. A consumer whose recorded ranges outgrow the replay budget is terminated with the cause rather than recovered. A compaction point recorded in History, past which earlier ranges need not be re-supplied, is the shape of the fix and is a design change of its own.
  • attachReplaySlices still returns early on a sticky task. A worker can hold a sticky queue and have evicted the workflow, and the recovery is the SDK's: it fails the task, the next one is dispatched on the normal queue with full history, and that path does attach the slices. An SDK that instead refetches through GetWorkflowExecutionHistory gets none, because that RPC has nowhere to carry them.
  • TestARefusedTruncateProbesThePinsItWasRefusedFor covers a Added the stream service and its frontend wiring. #5 change, and sits here because the harness it needs arrives in Added the workflow publish and subscribe commands. #6.

Review round: api surface

The subscribe commands in the delivery tests name their stream through StreamNameOrId. The service inputs in the same files keep their own StreamId, since those protos did not change. A finish record that used -1 to say it was unnumbered now carries a real sequence, because zero is the unnumbered value.

Brought onto current main (d8f9c6d86b2c) by merge; the biggest adaptation was regenerating the historyservice and matchingservice protos against main, along with a mockgen-generated mutable state mock and the Go 1.27 shapes in the slice reader and onebox.

Merged the activity-owned stream work up from #5 and #6; nothing on this layer changed beyond resolving the merge.

Start positions

Merged the start position work. The wake test that subscribed with start_offset: -1 now asks for start_position.tail.

Records ride the task response and the completion records only the offset range consumed, so replay re-reads the payloads from the stream instead of History carrying them. A range that cannot be served fails the task with its own cause, and a non-sticky query asks History for the same ranges after it has fetched the history.
A discarded speculative task commits no cursor, so the range it was given came
back on the next task while the worker had already consumed it. The replay
budget is dynamic config now, a consumer that stays over it is terminated with
the cause instead of failing the same task forever, and a query only asks
History for ranges when the workflow actually holds a cursor.
…-5-delivery-replay

# Conflicts:
#	service/history/api/recordworkflowtaskstarted/stream_routing.go
A workflow may consume as many streams as the subscription limit allows, and
the execution's lock is held across every read.
The subscribe command names a stream by name or by id, and an unnumbered record no longer rides a sentinel.
…-5-delivery-replay

# Conflicts:
#	api/historyservice/v1/request_response.pb.go
#	api/matchingservice/v1/request_response.pb.go
The ConsumesStreams entry had been hand-placed out of alphabetical order, so
a later mockgen run would have rewritten it.
Main now builds under Go 1.27, where go fix rewrites errors.As and reverse
index loops, so make fmt was left dirty on these two files.
…-5-delivery-replay

# Conflicts:
#	service/history/workflow/mutable_state_impl.go
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant