Repository navigation
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
…-5-delivery-replay
…-5-delivery-replay # Conflicts: # chasm/lib/stream/config.go
…-5-delivery-replay
…n the wire. Paging the re-supply needs a short flag on the poll response or the last slice, which is a wire change; the package doc says so next to the other wire-round notes so the gap is not lost.
…-5-delivery-replay
…-5-delivery-replay
A wake tells a running execution that a source it consumes moved. It folds into a wake no started task has carried, rides the next task, and is dropped when that task completes, so nothing about it enters History.
2 of 3 tasks
1 of 3 tasks
…-5-delivery-replay
The notification channel replaces it. The close hook stays, since the stream consumer's pending data and the channel's pending notifications still schedule a task through it.
…-5-delivery-replay
…-5-delivery-replay
…-5-delivery-replay
…-5-delivery-replay
Owner
Author
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 Workflow Tasks and re-supplies them on replay.
What changed?
recordworkflowtaskstarted/stream_slices.gostages a range per subscription and reads it from an owned stream or, through the routed client, a standalone one. Slices rideRecordWorkflowTaskStartedand the poll response. A duplicate request gets the same slice. Speculative tasks get no range, since a discarded one commits no cursor.WorkflowTaskCompleted. The cursor advance is folded in as that event is built, so both commit together. A subscription with offsets left schedules a task on transaction close, which also covers a fresh subscription to a full stream.stream.replayMaxRecords,stream.replayMaxBytesorstream.replayMaxPagesfails the task withSTREAM_RANGE_UNAVAILABLE. The second failure in a row terminates the workflow, since one retry covers a moving shard and a second isn't a race. There's no paged re-supply and no short flag on the wire yet, aschasm/lib/stream/doc.gonotes.stream.RoutedSetBudget, the deadline the completion path uses, because the execution's lock is held across them.GetStreamReplaySliceslets matching re-supply a non-sticky query task, which skipsRecordWorkflowTaskStarted. Matching calls it only whenGetMutableStateResponsesaysconsumes_streams, so other queries take no history lease.attachReplaySlicesreturns early on a sticky task, so an SDK that refetches history after eviction gets no slices.tests/testcore/onebox.gostarts more than one history host for the cross-host test.tests/sdk_server_host_test.goandtests/stream_publish_cost_test.goare opt-in.TestARefusedTruncateProbesThePinsItWasRefusedFortests a Added the stream service and its frontend wiring. #5 change. It lands here with the consumer harness it needs.Part of AI-198 (epic AI-37).
Why?
This is the half of the design that has to hold under replay. A cold worker sees what the original execution saw, and History never holds a payload.
How did you test it?
go build ./...,make protoandmake fmtare clean. Unit tests ran for./chasm/lib/...,./service/history/...,./service/matching/and./common/.... The functional tests ran with-tags test_dep: delivery, continue-as-new, replay past the first page, a deleted stream, the query path, cross-host routing and speculative tasks.