Conversation
The poll response carries the records for the task about to run and re-supplies what earlier tasks consumed, keyed by the completion that recorded each range. Core hands them over in a fixed order, gives a consuming task its own activation on replay, and fails a task whose recorded range arrived without its records.
The cases cover live delivery, replay across history pages, the lookahead to the closing completion, missing and mismatched records, and notifications ordered ahead of a stream range.
This was referenced Oct 3, 2026
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 hands a workflow the stream ranges it consumed, both live and on replay.
What changed?
DeliverStreamRecordsjobs in a fixed order by stream and offset, after the channel notifications.MissingRecordsbefore the workflow runs. A legacy query sent that way goes unanswered, so the server retries it on the normal task queue, where the records travel with it.HistoryForReplay::with_stream_sliceslets a language replayer pass the records it fetched from the stream service.Part of AI-198 (epic AI-37).
Why?
History holds only the offsets a task consumed, so replay depends on the server sending the bytes back. Replaying a task on less input than it ran on produces different commands, and that would surface later as an unrelated nondeterminism error. Failing early with a clear cause is better than failing late. The order ranges arrive in is something the workflow can branch on, so Core fixes it instead of taking whatever order the server produced.
How did you test it?
Link to a test plan if any -
core_tests::streamscovers live delivery, empty ranges, replay order, the lookahead, a page boundary, data-only tasks, a cached run, missing and mismatched records, the legacy query and pushed histories. A channel case checks that notifications come ahead of a stream range in the same activation. Every workspace lib suite passes, the lints and fmt are clean, and each commit builds.