Conversation
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.
This was referenced Sep 25, 2026
…-5-delivery-replay
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
…-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.
…-5-delivery-replay
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.
…-5-delivery-replay
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
…-5-delivery-replay # Conflicts: # service/history/workflow/mutable_state_impl.go
…-5-delivery-replay
…-5-delivery-replay
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This PR delivers stream records on the workflow task and re-supplies them on replay.
A subscribed workflow gets its records as slices on
RecordWorkflowTaskStartedand on the poll response. The task records only the offset range it consumed, onWorkflowTaskCompleted. 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.chasm/lib/workflow/workflow.goandservice/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.GetStreamReplaySliceson 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 withoutRecordWorkflowTaskStarted, so nothing else re-supplies its ranges and a cold worker could not replay a consuming workflow to answer.tests/testcore/onebox.gocan start more than one history host, which the cross-host routing test needs.tests/sdk_server_host_test.gois an opt-in host that serves a real SDK against this build.tests/stream_publish_cost_test.gois 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?
All green.
Series
Server 5 of 6. Replaces #2.
Previous: #6
moe/AI-198-srv-4-workflow-commands. Next: #8moe/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,
TestSpeculativeTaskConsumesNoStreamRangealso 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.replayMaxBytesandstream.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 consecutiveSTREAM_RANGE_UNAVAILABLEterminates 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.RoutedSetBudgetthe 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.
GetMutableStateResponsecarriesconsumes_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:
attachReplaySlicesstill 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 throughGetWorkflowExecutionHistorygets none, because that RPC has nowhere to carry them.TestARefusedTruncateProbesThePinsItWasRefusedForcovers 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 ownStreamId, 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 thehistoryserviceandmatchingserviceprotos 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: -1now asks forstart_position.tail.