Conversation
2 of 3 tasks
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.
moedash
force-pushed
the
moe/AI-198-fix-replay-slice-lookahead
branch
from
September 16, 2026 22:06
67ffc6d to
865805b
Compare
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.
This was referenced Sep 25, 2026
Owner
Author
|
Replaced by #5, #6 and #7, which split this branch's content into the proto layer, the two command machines, and the delivery path, each one compiling and passing its own tests on its own. The tip of #7 merges |
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.
What changed?
This PR adds server-side stream support to Core: the command state machines, the activation job that hands a consumed range to workflow code, and the replay repair that re-supplies recorded ranges before command matching. A recorded range now arrives in the same activation as its task across History page boundaries, and a replayed append or subscribe is checked against the recorded event rather than matched by type alone. The wire follows the api's record vocabulary: lang appends with
AppendStreamRecords, consumed ranges arrive asDeliverStreamRecordsjobs, and eachStreamRecordcarries its kind, producer id, attempt and sequence through Core untouched.A history fed to the replay worker can carry the stream records its tasks consumed (
HistoryForReplay::with_stream_slices). The replay worker puts them on its synthetic poll response, so the same delivery path and identity checks run as on a live cache miss, and a language replayer that fetched the records from the stream service can replay a consuming workflow.A task whose history records a consumed range with content that the response carried no records for now fails when the lookahead sees the range, before the workflow runs on less input than it had. That is the sticky task, or the sticky legacy query, handed to a worker that no longer holds the run and fetched the history itself. The failure is its own kind,
WFMachinesError::MissingRecords, treated like a failed history fetch: the task fails as an unhandled worker failure, never as nondeterminism, and a legacy query goes unanswered rather than answered from the wrong state, so the server's sticky attempt times out and the query is retried on the normal queue, where the records travel with it. The broken run is evicted even when the query failure is withheld.Fixes AI-198
Why?
Core matches commands to events positionally, so stream commands and range delivery need real state machines. Without the replay repair a consumer replays without the data it originally saw. History alone holds only the offsets a task consumed, so a replayer needs a way to hand Core the records with the history, and a worker that fetched a history itself has no way to get them at all.
How did you test it?
Link to a test plan if any -
The lib tests cover the stream machines and replay, the replay worker fed a read-then-publish history with its slices, with a wrong slice, and with none, and the missing re-supply failing before the first activation. Live coverage comes from the Python native provider running against a server built from the stream branch, including a sticky query against a worker that evicted the run.