Skip to content

Let a workflow subscribe to a stream and append records to it. - #6

Closed
moedash wants to merge 25 commits into
moe/AI-198-core-1-protosfrom
moe/AI-198-core-2-machines
Closed

moedash wants to merge 25 commits into
moe/AI-198-core-1-protosfrom
moe/AI-198-core-2-machines

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR lets a workflow subscribe to a server-side stream and append records to one.

What changed?

  • SubscribeStream and AppendStreamRecords on 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.
  • The lang fields match the wire: stream_name on an append, stream_name_or_id and start_position on a subscribe. The From conversion passes the position through untouched.
  • Two command machines, append_stream_records_state_machine.rs and subscribe_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.
  • Record bodies travel on the command. History gets one fixed-size event naming the offset range, so batch size doesn't grow History.
  • A reissued command is held against its recorded event. An append matches on stream and record count, and an unnamed append on the stream the run's earlier unnamed appends resolved to. A subscribe matches on its stream only.
  • Both machines are registered with the transition coverage reporter, and the unreachable adapt_response arms say fatal!.
  • SubscribeNotificationChannel on the lang command message at 31 and NotificationsReceived on the activation job at 22. subscribe_notification_channel_state_machine.rs holds 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.
  • UnsubscribeNotificationChannel is 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 --check and cargo 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.

  • Unit Tests
  • Staging
  • End to End Tests

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.
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.
@moedash moedash added the skip-changelog Changelog entry rides another PR label Sep 26, 2026
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #8, #9, #10, #11, #12, #13, #14, #15, #16, #17, #18, #19, #20, #21, #22 and #23.

Same content, split into 16 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

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

Labels

skip-changelog Changelog entry rides another PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant