Skip to content

Delivered server-side stream ranges to workflow activations. - #1

Closed
moedash wants to merge 20 commits into
mainfrom
moe/AI-198-fix-replay-slice-lookahead
Closed

moedash wants to merge 20 commits into
mainfrom
moe/AI-198-fix-replay-slice-lookahead

Conversation

@moedash

@moedash moedash commented Sep 14, 2026 •

Copy link
Copy Markdown
Owner

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 as DeliverStreamRecords jobs, and each StreamRecord carries 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 -

  • Unit Tests
  • Staging
  • End to End Tests

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.

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
moedash force-pushed the moe/AI-198-fix-replay-slice-lookahead branch from 67ffc6d to 865805b Compare September 16, 2026 22:06
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.
@moedash

moedash commented Sep 25, 2026

Copy link
Copy Markdown
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 32c66693f072, this branch's head, with no content of its own, so the pins and the downstream PRs that depend on this head keep reconciling; the branch stays.

@moedash moedash closed this Sep 25, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant