Conversation
History records the offsets a task consumed and never the payloads, so the server sends the bytes on the poll response: untagged for the task about to run, and tagged with a WorkflowTaskCompleted event id when re-supplying what an earlier task consumed. Core partitions the two and emits recorded ranges in event order before the live one, so a replaying workflow observes them exactly as it did the first time. An empty range is delivered rather than dropped: a task where the subscription saw nothing is a fact replay has to reproduce. The Rust SDKs have no stream API, so they fail loudly on this job instead of ignoring it. The server has already recorded the range as consumed and will not send it again, so dropping it would lose data silently.
The command has a history event, so it fits the matching every SDK's replay depends on: commands are popped from a queue as command-generated events arrive, and one producing no event would put that out of step. The machine never resolves. The event records the subscription and hands nothing back; the ranges arrive later as their own activation jobs.
The reporter panics on any machine name it has no visualizer for, so every state machine has to be listed there.
`complete_workflow_task` takes a shutdown token, so the `mockall` closures need a second parameter.
The matrix asks for 40 minutes on `macos-intel`, but the job never read that value, so the runner stayed on the 25-minute default it exceeds.
The nexgen the Python SDK drives rejects @nexus.type on a native declaration, so generation failed for every consumer. Upstream declares both policies as placeholders, which keeps the annotations valid.
The range a task consumed is recorded on the completion that closes it, and the machines look ahead to that event while replaying the task. An update cut at a WFT started event left the completion on the retained tail or on an unfetched page, so the range arrived one activation late.
A publish is matched on its stream and message count, a subscription on its stream. Replay sends no commands, so a reissued command that differs from the record would otherwise be accepted and the workflow's state would diverge silently.
A completion that records consumed stream cursors ends its task sequence. Such a task ran as one activation live, and folding it into a heartbeat chain would hand several ranges over at once on replay.
Ranges for one task are handed over sorted by stream and offset, so a workflow waiting on several streams sees the same order live and on replay. A re-supplied slice has to cover exactly the recorded offsets, and a range that observed nothing is rebuilt from the cursor rather than demanded from the server.
The external stream family, developed alongside this one, takes 17 to 20 in the activation job and 23 to 28 in the command. Both trees now encode the native stream messages at 21, 29 and 30, so an artifact generated from either decodes correctly against the other.
…lary. The api names an entry a record and carries its kind, producer, attempt and sequence on it, so the vendored stream protos follow the api head and lang appends with AppendStreamRecords and receives DeliverStreamRecords. Core passes each record through untouched.
History records only the offsets a task consumed, so a language replayer that fetched the records from the stream service needs a way to hand them to Core with the history. The replay worker puts them on its synthetic poll response, so the ordinary delivery path and identity checks run unchanged.
The bytes for a recorded range only travel on the response that carries the task, so a sticky task, or a sticky legacy query, handed to a worker that no longer holds the run and fetched the history itself cannot be replayed. Failing when the lookahead sees the range keeps the workflow off less input than it had, and lets a legacy query go unanswered so the server retries it on the normal queue, where the records travel with it.
A recorded range with content and no bytes for it is the worker's failure to reconstruct the run, not the workflow's nondeterminism, so it is its own error kind and treated like a failed history fetch: the task fails as an unhandled worker failure and a legacy query goes unanswered, so the server retries both where the records travel. The broken run is evicted even when the query failure is withheld, so the retry starts from history.
`api_upstream` carried an older cut of the stream additions, so regenerating `temporalio/api` from this Core removed what the api branch had since moved and documented. The three files are the api branch's diff applied onto the api tree this Core already pins.
History records the offsets a task consumed and never the payloads, so the server sends the bytes on the poll response: untagged for the task about to run, and tagged with a WorkflowTaskCompleted event id when re-supplying what an earlier task consumed. Core partitions the two and emits recorded ranges in event order before the live one, so a replaying workflow observes them exactly as it did the first time. A range a task consumed is recorded on the completion that closes it, so the machines look ahead to that event while replaying the task, and the paginator keeps that completion reachable when a page boundary falls between the two.
A completion that records consumed stream cursors ends its task sequence. Such a task ran as one activation live, and folding it into a heartbeat chain would hand several ranges over at once on replay.
A recorded range with content and no bytes for it is the worker's failure to reconstruct the run, not the workflow's nondeterminism, so the task fails as an unhandled worker failure and a legacy query goes unanswered, and the server retries both where the records travel. The broken run is evicted even when the query failure is withheld, so the retry starts from history.
History records only the offsets a task consumed, so a language replayer that fetched the records from the stream service needs a way to hand them to Core with the history. The replay worker puts them on its synthetic poll response, so the ordinary delivery path and identity checks run unchanged.
The guards are what the delivery path is for: an empty range still arrives, a re-supplied slice has to cover exactly the recorded offsets, a cached run is not handed its own range twice, and a task owed records the response did not carry fails before the workflow runs on less input.
The merge adds no changes; it gives the chain the original head as an ancestor so downstream pins and merges reconcile.
This was referenced Sep 25, 2026
2 of 3 tasks
# Conflicts: # crates/sdk-core/src/core_tests/streams.rs
A task owed two streams and sent one came out as nondeterminism, which a worker configured for it fails the execution over. Both sides of that comparison come from the server, so no arm of it is the workflow's fault.
The flag that says a completion consumed a range was written into the two that carry across the whole scan, which only works while every path after it returns. A local one says what it means on its own.
The fetch costs a page for any run whose page boundary lands there. A subscription made through the stream service records no event, so there is no telling a stream run from any other before the ranges arrive.
The delivery tests build their own commands and histories, so they name the fields directly. One recorded append passed a count where the event now wants an end offset.
# Conflicts: # crates/sdk-core/src/worker/workflow/managed_run.rs
The main merge introduced the released headings above them and the resolution left the entries inside 0.8.0. The checkpoint requires additions under Unreleased.
# Conflicts: # crates/sdk-core/src/core_tests/streams.rs
The server records an event for every unsubscribe, also for a channel the run never subscribed to, so the machine matches by channel alone. The notification job stays as it is, since Core keeps no view of the run's subscriptions and lang drops a late notification for a closed one.
3 tasks
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 server-side stream ranges to a workflow, live and on replay.
What changed?
PreparedWFTand the run, into the machines, and out asDeliverStreamRecordsjobs. An empty range is delivered, not dropped.WorkflowTaskCompletedevent id re-supply what an earlier task consumed. Core emits recorded ranges in event order before the live one, sorted by stream and offset within a task.history_update.rs: a page that ends on a task started event fetches the next page, so the range recorded on that task's completion isn't found a task late.MissingRecords: History records a range with content and the response carried none. The task fails as a worker failure, and a sticky legacy query goes unanswered. The server retries both on the normal queue, where the records travel. A partial re-supply fails the same way.HistoryForReplay::with_stream_slicesputs fetched records on the replay worker's synthetic poll response. A cached run doesn't take its own range again from a re-supply.NotificationsReceivedcomes from thenotificationson the scheduled events of a task, folded per channel with the highest counter kept, live and on replay, ordered with signals and ahead of the stream ranges. A failed task's notifications reach lang with its retry's in that one job. A task scheduled with notifications ends a heartbeat chain inhistory_update.rs, so it keeps its own replay activation.UnsubscribeNotificationChannelon the command machine next to the subscribe: sent on the completion, held to itsWorkflowNotificationChannelUnsubscribedevent by channel on replay, and a different channel fails the task as nondeterminism. TheNotificationsReceivedjob is unchanged, so a notification the server already put on a scheduled event still reaches lang, which drops it for a closed subscription.This is the top of the native chain. #4 replays this content onto the unified Core.
Part of AI-198 (epic AI-37).
Why?
History records the offsets a task consumed, never the payloads, so replay has to get the same bytes from the server and hand them over at the same point. A range that arrives late, twice or short is a silent divergence that surfaces later as an unrelated nondeterminism error, and each case has its own guard and test. A mismatch between History and the poll response is a server-side fact, so it fails as the worker's failure, not as nondeterminism.
How did you test it?
cargo check --workspace,cargo test -p temporalio-sdk-core --lib core_tests::streams, the wholecargo test -p temporalio-sdk-core --lib,cargo fmt --all --checkandcargo test-lint. The page-boundary case is anrstestover both boundary positions.