Conversation
The external stream family, developed alongside this one, takes 23 to 28 in the command message. Both trees encode the native stream commands at 29 and 30, so an artifact generated from either decodes correctly against the other.
Each 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. Neither machine resolves, because the event records what happened and hands nothing back; the ranges a subscription brings arrive later as their own activation jobs. A reissued command is held against the recorded event, a publish on its stream and record count and a subscription on its stream, since replay sends no commands and a divergence would otherwise pass unnoticed.
The reporter panics on any machine name it has no visualizer for, so every state machine has to be listed there.
The events the two commands produce are added to the test history builder, so a recorded command can be reissued against them and the check that holds it to the record is exercised.
This was referenced Sep 25, 2026
A subscription's explicit start offset and an unnamed append's resolved stream both went unchecked, so a replay that asked for something else passed. The offset check skips a repeat subscribe, which the server records at the cursor rather than at what the command asked for.
Both command enums are uninhabited, so neither arm runs. Nondeterminism is the wrong label for an internal invariant when the rest of the series works to keep worker failures out of that bucket.
Comparing it is only right for a run's first subscribe to a stream, and a subscription made through the stream service leaves no event, so which one is first cannot be told. Failing a sound run costs more than the drift.
The lang command now says stream_name for an append and stream_name_or_id for a subscribe, so one vocabulary runs from lang through to History. The append machine reads the batch size off the event's offset range.
Lang can now ask for the earliest record, the tail or the last N instead of a negative offset. Core only forwards it: the server resolves it and records the offset, which is all replay matches against.
# Conflicts: # crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto
# Conflicts: # crates/protos/src/protos/mod.rs # crates/sdk-core/src/worker/workflow/mod.rs
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 lets a workflow subscribe to a server-side stream and append records to one.
What changed?
SubscribeStreamandAppendStreamRecordson the lang command message at 29 and 30. The external family takes 23 to 28, so both trees encode the native commands at the same numbers.stream_nameon an append,stream_name_or_idandstart_positionon a subscribe. TheFromconversion passes the position through untouched.append_stream_records_state_machine.rsandsubscribe_stream_state_machine.rs. Neither resolves. A subscription's ranges arrive later as their own jobs, and that's Delivered server-side stream ranges to workflows. #7.adapt_responsearms sayfatal!.SubscribeNotificationChannelon the lang command message at 31 andNotificationsReceivedon the activation job at 22.subscribe_notification_channel_state_machine.rsholds a reissued subscribe to its recorded event by channel, and the Rust SDKs ignore the job since they cannot subscribe and a notification carries no data.UnsubscribeNotificationChannelis still refused on this layer. Its machine lands with Delivered server-side stream ranges to workflows. #7, next to the delivery it is paired with.Part of AI-198 (epic AI-37).
Why?
Both are ordinary commands with a history event, so they fit the command matching replay relies on. The server resolves a subscription's addressing and start, since the workflow can't look them up without I/O, so the replay check doesn't compare the recorded start offset. A subscription made through the stream service leaves no event, so history can't tell a run's first subscribe, and a test pins that non-check.
How did you test it?
cargo check --workspace,cargo test -p temporalio-sdk-core --lib core_tests::streams,cargo fmt --all --checkandcargo test-lint. The stream tests cover both commands reaching the server as lang asked, each start position arm, both round-tripping through replay, and a reissue to another stream or with another batch size failing the task.