From 4e7910bca52e34162ec46b6bdfd195528e28aef6 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 08:07:47 -0700 Subject: [PATCH 01/10] Added the stream commands to the bridge protos. 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. --- .../workflow_commands/workflow_commands.proto | 32 +++++++++++++++ crates/protos/src/protos/mod.rs | 39 +++++++++++++++++++ 2 files changed, 71 insertions(+) diff --git a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto index 03e172216..456055d6c 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto @@ -15,6 +15,7 @@ import "temporal/api/common/v1/message.proto"; import "temporal/api/enums/v1/workflow.proto"; import "temporal/api/failure/v1/message.proto"; import "temporal/api/sdk/v1/user_metadata.proto"; +import "temporal/api/stream/v1/message.proto"; import "temporal/api/sdk/v1/event_group_marker.proto"; import "temporal/sdk/core/child_workflow/child_workflow.proto"; import "temporal/sdk/core/nexus/nexus.proto"; @@ -54,9 +55,40 @@ message WorkflowCommand { UpdateResponse update_response = 20; ScheduleNexusOperation schedule_nexus_operation = 21; RequestCancelNexusOperation request_cancel_nexus_operation = 22; + // 23 to 28 are taken by the external stream commands, which are developed + // alongside these and share this message. The two numbers below are fixed + // with that family and must not be reused. + SubscribeStream subscribe_stream = 29; + AppendStreamRecords append_stream_records = 30; } } +// Append a batch of records to a stream this workflow owns. +// +// The bodies go to the stream's own log rather than into History, which gets +// one fixed-size event naming the offset range. That is what makes the batch +// size free: a thousand records cost the same in History as one. The server +// stores each record with an empty producer id, because the workflow is the +// producer here. +message AppendStreamRecords { + // Empty means the workflow's default output stream. + string stream_id = 1; + repeated temporal.api.stream.v1.StreamRecord records = 2; +} + +// Subscribe this workflow to a stream, so later Workflow Tasks carry the +// ranges it has not consumed yet. +// +// The stream's addressing is resolved by the server. A workflow cannot look it +// up without doing I/O, and a value it carried would be a reading rather than a +// fact, so it could differ on replay. +message SubscribeStream { + string stream_id = 1; + // Negative means from wherever the stream is when the subscription is + // registered. The server resolves that once and records it. + int64 start_offset = 2; +} + message StartTimer { // Lang's incremental sequence number, used as the operation identifier uint32 seq = 1; diff --git a/crates/protos/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index 4af74a5f4..3cf466a9a 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1493,6 +1493,23 @@ pub mod coresdk { } } + impl Display for SubscribeStream { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "SubscribeStream({})", self.stream_id) + } + } + + impl Display for AppendStreamRecords { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!( + f, + "AppendStreamRecords({}, {} records)", + self.stream_id, + self.records.len() + ) + } + } + impl Display for StartTimer { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "StartTimer({})", self.seq) @@ -1891,6 +1908,28 @@ pub mod temporal { } } + impl From for command::Attributes { + fn from(s: workflow_commands::AppendStreamRecords) -> Self { + Self::AppendStreamRecordsCommandAttributes( + AppendStreamRecordsCommandAttributes { + stream_id: s.stream_id, + records: s.records, + }, + ) + } + } + + impl From for command::Attributes { + fn from(s: workflow_commands::SubscribeStream) -> Self { + Self::SubscribeStreamCommandAttributes( + SubscribeStreamCommandAttributes { + stream_id: s.stream_id, + start_offset: s.start_offset, + }, + ) + } + } + impl From for command::Attributes { fn from(s: workflow_commands::StartTimer) -> Self { Self::StartTimerCommandAttributes(StartTimerCommandAttributes { From e3e429ce660df8f626d4f720e6db0c0b27e3cddf Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 08:07:58 -0700 Subject: [PATCH 02/10] Let a workflow subscribe to a stream and append records to it. 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. --- .../append_stream_records_state_machine.rs | 144 ++++++++++++++++++ .../src/worker/workflow/machines/mod.rs | 6 + .../subscribe_stream_state_machine.rs | 136 +++++++++++++++++ .../workflow/machines/workflow_machines.rs | 24 ++- crates/sdk-core/src/worker/workflow/mod.rs | 6 + 5 files changed, 315 insertions(+), 1 deletion(-) create mode 100644 crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs create mode 100644 crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs diff --git a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs new file mode 100644 index 000000000..5acdc7eda --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs @@ -0,0 +1,144 @@ +use super::{ + NewMachineWithCommand, StateMachine, TransitionResult, fsm, workflow_machines::MachineResponse, +}; +use crate::worker::workflow::{ + WFMachinesError, fatal, + machines::{EventInfo, HistEventData, WFMachinesAdapter}, + nondeterminism, +}; +use temporalio_common::protos::{ + coresdk::workflow_commands::AppendStreamRecords, + temporal::api::{ + enums::v1::{CommandType, EventType}, + history::v1::{WorkflowStreamRecordsAppendedEventAttributes, history_event}, + }, +}; + +fsm! { + pub(super) name AppendStreamRecordsMachine; + command AppendStreamRecordsMachineCommand; + error WFMachinesError; + shared_state SharedState; + + Created --(CommandScheduled) --> CommandIssued; + CommandIssued --(CommandRecorded(WorkflowStreamRecordsAppendedEventAttributes), + shared on_command_recorded) --> Done; +} + +/// What the command claimed, kept so the recorded event can be held against it. +#[derive(Default, Clone)] +pub(super) struct SharedState { + stream_id: String, + record_count: i64, +} + +/// Append a batch of records to a stream this workflow owns. +/// +/// The bodies go to the stream's own log, and History gets one event naming the +/// offset range the batch landed at. The offsets are assigned by the server, so +/// nothing here predicts them. +pub(super) fn append_stream_records(lang_cmd: AppendStreamRecords) -> NewMachineWithCommand { + let sm = AppendStreamRecordsMachine::from_parts( + Created {}.into(), + SharedState { + stream_id: lang_cmd.stream_id.clone(), + record_count: lang_cmd.records.len() as i64, + }, + ); + NewMachineWithCommand { + command: lang_cmd.into(), + machine: sm.into(), + } +} + +#[derive(Debug, derive_more::Display)] +pub(super) enum AppendStreamRecordsMachineCommand {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Created {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct CommandIssued {} + +impl CommandIssued { + pub(super) fn on_command_recorded( + self, + dat: &mut SharedState, + attrs: WorkflowStreamRecordsAppendedEventAttributes, + ) -> AppendStreamRecordsMachineTransition { + // An empty id names the workflow's default stream, and the server is the + // one that resolves that name, so only a named stream can be compared. + let same_stream = dat.stream_id.is_empty() || dat.stream_id == attrs.stream_id; + if same_stream && dat.record_count == attrs.record_count { + TransitionResult::default() + } else { + TransitionResult::Err(nondeterminism!( + "Recorded append of {} records to stream {:?} does not match the reissued \ + append of {} records to stream {:?}", + attrs.record_count, + attrs.stream_id, + dat.record_count, + dat.stream_id + )) + } + } +} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Done {} + +impl WFMachinesAdapter for AppendStreamRecordsMachine { + fn adapt_response( + &self, + _my_command: Self::Command, + _event_info: Option, + ) -> Result, Self::Error> { + Err(Self::Error::Nondeterminism( + "AppendStreamRecords does not use state machine commands".to_string(), + )) + } +} + +impl TryFrom for AppendStreamRecordsMachineEvents { + type Error = WFMachinesError; + + fn try_from(e: HistEventData) -> Result { + let e = e.event; + match e.event_type() { + EventType::WorkflowStreamRecordsAppended => { + if let Some( + history_event::Attributes::WorkflowStreamRecordsAppendedEventAttributes(attrs), + ) = e.attributes + { + Ok(AppendStreamRecordsMachineEvents::CommandRecorded(attrs)) + } else { + Err(fatal!("Stream records appended attributes were unset: {e}")) + } + } + _ => Err(Self::Error::Nondeterminism(format!( + "AppendStreamRecordsMachine does not handle {e}" + ))), + } + } +} + +impl TryFrom for AppendStreamRecordsMachineEvents { + type Error = WFMachinesError; + + fn try_from(c: CommandType) -> Result { + match c { + CommandType::AppendStreamRecords => { + Ok(AppendStreamRecordsMachineEvents::CommandScheduled) + } + _ => Err(Self::Error::Nondeterminism(format!( + "AppendStreamRecordsMachine does not handle command type {c:?}" + ))), + } + } +} + +impl From for CommandIssued { + fn from(_: Created) -> Self { + Self {} + } +} diff --git a/crates/sdk-core/src/worker/workflow/machines/mod.rs b/crates/sdk-core/src/worker/workflow/machines/mod.rs index bc5e3fcb9..364897465 100644 --- a/crates/sdk-core/src/worker/workflow/machines/mod.rs +++ b/crates/sdk-core/src/worker/workflow/machines/mod.rs @@ -1,3 +1,4 @@ +mod append_stream_records_state_machine; mod workflow_machines; mod activity_state_machine; @@ -15,6 +16,7 @@ mod modify_workflow_properties_state_machine; mod nexus_operation_state_machine; mod patch_state_machine; mod signal_external_state_machine; +mod subscribe_stream_state_machine; mod timer_state_machine; mod update_state_machine; mod upsert_search_attributes_state_machine; @@ -31,6 +33,7 @@ use crate::{ worker::workflow::{WFMachinesError, fatal, nondeterminism}, }; use activity_state_machine::ActivityMachine; +use append_stream_records_state_machine::AppendStreamRecordsMachine; use cancel_external_state_machine::CancelExternalMachine; use cancel_workflow_state_machine::CancelWorkflowMachine; use child_workflow_state_machine::ChildWorkflowMachine; @@ -46,6 +49,7 @@ use std::{ convert::{TryFrom, TryInto}, fmt::{Debug, Display}, }; +use subscribe_stream_state_machine::SubscribeStreamMachine; use temporalio_common::{ fsm_trait::{StateMachine, TransitionResult}, protos::temporal::api::{ @@ -80,6 +84,8 @@ enum Machines { WorkflowTaskMachine, UpsertSearchAttributesMachine, ModifyWorkflowPropertiesMachine, + SubscribeStreamMachine, + AppendStreamRecordsMachine, UpdateMachine, NexusOperationMachine, } diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs new file mode 100644 index 000000000..3fc08e768 --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -0,0 +1,136 @@ +use super::{ + NewMachineWithCommand, StateMachine, TransitionResult, fsm, workflow_machines::MachineResponse, +}; +use crate::worker::workflow::{ + WFMachinesError, fatal, + machines::{EventInfo, HistEventData, WFMachinesAdapter}, + nondeterminism, +}; +use temporalio_common::protos::{ + coresdk::workflow_commands::SubscribeStream, + temporal::api::{ + enums::v1::{CommandType, EventType}, + history::v1::{WorkflowStreamSubscribedEventAttributes, history_event}, + }, +}; + +fsm! { + pub(super) name SubscribeStreamMachine; + command SubscribeStreamMachineCommand; + error WFMachinesError; + shared_state SharedState; + + Created --(CommandScheduled) --> CommandIssued; + CommandIssued --(CommandRecorded(WorkflowStreamSubscribedEventAttributes), + shared on_command_recorded) --> Done; +} + +/// The stream the command named, kept so the recorded event can be held against it. +/// The start offset is not kept: the server resolves it, so the recorded value is +/// its answer rather than what the command said. +#[derive(Default, Clone)] +pub(super) struct SharedState { + stream_id: String, +} + +/// Subscribe this workflow to a stream. The command carries only the stream id +/// and a start offset; the server resolves the addressing, because a workflow +/// cannot look it up without doing I/O and a value it carried would be a +/// reading rather than a fact. +pub(super) fn subscribe_stream(lang_cmd: SubscribeStream) -> NewMachineWithCommand { + let sm = SubscribeStreamMachine::from_parts( + Created {}.into(), + SharedState { + stream_id: lang_cmd.stream_id.clone(), + }, + ); + NewMachineWithCommand { + command: lang_cmd.into(), + machine: sm.into(), + } +} + +#[derive(Debug, derive_more::Display)] +pub(super) enum SubscribeStreamMachineCommand {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Created {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct CommandIssued {} + +impl CommandIssued { + pub(super) fn on_command_recorded( + self, + dat: &mut SharedState, + attrs: WorkflowStreamSubscribedEventAttributes, + ) -> SubscribeStreamMachineTransition { + if dat.stream_id == attrs.stream_id { + TransitionResult::default() + } else { + TransitionResult::Err(nondeterminism!( + "Recorded subscription to stream {:?} does not match the reissued subscription \ + to stream {:?}", + attrs.stream_id, + dat.stream_id + )) + } + } +} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Done {} + +impl WFMachinesAdapter for SubscribeStreamMachine { + fn adapt_response( + &self, + _my_command: Self::Command, + _event_info: Option, + ) -> Result, Self::Error> { + Err(Self::Error::Nondeterminism( + "SubscribeStream does not use state machine commands".to_string(), + )) + } +} + +impl TryFrom for SubscribeStreamMachineEvents { + type Error = WFMachinesError; + + fn try_from(e: HistEventData) -> Result { + let e = e.event; + match e.event_type() { + EventType::WorkflowStreamSubscribed => { + if let Some(history_event::Attributes::WorkflowStreamSubscribedEventAttributes( + attrs, + )) = e.attributes + { + Ok(SubscribeStreamMachineEvents::CommandRecorded(attrs)) + } else { + Err(fatal!("Stream subscribed attributes were unset: {e}")) + } + } + _ => Err(Self::Error::Nondeterminism(format!( + "SubscribeStreamMachine does not handle {e}" + ))), + } + } +} + +impl TryFrom for SubscribeStreamMachineEvents { + type Error = WFMachinesError; + + fn try_from(c: CommandType) -> Result { + match c { + CommandType::SubscribeStream => Ok(SubscribeStreamMachineEvents::CommandScheduled), + _ => Err(Self::Error::Nondeterminism(format!( + "SubscribeStreamMachine does not handle command type {c:?}" + ))), + } + } +} + +impl From for CommandIssued { + fn from(_: Created) -> Self { + Self {} + } +} diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index 8d4b8f54e..7648aed81 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -2,13 +2,15 @@ mod local_acts; use super::{ Machines, NewMachineWithCommand, TemporalStateMachine, + append_stream_records_state_machine::append_stream_records, cancel_external_state_machine::new_external_cancel, cancel_workflow_state_machine::cancel_workflow, complete_workflow_state_machine::complete_workflow, continue_as_new_workflow_state_machine::continue_as_new, fail_workflow_state_machine::fail_workflow, local_activity_state_machine::new_local_activity, patch_state_machine::has_change, signal_external_state_machine::new_external_signal, - timer_state_machine::new_timer, upsert_search_attributes_state_machine::upsert_search_attrs, + subscribe_stream_state_machine::subscribe_stream, timer_state_machine::new_timer, + upsert_search_attributes_state_machine::upsert_search_attrs, workflow_machines::local_acts::LocalActivityData, workflow_task_state_machine::WorkflowTaskMachine, }; @@ -1526,6 +1528,26 @@ impl WorkflowMachines { CommandIdKind::NeverResolves, ); } + WFCommandVariant::AppendStreamRecords(attrs) => { + // Never resolves: the event names the offset range the + // server assigned and hands nothing back. A workflow that + // wants to know where its batch landed reads the stream. + self.add_cmd_to_wf_task( + append_stream_records(attrs), + annotations, + CommandIdKind::NeverResolves, + ); + } + WFCommandVariant::SubscribeStream(attrs) => { + // Never resolves: the event it produces records the + // subscription and hands nothing back to the workflow. The + // ranges arrive later as their own activation jobs. + self.add_cmd_to_wf_task( + subscribe_stream(attrs), + annotations, + CommandIdKind::NeverResolves, + ); + } WFCommandVariant::UpdateResponse(ur) => { let m_key = self.get_machine_by_msg(&ur.protocol_instance_id)?; let m = if let Machines::UpdateMachine(m) = self.machine_mut(m_key) { diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index b0da392a3..c9f1dda4b 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -1527,6 +1527,8 @@ enum WFCommandVariant { UpdateResponse(UpdateResponse), ScheduleNexusOperation(ScheduleNexusOperation), RequestCancelNexusOperation(RequestCancelNexusOperation), + SubscribeStream(SubscribeStream), + AppendStreamRecords(AppendStreamRecords), } impl TryFrom for WFCommand { @@ -1535,6 +1537,10 @@ impl TryFrom for WFCommand { fn try_from(c: WorkflowCommand) -> result::Result { let variant = match c.variant.ok_or(EmptyWorkflowCommandErr)? { workflow_command::Variant::StartTimer(s) => WFCommandVariant::AddTimer(s), + workflow_command::Variant::SubscribeStream(s) => WFCommandVariant::SubscribeStream(s), + workflow_command::Variant::AppendStreamRecords(s) => { + WFCommandVariant::AppendStreamRecords(s) + } workflow_command::Variant::CancelTimer(s) => WFCommandVariant::CancelTimer(s), workflow_command::Variant::ScheduleActivity(s) => WFCommandVariant::AddActivity(s), workflow_command::Variant::RequestCancelActivity(s) => { From e0a6d1533a3dc091f7c752c5ca7bca9e8e1af29a Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 08:07:58 -0700 Subject: [PATCH 03/10] Registered the stream machines with the coverage reporter. The reporter panics on any machine name it has no visualizer for, so every state machine has to be listed there. --- .../worker/workflow/machines/transition_coverage.rs | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs index b266137fa..41c054450 100644 --- a/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs +++ b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs @@ -66,6 +66,7 @@ mod machine_coverage_report { use super::*; use crate::worker::workflow::machines::{ StateMachine, activity_state_machine::ActivityMachine, + append_stream_records_state_machine::AppendStreamRecordsMachine, cancel_external_state_machine::CancelExternalMachine, cancel_workflow_state_machine::CancelWorkflowMachine, child_workflow_state_machine::ChildWorkflowMachine, @@ -75,7 +76,8 @@ mod machine_coverage_report { local_activity_state_machine::LocalActivityMachine, modify_workflow_properties_state_machine::ModifyWorkflowPropertiesMachine, nexus_operation_state_machine::NexusOperationMachine, patch_state_machine::PatchMachine, - signal_external_state_machine::SignalExternalMachine, timer_state_machine::TimerMachine, + signal_external_state_machine::SignalExternalMachine, + subscribe_stream_state_machine::SubscribeStreamMachine, timer_state_machine::TimerMachine, update_state_machine::UpdateMachine, upsert_search_attributes_state_machine::UpsertSearchAttributesMachine, workflow_task_state_machine::WorkflowTaskMachine, @@ -117,6 +119,8 @@ mod machine_coverage_report { let mut modify_wf_props = ModifyWorkflowPropertiesMachine::visualizer().to_owned(); let mut update = UpdateMachine::visualizer().to_owned(); let mut nexus = NexusOperationMachine::visualizer().to_owned(); + let mut subscribe_stream = SubscribeStreamMachine::visualizer().to_owned(); + let mut append_stream_records = AppendStreamRecordsMachine::visualizer().to_owned(); // This isn't at all efficient but doesn't need to be. // Replace transitions in the vizzes with green color if they are covered. @@ -144,6 +148,12 @@ mod machine_coverage_report { } m @ "UpdateMachine" => cover_transitions(m, &mut update, coverage), m @ "NexusOperationMachine" => cover_transitions(m, &mut nexus, coverage), + m @ "SubscribeStreamMachine" => { + cover_transitions(m, &mut subscribe_stream, coverage) + } + m @ "AppendStreamRecordsMachine" => { + cover_transitions(m, &mut append_stream_records, coverage) + } m => panic!("Unknown machine {m}"), } } From 37500766fee0cd4d7689b962c842b6a94f3e7e83 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 08:07:58 -0700 Subject: [PATCH 04/10] Covered the stream commands and their replay checks. 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. --- crates/sdk-core/src/core_tests/mod.rs | 1 + crates/sdk-core/src/core_tests/streams.rs | 330 ++++++++++++++++++ crates/sdk-core/src/replay/history_builder.rs | 26 ++ 3 files changed, 357 insertions(+) create mode 100644 crates/sdk-core/src/core_tests/streams.rs diff --git a/crates/sdk-core/src/core_tests/mod.rs b/crates/sdk-core/src/core_tests/mod.rs index 6dd4479bb..4779016bb 100644 --- a/crates/sdk-core/src/core_tests/mod.rs +++ b/crates/sdk-core/src/core_tests/mod.rs @@ -2,6 +2,7 @@ mod activity_tasks; mod event_groups; mod queries; mod replay_flag; +mod streams; mod updates; mod workers; mod workflow_cancels; diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs new file mode 100644 index 000000000..0deefe2d4 --- /dev/null +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -0,0 +1,330 @@ +use crate::{ + Worker, + replay::TestHistoryBuilder, + test_help::{MockPollCfg, ResponseType, WorkerTestHelpers, build_mock_pollers, mock_worker}, + worker::client::mocks::mock_worker_client, +}; +use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, +}; +use temporalio_common::protos::{ + coresdk::{ + workflow_commands::{AppendStreamRecords, SubscribeStream}, + workflow_completion::WorkflowActivationCompletion, + }, + temporal::api::{ + command::v1::command, + enums::v1::{CommandType, EventType, WorkflowTaskFailedCause}, + stream::v1::StreamRecord, + workflowservice::v1::RespondWorkflowTaskCompletedResponse, + }, +}; + +/// A worker served one poll response, expected to fail that task as +/// nondeterministic. Returns the worker and the count of failures it reported, +/// which the test asserts itself: the mock only verifies call counts when it +/// is dropped, and a missing failure would otherwise go unnoticed. The task +/// stream is kept open so the eviction that follows the failure can be polled. +fn worker_expecting_one_nondeterminism_failure( + t: TestHistoryBuilder, + resp: ResponseType, +) -> (Worker, Arc) { + worker_expecting_one_failure(t, resp, WorkflowTaskFailedCause::NonDeterministicError) +} + +fn worker_expecting_one_failure( + t: TestHistoryBuilder, + resp: ResponseType, + expected: WorkflowTaskFailedCause, +) -> (Worker, Arc) { + let failures = Arc::new(AtomicUsize::new(0)); + let counted = failures.clone(); + let mut mock = MockPollCfg::from_resp_batches("wfid", t, [resp], mock_worker_client()); + mock.num_expected_fails = 1; + mock.expect_fail_wft_matcher = Box::new(move |_, cause, _| { + counted.fetch_add(1, Ordering::Relaxed); + *cause == expected + }); + let mut mock = build_mock_pollers(mock); + mock.make_wft_stream_interminable(); + (mock_worker(mock), failures) +} +// A workflow subscribing itself. The command exists at all because every SDK +// matches issued commands against command-generated events in order, so a +// command producing no event would put that matching out of step. This asserts +// the command goes out and that replaying its event does not trip that check. +#[tokio::test] +async fn subscribe_command_round_trips_through_replay() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_subscribed("s1", 4); + t.add_full_wf_task(); + + let mut mock_client = mock_worker_client(); + // A replayed task sends no commands, so there is nothing to assert on the + // completion here. Rejection is what this test watches for, and core + // signals that by failing the task. + mock_client + .expect_complete_workflow_task() + .returning(|_, _| Ok(RespondWorkflowTaskCompletedResponse::default())); + mock_client + .expect_fail_workflow_task() + .returning(|_, _, f| panic!("core rejected the reissued subscribe: {f:?}")); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + // Full history, so this activation replays the recorded subscription. Lang + // reissues the command, and core has to match it to that event rather than + // calling it nondeterministic. + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s1".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); +} + +/// The command has to reach the server carrying the bodies. +/// +/// Only a task that is not being replayed sends commands, so this drives a +/// single open Workflow Task rather than a recorded history. +#[tokio::test] +async fn publish_command_reaches_the_server_with_its_payloads() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .times(1) + .returning(|resp, _| { + let cmd = resp.commands.first().expect("a command was sent"); + assert_eq!(cmd.command_type(), CommandType::AppendStreamRecords); + match cmd.attributes.as_ref().unwrap() { + command::Attributes::AppendStreamRecordsCommandAttributes(a) => { + assert_eq!(a.stream_id, "s1"); + // The bodies are the half of the batch History never sees, + // so the command is the only thing that can carry them. + let bodies: Vec<_> = a + .records + .iter() + .map(|m| m.body.as_ref().unwrap().data.clone()) + .collect(); + assert_eq!(bodies, vec![b"one".to_vec(), b"two".to_vec()]); + } + other => panic!("wrong attributes: {other:?}"), + } + Ok(RespondWorkflowTaskCompletedResponse::default()) + }); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish_two("s1").into()], + )) + .await + .unwrap(); + core.shutdown().await; +} + +/// A publish reissued on replay has to match the event the original run wrote. +/// +/// This is the property the event exists for. Core pops one queued command per +/// command-generated event, so a publish producing none would leave every later +/// command matched against the wrong event. Core signals the mismatch by +/// failing the workflow task, which the mock turns into a panic. +#[tokio::test] +async fn publish_command_round_trips_through_replay() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_records_appended("s1", 0, 2); + t.add_full_wf_task(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .returning(|_, _| Ok(RespondWorkflowTaskCompletedResponse::default())); + mock_client + .expect_fail_workflow_task() + .returning(|_, _, f| panic!("core rejected the reissued publish: {f:?}")); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + // First activation replays the recorded publish: lang reissues it, and core + // has to match it to that event rather than calling it nondeterministic. + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish_two("s1").into()], + )) + .await + .unwrap(); + + core.shutdown().await; +} + +/// The subscribe command has to reach the server, which only a task that is +/// not being replayed will send. +#[tokio::test] +async fn subscribe_command_reaches_the_server() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .times(1) + .returning(|resp, _| { + let cmd = resp.commands.first().expect("a command was sent"); + assert_eq!(cmd.command_type(), CommandType::SubscribeStream); + match cmd.attributes.as_ref().unwrap() { + command::Attributes::SubscribeStreamCommandAttributes(a) => { + assert_eq!(a.stream_id, "s1"); + // Passed through unresolved: the server turns it into a + // real offset and records that. + assert_eq!(a.start_offset, -1); + } + other => panic!("wrong attributes: {other:?}"), + } + Ok(RespondWorkflowTaskCompletedResponse::default()) + }); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s1".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); + core.shutdown().await; +} + +fn publish_two(stream_id: &str) -> AppendStreamRecords { + AppendStreamRecords { + stream_id: stream_id.to_string(), + records: vec![ + StreamRecord { + body: Some(b"one".to_vec().into()), + ..Default::default() + }, + StreamRecord { + body: Some(b"two".to_vec().into()), + ..Default::default() + }, + ], + } +} +/// A publish reissued on replay is held against the recorded event, not only +/// against its type. Core sends no commands while replaying, so this check is +/// the only place a publish to the wrong stream can be noticed. +#[tokio::test] +async fn a_publish_reissued_to_a_different_stream_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_records_appended("s1", 0, 2); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish_two("s2").into()], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// The batch size is part of the record too: the event names how many records +/// landed, so a replay that publishes fewer has diverged from the original run. +#[tokio::test] +async fn a_publish_reissued_with_a_different_batch_size_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_records_appended("s1", 0, 2); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + AppendStreamRecords { + stream_id: "s1".to_string(), + records: vec![StreamRecord { + body: Some(b"one".to_vec().into()), + ..Default::default() + }], + } + .into(), + ], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// A subscription is checked on the stream alone. The recorded start offset is +/// the server's resolution of what the command asked for, so it is not the +/// command's to reproduce. +#[tokio::test] +async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_subscribed("s1", 4); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s2".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} diff --git a/crates/sdk-core/src/replay/history_builder.rs b/crates/sdk-core/src/replay/history_builder.rs index 418600403..2ad963a06 100644 --- a/crates/sdk-core/src/replay/history_builder.rs +++ b/crates/sdk-core/src/replay/history_builder.rs @@ -129,6 +129,32 @@ impl TestHistoryBuilder { self.previous_task_completed_id = id; } + /// Add the event a subscribe-stream command produces. + pub fn add_stream_subscribed(&mut self, stream_id: &str, start_offset: i64) -> i64 { + let attrs = WorkflowStreamSubscribedEventAttributes { + workflow_task_completed_event_id: self.previous_task_completed_id, + stream_id: stream_id.to_string(), + start_offset, + }; + self.add(attrs) + } + + /// Add the event an append-stream-records command produces. + pub fn add_stream_records_appended( + &mut self, + stream_id: &str, + first_offset: i64, + record_count: i64, + ) -> i64 { + let attrs = WorkflowStreamRecordsAppendedEventAttributes { + workflow_task_completed_event_id: self.previous_task_completed_id, + stream_id: stream_id.to_string(), + first_offset, + record_count, + }; + self.add(attrs) + } + /// Add a workflow task timed out event. pub fn add_workflow_task_timed_out(&mut self) { let attrs = WorkflowTaskTimedOutEventAttributes { From 36f218badaddc45af5a8606b71f40748e0f1d61d Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 13:41:49 -0700 Subject: [PATCH 05/10] Held the stream commands to more of what their events record. 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. --- crates/sdk-core/src/core_tests/streams.rs | 137 +++++++++++++++++- .../append_stream_records_state_machine.rs | 33 ++++- .../subscribe_stream_state_machine.rs | 38 +++-- .../workflow/machines/workflow_machines.rs | 31 +++- 4 files changed, 215 insertions(+), 24 deletions(-) diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs index 0deefe2d4..fba4cbec5 100644 --- a/crates/sdk-core/src/core_tests/streams.rs +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -50,6 +50,21 @@ fn worker_expecting_one_failure( mock.make_wft_stream_interminable(); (mock_worker(mock), failures) } + +/// A worker whose only assertion is that nothing is rejected. Core signals a +/// reissued command it will not accept by failing the workflow task, so turning +/// that into a panic is how a test says the command was accepted. +fn worker_rejecting_any_failure(t: TestHistoryBuilder) -> Worker { + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .returning(|_, _| Ok(RespondWorkflowTaskCompletedResponse::default())); + mock_client + .expect_fail_workflow_task() + .returning(|_, _, f| panic!("core rejected a reissued command: {f:?}")); + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + mock_worker(build_mock_pollers(mock)) +} // A workflow subscribing itself. The command exists at all because every SDK // matches issued commands against command-generated events in order, so a // command producing no event would put that matching out of step. This asserts @@ -298,9 +313,35 @@ async fn a_publish_reissued_with_a_different_batch_size_fails_the_task() { core.shutdown().await; } -/// A subscription is checked on the stream alone. The recorded start offset is -/// the server's resolution of what the command asked for, so it is not the -/// command's to reproduce. +/// A reissued append that names no stream is held to the name the run's earlier +/// unnamed appends resolved to. Without that, an empty id would match any +/// recorded stream and a workflow that moved its output would go unnoticed. +#[tokio::test] +async fn a_default_publish_reissued_against_another_stream_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + // The server resolved the run's unnamed appends to this name. + t.add_stream_records_appended("output", 0, 2); + t.add_stream_records_appended("elsewhere", 0, 2); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish_two("").into(), publish_two("").into()], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// The stream a subscription names is part of what the recorded event holds it +/// to. #[tokio::test] async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { let mut t = TestHistoryBuilder::default(); @@ -328,3 +369,93 @@ async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { assert_eq!(failures.load(Ordering::Relaxed), 1); core.shutdown().await; } + +/// An explicit start offset is a value the workflow chose, so the recorded +/// event holds the reissued command to it. Without the check a replay that +/// asked to read from the top would be served from wherever the original run +/// began, and nothing would say so. +#[tokio::test] +async fn a_subscribe_reissued_with_a_different_offset_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_subscribed("s1", 100); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s1".to_string(), + start_offset: 0, + } + .into(), + ], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// A negative offset asks the server where the stream stands, so the recorded +/// answer is its own and the reissued command is not held to it. +#[tokio::test] +async fn a_subscribe_from_the_tail_is_not_held_to_the_recorded_offset() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_subscribed("s1", 100); + t.add_full_wf_task(); + + let core = worker_rejecting_any_failure(t); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s1".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); + core.shutdown().await; +} + +/// A second subscribe to the same stream registers nothing: the server records +/// where the cursor has already reached, which is not what the command asked +/// for. Holding the repeat to its offset would fail a run that did nothing +/// wrong. +#[tokio::test] +async fn a_repeat_subscribe_is_not_held_to_the_offset_it_asked_for() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_subscribed("s1", 100); + // The cursor had moved on by the time the second command was handled. + t.add_stream_subscribed("s1", 140); + t.add_full_wf_task(); + + let core = worker_rejecting_any_failure(t); + + let subscribe = SubscribeStream { + stream_id: "s1".to_string(), + start_offset: 100, + }; + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![subscribe.clone().into(), subscribe.into()], + )) + .await + .unwrap(); + core.shutdown().await; +} diff --git a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs index 5acdc7eda..2d1e3d7d7 100644 --- a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs @@ -6,6 +6,7 @@ use crate::worker::workflow::{ machines::{EventInfo, HistEventData, WFMachinesAdapter}, nondeterminism, }; +use std::{cell::RefCell, rc::Rc}; use temporalio_common::protos::{ coresdk::workflow_commands::AppendStreamRecords, temporal::api::{ @@ -26,23 +27,36 @@ fsm! { } /// What the command claimed, kept so the recorded event can be held against it. +/// +/// The default stream's resolved name is shared with the run's other appends, +/// because only a recorded event carries it and one append's event is what tells +/// the next what the name is. #[derive(Default, Clone)] pub(super) struct SharedState { stream_id: String, record_count: i64, + default_stream_id: DefaultStreamIdRef, } +/// The name the server resolved this run's unnamed appends to, once one of them +/// has been recorded. +pub(super) type DefaultStreamIdRef = Rc>>; + /// Append a batch of records to a stream this workflow owns. /// /// The bodies go to the stream's own log, and History gets one event naming the /// offset range the batch landed at. The offsets are assigned by the server, so /// nothing here predicts them. -pub(super) fn append_stream_records(lang_cmd: AppendStreamRecords) -> NewMachineWithCommand { +pub(super) fn append_stream_records( + lang_cmd: AppendStreamRecords, + default_stream_id: DefaultStreamIdRef, +) -> NewMachineWithCommand { let sm = AppendStreamRecordsMachine::from_parts( Created {}.into(), SharedState { stream_id: lang_cmd.stream_id.clone(), record_count: lang_cmd.records.len() as i64, + default_stream_id, }, ); NewMachineWithCommand { @@ -67,9 +81,18 @@ impl CommandIssued { attrs: WorkflowStreamRecordsAppendedEventAttributes, ) -> AppendStreamRecordsMachineTransition { // An empty id names the workflow's default stream, and the server is the - // one that resolves that name, so only a named stream can be compared. - let same_stream = dat.stream_id.is_empty() || dat.stream_id == attrs.stream_id; - if same_stream && dat.record_count == attrs.record_count { + // one that resolves that name. The resolved name is on the event, so the + // run's first unnamed append is what teaches it and every later one is + // held to it. + let expected = if dat.stream_id.is_empty() { + dat.default_stream_id + .borrow_mut() + .get_or_insert_with(|| attrs.stream_id.clone()) + .clone() + } else { + dat.stream_id.clone() + }; + if expected == attrs.stream_id && dat.record_count == attrs.record_count { TransitionResult::default() } else { TransitionResult::Err(nondeterminism!( @@ -78,7 +101,7 @@ impl CommandIssued { attrs.record_count, attrs.stream_id, dat.record_count, - dat.stream_id + expected )) } } diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs index 3fc08e768..36e55c5fe 100644 --- a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -25,23 +25,33 @@ fsm! { shared on_command_recorded) --> Done; } -/// The stream the command named, kept so the recorded event can be held against it. -/// The start offset is not kept: the server resolves it, so the recorded value is -/// its answer rather than what the command said. +/// What the command claimed, kept so the recorded event can be held against it. +/// +/// The offset is only kept when it is one the server echoes back rather than one +/// it works out for itself. A negative offset asks the server where the stream +/// stands, and a repeat subscription is recorded at wherever the cursor has +/// already reached, so in both cases the recorded value is the server's answer +/// rather than what the command said. #[derive(Default, Clone)] pub(super) struct SharedState { stream_id: String, + start_offset: Option, } /// Subscribe this workflow to a stream. The command carries only the stream id /// and a start offset; the server resolves the addressing, because a workflow /// cannot look it up without doing I/O and a value it carried would be a /// reading rather than a fact. -pub(super) fn subscribe_stream(lang_cmd: SubscribeStream) -> NewMachineWithCommand { +pub(super) fn subscribe_stream( + lang_cmd: SubscribeStream, + first_for_stream: bool, +) -> NewMachineWithCommand { let sm = SubscribeStreamMachine::from_parts( Created {}.into(), SharedState { stream_id: lang_cmd.stream_id.clone(), + start_offset: (first_for_stream && lang_cmd.start_offset >= 0) + .then_some(lang_cmd.start_offset), }, ); NewMachineWithCommand { @@ -65,16 +75,26 @@ impl CommandIssued { dat: &mut SharedState, attrs: WorkflowStreamSubscribedEventAttributes, ) -> SubscribeStreamMachineTransition { - if dat.stream_id == attrs.stream_id { - TransitionResult::default() - } else { - TransitionResult::Err(nondeterminism!( + if dat.stream_id != attrs.stream_id { + return TransitionResult::Err(nondeterminism!( "Recorded subscription to stream {:?} does not match the reissued subscription \ to stream {:?}", attrs.stream_id, dat.stream_id - )) + )); } + if let Some(asked_for) = dat.start_offset + && asked_for != attrs.start_offset + { + return TransitionResult::Err(nondeterminism!( + "Recorded subscription to stream {:?} starts at offset {}, and the reissued \ + subscription asked for offset {}", + attrs.stream_id, + attrs.start_offset, + asked_for + )); + } + TransitionResult::default() } } diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index 7648aed81..d3c8e7a37 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -2,14 +2,17 @@ mod local_acts; use super::{ Machines, NewMachineWithCommand, TemporalStateMachine, - append_stream_records_state_machine::append_stream_records, + append_stream_records_state_machine::{DefaultStreamIdRef, append_stream_records}, cancel_external_state_machine::new_external_cancel, cancel_workflow_state_machine::cancel_workflow, complete_workflow_state_machine::complete_workflow, continue_as_new_workflow_state_machine::continue_as_new, - fail_workflow_state_machine::fail_workflow, local_activity_state_machine::new_local_activity, - patch_state_machine::has_change, signal_external_state_machine::new_external_signal, - subscribe_stream_state_machine::subscribe_stream, timer_state_machine::new_timer, + fail_workflow_state_machine::fail_workflow, + local_activity_state_machine::new_local_activity, + patch_state_machine::has_change, + signal_external_state_machine::new_external_signal, + subscribe_stream_state_machine::subscribe_stream, + timer_state_machine::new_timer, upsert_search_attributes_state_machine::upsert_search_attrs, workflow_machines::local_acts::LocalActivityData, workflow_task_state_machine::WorkflowTaskMachine, @@ -47,7 +50,7 @@ use siphasher::sip::SipHasher13; use slotmap::{SlotMap, SparseSecondaryMap}; use std::{ cell::RefCell, - collections::{HashMap, VecDeque}, + collections::{HashMap, HashSet, VecDeque}, convert::TryInto, hash::{Hash, Hasher}, iter::Peekable, @@ -165,6 +168,16 @@ pub(crate) struct WorkflowMachines { /// Contains extra local-activity related data local_activity_data: LocalActivityData, + /// Streams this run has already issued a subscribe command for. A repeat + /// subscription registers nothing and is recorded at wherever the cursor has + /// reached, so only the first one can be held to the offset it asked for. A + /// cursor put on this run out of band through the stream service leaves no + /// event, so there is no seeing that one from here. + subscribed_stream_ids: HashSet, + /// What the server resolved this run's unnamed appends to, learned from the + /// first one it recorded and shared with the machines that follow. + default_stream_id: DefaultStreamIdRef, + /// The workflow that is being driven by this instance of the machines drive_me: DrivenWorkflow, @@ -304,6 +317,8 @@ impl WorkflowMachines { message_outbox: Default::default(), encountered_patch_markers: Default::default(), local_activity_data: LocalActivityData::default(), + subscribed_stream_ids: Default::default(), + default_stream_id: Default::default(), have_seen_terminal_event: false, worker_config: basics.worker_config, } @@ -1533,7 +1548,7 @@ impl WorkflowMachines { // server assigned and hands nothing back. A workflow that // wants to know where its batch landed reads the stream. self.add_cmd_to_wf_task( - append_stream_records(attrs), + append_stream_records(attrs, self.default_stream_id.clone()), annotations, CommandIdKind::NeverResolves, ); @@ -1542,8 +1557,10 @@ impl WorkflowMachines { // Never resolves: the event it produces records the // subscription and hands nothing back to the workflow. The // ranges arrive later as their own activation jobs. + let first_for_stream = + self.subscribed_stream_ids.insert(attrs.stream_id.clone()); self.add_cmd_to_wf_task( - subscribe_stream(attrs), + subscribe_stream(attrs, first_for_stream), annotations, CommandIdKind::NeverResolves, ); From 1e3532a561c5308d9c8b007d3c06bd82d4a13147 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 13:41:56 -0700 Subject: [PATCH 06/10] Named the unreachable stream machine responses fatal. 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. --- .../workflow/machines/append_stream_records_state_machine.rs | 4 ++-- .../workflow/machines/subscribe_stream_state_machine.rs | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs index 2d1e3d7d7..3e83ea387 100644 --- a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs @@ -116,8 +116,8 @@ impl WFMachinesAdapter for AppendStreamRecordsMachine { _my_command: Self::Command, _event_info: Option, ) -> Result, Self::Error> { - Err(Self::Error::Nondeterminism( - "AppendStreamRecords does not use state machine commands".to_string(), + Err(fatal!( + "AppendStreamRecords does not use state machine commands" )) } } diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs index 36e55c5fe..4f986db48 100644 --- a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -107,8 +107,8 @@ impl WFMachinesAdapter for SubscribeStreamMachine { _my_command: Self::Command, _event_info: Option, ) -> Result, Self::Error> { - Err(Self::Error::Nondeterminism( - "SubscribeStream does not use state machine commands".to_string(), + Err(fatal!( + "SubscribeStream does not use state machine commands" )) } } From 99b760913ed406f8ff799b03f5de49536a426bbe Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 14:36:51 -0700 Subject: [PATCH 07/10] Dropped the subscribe start offset check as unsound. 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. --- crates/sdk-core/src/core_tests/streams.rs | 52 ++++--------------- .../subscribe_stream_state_machine.rs | 41 +++++---------- .../workflow/machines/workflow_machines.rs | 13 +---- 3 files changed, 26 insertions(+), 80 deletions(-) diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs index fba4cbec5..fbe71e482 100644 --- a/crates/sdk-core/src/core_tests/streams.rs +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -370,42 +370,13 @@ async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { core.shutdown().await; } -/// An explicit start offset is a value the workflow chose, so the recorded -/// event holds the reissued command to it. Without the check a replay that -/// asked to read from the top would be served from wherever the original run -/// began, and nothing would say so. +/// The start offset is deliberately not compared, even when the command names +/// one itself. The comparison would only be sound for the 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 from history. Failing a run that +/// did nothing wrong costs more than the drift the check would catch. #[tokio::test] -async fn a_subscribe_reissued_with_a_different_offset_fails_the_task() { - let mut t = TestHistoryBuilder::default(); - t.add_by_type(EventType::WorkflowExecutionStarted); - t.add_full_wf_task(); - t.add_stream_subscribed("s1", 100); - t.add_full_wf_task(); - - let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); - - let task = core.poll_workflow_activation().await.unwrap(); - core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( - task.run_id, - vec![ - SubscribeStream { - stream_id: "s1".to_string(), - start_offset: 0, - } - .into(), - ], - )) - .await - .unwrap(); - core.handle_eviction().await; - assert_eq!(failures.load(Ordering::Relaxed), 1); - core.shutdown().await; -} - -/// A negative offset asks the server where the stream stands, so the recorded -/// answer is its own and the reissued command is not held to it. -#[tokio::test] -async fn a_subscribe_from_the_tail_is_not_held_to_the_recorded_offset() { +async fn a_subscribe_reissued_with_a_different_offset_is_accepted() { let mut t = TestHistoryBuilder::default(); t.add_by_type(EventType::WorkflowExecutionStarted); t.add_full_wf_task(); @@ -420,7 +391,7 @@ async fn a_subscribe_from_the_tail_is_not_held_to_the_recorded_offset() { vec![ SubscribeStream { stream_id: "s1".to_string(), - start_offset: -1, + start_offset: 0, } .into(), ], @@ -430,12 +401,11 @@ async fn a_subscribe_from_the_tail_is_not_held_to_the_recorded_offset() { core.shutdown().await; } -/// A second subscribe to the same stream registers nothing: the server records -/// where the cursor has already reached, which is not what the command asked -/// for. Holding the repeat to its offset would fail a run that did nothing -/// wrong. +/// A second subscribe to the same stream registers nothing on the server, which +/// records the event at wherever the cursor has already reached. The reissued +/// command still has to be accepted against it. #[tokio::test] -async fn a_repeat_subscribe_is_not_held_to_the_offset_it_asked_for() { +async fn a_repeat_subscribe_is_accepted() { let mut t = TestHistoryBuilder::default(); t.add_by_type(EventType::WorkflowExecutionStarted); t.add_full_wf_task(); diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs index 4f986db48..f6b4b3099 100644 --- a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -25,33 +25,28 @@ fsm! { shared on_command_recorded) --> Done; } -/// What the command claimed, kept so the recorded event can be held against it. +/// The stream the command named, kept so the recorded event can be held against it. /// -/// The offset is only kept when it is one the server echoes back rather than one -/// it works out for itself. A negative offset asks the server where the stream -/// stands, and a repeat subscription is recorded at wherever the cursor has -/// already reached, so in both cases the recorded value is the server's answer -/// rather than what the command said. +/// The offset is not kept. Comparing it would only be sound for the run's first +/// subscribe to a stream, since the server records a later one at wherever the +/// cursor has already reached, and which one is first cannot be told from here: +/// a subscription made through the stream service leaves no event at all. A +/// check that can fire on a run that did nothing wrong costs more than the drift +/// it would catch. #[derive(Default, Clone)] pub(super) struct SharedState { stream_id: String, - start_offset: Option, } /// Subscribe this workflow to a stream. The command carries only the stream id /// and a start offset; the server resolves the addressing, because a workflow /// cannot look it up without doing I/O and a value it carried would be a /// reading rather than a fact. -pub(super) fn subscribe_stream( - lang_cmd: SubscribeStream, - first_for_stream: bool, -) -> NewMachineWithCommand { +pub(super) fn subscribe_stream(lang_cmd: SubscribeStream) -> NewMachineWithCommand { let sm = SubscribeStreamMachine::from_parts( Created {}.into(), SharedState { stream_id: lang_cmd.stream_id.clone(), - start_offset: (first_for_stream && lang_cmd.start_offset >= 0) - .then_some(lang_cmd.start_offset), }, ); NewMachineWithCommand { @@ -75,26 +70,16 @@ impl CommandIssued { dat: &mut SharedState, attrs: WorkflowStreamSubscribedEventAttributes, ) -> SubscribeStreamMachineTransition { - if dat.stream_id != attrs.stream_id { - return TransitionResult::Err(nondeterminism!( + if dat.stream_id == attrs.stream_id { + TransitionResult::default() + } else { + TransitionResult::Err(nondeterminism!( "Recorded subscription to stream {:?} does not match the reissued subscription \ to stream {:?}", attrs.stream_id, dat.stream_id - )); + )) } - if let Some(asked_for) = dat.start_offset - && asked_for != attrs.start_offset - { - return TransitionResult::Err(nondeterminism!( - "Recorded subscription to stream {:?} starts at offset {}, and the reissued \ - subscription asked for offset {}", - attrs.stream_id, - attrs.start_offset, - asked_for - )); - } - TransitionResult::default() } } diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index d3c8e7a37..a43b75b0f 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -50,7 +50,7 @@ use siphasher::sip::SipHasher13; use slotmap::{SlotMap, SparseSecondaryMap}; use std::{ cell::RefCell, - collections::{HashMap, HashSet, VecDeque}, + collections::{HashMap, VecDeque}, convert::TryInto, hash::{Hash, Hasher}, iter::Peekable, @@ -168,12 +168,6 @@ pub(crate) struct WorkflowMachines { /// Contains extra local-activity related data local_activity_data: LocalActivityData, - /// Streams this run has already issued a subscribe command for. A repeat - /// subscription registers nothing and is recorded at wherever the cursor has - /// reached, so only the first one can be held to the offset it asked for. A - /// cursor put on this run out of band through the stream service leaves no - /// event, so there is no seeing that one from here. - subscribed_stream_ids: HashSet, /// What the server resolved this run's unnamed appends to, learned from the /// first one it recorded and shared with the machines that follow. default_stream_id: DefaultStreamIdRef, @@ -317,7 +311,6 @@ impl WorkflowMachines { message_outbox: Default::default(), encountered_patch_markers: Default::default(), local_activity_data: LocalActivityData::default(), - subscribed_stream_ids: Default::default(), default_stream_id: Default::default(), have_seen_terminal_event: false, worker_config: basics.worker_config, @@ -1557,10 +1550,8 @@ impl WorkflowMachines { // Never resolves: the event it produces records the // subscription and hands nothing back to the workflow. The // ranges arrive later as their own activation jobs. - let first_for_stream = - self.subscribed_stream_ids.insert(attrs.stream_id.clone()); self.add_cmd_to_wf_task( - subscribe_stream(attrs, first_for_stream), + subscribe_stream(attrs), annotations, CommandIdKind::NeverResolves, ); From 5b4088ece815686ef65e471110d7f0954ae953b4 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Fri, 25 Sep 2026 15:38:50 -0700 Subject: [PATCH 08/10] Renamed the lang stream command fields to match the wire. 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. --- .../workflow_commands/workflow_commands.proto | 21 +++++++++--- crates/protos/src/protos/mod.rs | 8 ++--- crates/sdk-core/src/core_tests/streams.rs | 25 +++++++------- crates/sdk-core/src/replay/history_builder.rs | 11 ++++--- .../append_stream_records_state_machine.rs | 33 ++++++++++--------- .../subscribe_stream_state_machine.rs | 10 +++--- .../workflow/machines/workflow_machines.rs | 8 ++--- 7 files changed, 66 insertions(+), 50 deletions(-) diff --git a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto index 456055d6c..6c3e7832b 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto @@ -71,8 +71,11 @@ message WorkflowCommand { // stores each record with an empty producer id, because the workflow is the // producer here. message AppendStreamRecords { - // Empty means the workflow's default output stream. - string stream_id = 1; + // Name of a stream this workflow owns, scoped to the workflow. Created on + // first use. Empty means the workflow's default output stream. A workflow + // cannot append to a stream in another execution, so this is never the id + // of a standalone stream. + string stream_name = 1; repeated temporal.api.stream.v1.StreamRecord records = 2; } @@ -83,9 +86,17 @@ message AppendStreamRecords { // up without doing I/O, and a value it carried would be a reading rather than a // fact, so it could differ on replay. message SubscribeStream { - string stream_id = 1; - // Negative means from wherever the stream is when the subscription is - // registered. The server resolves that once and records it. + // Stream to consume, named either way round: a stream this workflow owns + // by the name it appends under, a stream in another execution by its id. + // The server tries them in that order, so a workflow that owns a stream + // under this name cannot reach a standalone stream with the same id. When + // neither exists the workflow gets a stream of its own by that name, which + // is how a reader subscribes before the first record is written. + string stream_name_or_id = 1; + // Where to start, as an absolute offset. Any negative value means the head + // of the stream as of registration, and they all mean the same thing. The + // server resolves it and records the result, so replay does not resolve it + // again. int64 start_offset = 2; } diff --git a/crates/protos/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index 3cf466a9a..34193d501 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1495,7 +1495,7 @@ pub mod coresdk { impl Display for SubscribeStream { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { - write!(f, "SubscribeStream({})", self.stream_id) + write!(f, "SubscribeStream({})", self.stream_name_or_id) } } @@ -1504,7 +1504,7 @@ pub mod coresdk { write!( f, "AppendStreamRecords({}, {} records)", - self.stream_id, + self.stream_name, self.records.len() ) } @@ -1912,7 +1912,7 @@ pub mod temporal { fn from(s: workflow_commands::AppendStreamRecords) -> Self { Self::AppendStreamRecordsCommandAttributes( AppendStreamRecordsCommandAttributes { - stream_id: s.stream_id, + stream_name: s.stream_name, records: s.records, }, ) @@ -1923,7 +1923,7 @@ pub mod temporal { fn from(s: workflow_commands::SubscribeStream) -> Self { Self::SubscribeStreamCommandAttributes( SubscribeStreamCommandAttributes { - stream_id: s.stream_id, + stream_name_or_id: s.stream_name_or_id, start_offset: s.start_offset, }, ) diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs index fbe71e482..debd8b950 100644 --- a/crates/sdk-core/src/core_tests/streams.rs +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -99,7 +99,7 @@ async fn subscribe_command_round_trips_through_replay() { task.run_id, vec![ SubscribeStream { - stream_id: "s1".to_string(), + stream_name_or_id: "s1".to_string(), start_offset: -1, } .into(), @@ -128,7 +128,7 @@ async fn publish_command_reaches_the_server_with_its_payloads() { assert_eq!(cmd.command_type(), CommandType::AppendStreamRecords); match cmd.attributes.as_ref().unwrap() { command::Attributes::AppendStreamRecordsCommandAttributes(a) => { - assert_eq!(a.stream_id, "s1"); + assert_eq!(a.stream_name, "s1"); // The bodies are the half of the batch History never sees, // so the command is the only thing that can carry them. let bodies: Vec<_> = a @@ -211,7 +211,7 @@ async fn subscribe_command_reaches_the_server() { assert_eq!(cmd.command_type(), CommandType::SubscribeStream); match cmd.attributes.as_ref().unwrap() { command::Attributes::SubscribeStreamCommandAttributes(a) => { - assert_eq!(a.stream_id, "s1"); + assert_eq!(a.stream_name_or_id, "s1"); // Passed through unresolved: the server turns it into a // real offset and records that. assert_eq!(a.start_offset, -1); @@ -229,7 +229,7 @@ async fn subscribe_command_reaches_the_server() { task.run_id, vec![ SubscribeStream { - stream_id: "s1".to_string(), + stream_name_or_id: "s1".to_string(), start_offset: -1, } .into(), @@ -240,9 +240,9 @@ async fn subscribe_command_reaches_the_server() { core.shutdown().await; } -fn publish_two(stream_id: &str) -> AppendStreamRecords { +fn publish_two(stream_name: &str) -> AppendStreamRecords { AppendStreamRecords { - stream_id: stream_id.to_string(), + stream_name: stream_name.to_string(), records: vec![ StreamRecord { body: Some(b"one".to_vec().into()), @@ -280,8 +280,9 @@ async fn a_publish_reissued_to_a_different_stream_fails_the_task() { core.shutdown().await; } -/// The batch size is part of the record too: the event names how many records -/// landed, so a replay that publishes fewer has diverged from the original run. +/// The batch size is part of the record too: the event names the offset range +/// the batch landed at, so a replay that publishes fewer has diverged from the +/// original run. #[tokio::test] async fn a_publish_reissued_with_a_different_batch_size_fails_the_task() { let mut t = TestHistoryBuilder::default(); @@ -297,7 +298,7 @@ async fn a_publish_reissued_with_a_different_batch_size_fails_the_task() { task.run_id, vec![ AppendStreamRecords { - stream_id: "s1".to_string(), + stream_name: "s1".to_string(), records: vec![StreamRecord { body: Some(b"one".to_vec().into()), ..Default::default() @@ -357,7 +358,7 @@ async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { task.run_id, vec![ SubscribeStream { - stream_id: "s2".to_string(), + stream_name_or_id: "s2".to_string(), start_offset: -1, } .into(), @@ -390,7 +391,7 @@ async fn a_subscribe_reissued_with_a_different_offset_is_accepted() { task.run_id, vec![ SubscribeStream { - stream_id: "s1".to_string(), + stream_name_or_id: "s1".to_string(), start_offset: 0, } .into(), @@ -417,7 +418,7 @@ async fn a_repeat_subscribe_is_accepted() { let core = worker_rejecting_any_failure(t); let subscribe = SubscribeStream { - stream_id: "s1".to_string(), + stream_name_or_id: "s1".to_string(), start_offset: 100, }; let task = core.poll_workflow_activation().await.unwrap(); diff --git a/crates/sdk-core/src/replay/history_builder.rs b/crates/sdk-core/src/replay/history_builder.rs index 2ad963a06..2d88d1b7f 100644 --- a/crates/sdk-core/src/replay/history_builder.rs +++ b/crates/sdk-core/src/replay/history_builder.rs @@ -139,18 +139,19 @@ impl TestHistoryBuilder { self.add(attrs) } - /// Add the event an append-stream-records command produces. + /// Add the event an append-stream-records command produces. The range is + /// half-open, as it is on the event. pub fn add_stream_records_appended( &mut self, stream_id: &str, - first_offset: i64, - record_count: i64, + from_offset: i64, + to_offset: i64, ) -> i64 { let attrs = WorkflowStreamRecordsAppendedEventAttributes { workflow_task_completed_event_id: self.previous_task_completed_id, stream_id: stream_id.to_string(), - first_offset, - record_count, + from_offset, + to_offset, }; self.add(attrs) } diff --git a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs index 3e83ea387..e7c2ad302 100644 --- a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs @@ -33,14 +33,14 @@ fsm! { /// the next what the name is. #[derive(Default, Clone)] pub(super) struct SharedState { - stream_id: String, + stream_name: String, record_count: i64, - default_stream_id: DefaultStreamIdRef, + default_stream_name: DefaultStreamNameRef, } /// The name the server resolved this run's unnamed appends to, once one of them /// has been recorded. -pub(super) type DefaultStreamIdRef = Rc>>; +pub(super) type DefaultStreamNameRef = Rc>>; /// Append a batch of records to a stream this workflow owns. /// @@ -49,14 +49,14 @@ pub(super) type DefaultStreamIdRef = Rc>>; /// nothing here predicts them. pub(super) fn append_stream_records( lang_cmd: AppendStreamRecords, - default_stream_id: DefaultStreamIdRef, + default_stream_name: DefaultStreamNameRef, ) -> NewMachineWithCommand { let sm = AppendStreamRecordsMachine::from_parts( Created {}.into(), SharedState { - stream_id: lang_cmd.stream_id.clone(), + stream_name: lang_cmd.stream_name.clone(), record_count: lang_cmd.records.len() as i64, - default_stream_id, + default_stream_name, }, ); NewMachineWithCommand { @@ -80,25 +80,28 @@ impl CommandIssued { dat: &mut SharedState, attrs: WorkflowStreamRecordsAppendedEventAttributes, ) -> AppendStreamRecordsMachineTransition { - // An empty id names the workflow's default stream, and the server is the - // one that resolves that name. The resolved name is on the event, so the - // run's first unnamed append is what teaches it and every later one is - // held to it. - let expected = if dat.stream_id.is_empty() { - dat.default_stream_id + // An empty name means the workflow's default stream, and the server is + // the one that resolves that name. The resolved name is on the event, so + // the run's first unnamed append is what teaches it and every later one + // is held to it. + let expected = if dat.stream_name.is_empty() { + dat.default_stream_name .borrow_mut() .get_or_insert_with(|| attrs.stream_id.clone()) .clone() } else { - dat.stream_id.clone() + dat.stream_name.clone() }; - if expected == attrs.stream_id && dat.record_count == attrs.record_count { + // The event names a half-open offset range rather than a count, and the + // batch size is what this machine can hold a reissued command to. + let recorded_count = attrs.to_offset - attrs.from_offset; + if expected == attrs.stream_id && dat.record_count == recorded_count { TransitionResult::default() } else { TransitionResult::Err(nondeterminism!( "Recorded append of {} records to stream {:?} does not match the reissued \ append of {} records to stream {:?}", - attrs.record_count, + recorded_count, attrs.stream_id, dat.record_count, expected diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs index f6b4b3099..2e547d0d3 100644 --- a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -35,10 +35,10 @@ fsm! { /// it would catch. #[derive(Default, Clone)] pub(super) struct SharedState { - stream_id: String, + stream_name_or_id: String, } -/// Subscribe this workflow to a stream. The command carries only the stream id +/// Subscribe this workflow to a stream. The command carries only the name or id /// and a start offset; the server resolves the addressing, because a workflow /// cannot look it up without doing I/O and a value it carried would be a /// reading rather than a fact. @@ -46,7 +46,7 @@ pub(super) fn subscribe_stream(lang_cmd: SubscribeStream) -> NewMachineWithComma let sm = SubscribeStreamMachine::from_parts( Created {}.into(), SharedState { - stream_id: lang_cmd.stream_id.clone(), + stream_name_or_id: lang_cmd.stream_name_or_id.clone(), }, ); NewMachineWithCommand { @@ -70,14 +70,14 @@ impl CommandIssued { dat: &mut SharedState, attrs: WorkflowStreamSubscribedEventAttributes, ) -> SubscribeStreamMachineTransition { - if dat.stream_id == attrs.stream_id { + if dat.stream_name_or_id == attrs.stream_id { TransitionResult::default() } else { TransitionResult::Err(nondeterminism!( "Recorded subscription to stream {:?} does not match the reissued subscription \ to stream {:?}", attrs.stream_id, - dat.stream_id + dat.stream_name_or_id )) } } diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index a43b75b0f..d96ae0eaf 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -2,7 +2,7 @@ mod local_acts; use super::{ Machines, NewMachineWithCommand, TemporalStateMachine, - append_stream_records_state_machine::{DefaultStreamIdRef, append_stream_records}, + append_stream_records_state_machine::{DefaultStreamNameRef, append_stream_records}, cancel_external_state_machine::new_external_cancel, cancel_workflow_state_machine::cancel_workflow, complete_workflow_state_machine::complete_workflow, @@ -170,7 +170,7 @@ pub(crate) struct WorkflowMachines { /// What the server resolved this run's unnamed appends to, learned from the /// first one it recorded and shared with the machines that follow. - default_stream_id: DefaultStreamIdRef, + default_stream_name: DefaultStreamNameRef, /// The workflow that is being driven by this instance of the machines drive_me: DrivenWorkflow, @@ -311,7 +311,7 @@ impl WorkflowMachines { message_outbox: Default::default(), encountered_patch_markers: Default::default(), local_activity_data: LocalActivityData::default(), - default_stream_id: Default::default(), + default_stream_name: Default::default(), have_seen_terminal_event: false, worker_config: basics.worker_config, } @@ -1541,7 +1541,7 @@ impl WorkflowMachines { // server assigned and hands nothing back. A workflow that // wants to know where its batch landed reads the stream. self.add_cmd_to_wf_task( - append_stream_records(attrs, self.default_stream_id.clone()), + append_stream_records(attrs, self.default_stream_name.clone()), annotations, CommandIdKind::NeverResolves, ); From 07fcbfadd80cbdc6c7502ec9b72f6655907766d1 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Mon, 28 Sep 2026 16:34:22 -0700 Subject: [PATCH 09/10] Passed the subscribe start position from lang to the server. 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. --- .../workflow_commands/workflow_commands.proto | 10 ++-- crates/protos/src/protos/mod.rs | 1 + crates/sdk-core/src/core_tests/streams.rs | 49 ++++++++++++++----- 3 files changed, 43 insertions(+), 17 deletions(-) diff --git a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto index b58816481..81c45b64c 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto @@ -93,11 +93,13 @@ message SubscribeStream { // neither exists the workflow gets a stream of its own by that name, which // is how a reader subscribes before the first record is written. string stream_name_or_id = 1; - // Where to start, as an absolute offset. Any negative value means the head - // of the stream as of registration, and they all mean the same thing. The - // server resolves it and records the result, so replay does not resolve it - // again. + // Where to start, as an absolute offset. Read only when start_position is + // unset. The server refuses a negative value. int64 start_offset = 2; + // Where to start: an offset, the earliest record held, the tail, or the + // last N records. The server resolves it and records the resulting offset, + // so replay does not resolve it again and Core does not keep it. + temporal.api.stream.v1.StreamStartPosition start_position = 3; } message StartTimer { diff --git a/crates/protos/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index f835b32e8..6f28b8296 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1925,6 +1925,7 @@ pub mod temporal { SubscribeStreamCommandAttributes { stream_name_or_id: s.stream_name_or_id, start_offset: s.start_offset, + start_position: s.start_position, }, ) } diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs index debd8b950..fa6514a75 100644 --- a/crates/sdk-core/src/core_tests/streams.rs +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -16,7 +16,7 @@ use temporalio_common::protos::{ temporal::api::{ command::v1::command, enums::v1::{CommandType, EventType, WorkflowTaskFailedCause}, - stream::v1::StreamRecord, + stream::v1::{StreamRecord, StreamStartPosition, stream_start_position::Position}, workflowservice::v1::RespondWorkflowTaskCompletedResponse, }, }; @@ -100,7 +100,8 @@ async fn subscribe_command_round_trips_through_replay() { vec![ SubscribeStream { stream_name_or_id: "s1".to_string(), - start_offset: -1, + start_offset: 0, + start_position: start(Position::Tail(true)), } .into(), ], @@ -195,26 +196,38 @@ async fn publish_command_round_trips_through_replay() { } /// The subscribe command has to reach the server, which only a task that is -/// not being replayed will send. +/// not being replayed will send. The start position goes out as lang gave it: +/// the server resolves it and records the offset, so Core never holds one. #[tokio::test] -async fn subscribe_command_reaches_the_server() { +async fn subscribe_command_reaches_the_server_with_each_start_position() { + for position in [ + Position::Offset(7), + Position::LastN(3), + Position::Earliest(true), + Position::Tail(true), + ] { + subscribe_reaches_the_server(start(position)).await; + } +} + +async fn subscribe_reaches_the_server(start_position: Option) { let mut t = TestHistoryBuilder::default(); t.add_by_type(EventType::WorkflowExecutionStarted); t.add_workflow_task_scheduled_and_started(); + let expected = start_position; let mut mock_client = mock_worker_client(); mock_client .expect_complete_workflow_task() .times(1) - .returning(|resp, _| { + .returning(move |resp, _| { let cmd = resp.commands.first().expect("a command was sent"); assert_eq!(cmd.command_type(), CommandType::SubscribeStream); match cmd.attributes.as_ref().unwrap() { command::Attributes::SubscribeStreamCommandAttributes(a) => { assert_eq!(a.stream_name_or_id, "s1"); - // Passed through unresolved: the server turns it into a - // real offset and records that. - assert_eq!(a.start_offset, -1); + assert_eq!(a.start_offset, 0); + assert_eq!(a.start_position, expected); } other => panic!("wrong attributes: {other:?}"), } @@ -230,7 +243,8 @@ async fn subscribe_command_reaches_the_server() { vec![ SubscribeStream { stream_name_or_id: "s1".to_string(), - start_offset: -1, + start_offset: 0, + start_position, } .into(), ], @@ -240,6 +254,12 @@ async fn subscribe_command_reaches_the_server() { core.shutdown().await; } +fn start(position: Position) -> Option { + Some(StreamStartPosition { + position: Some(position), + }) +} + fn publish_two(stream_name: &str) -> AppendStreamRecords { AppendStreamRecords { stream_name: stream_name.to_string(), @@ -359,7 +379,8 @@ async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { vec![ SubscribeStream { stream_name_or_id: "s2".to_string(), - start_offset: -1, + start_offset: 0, + start_position: start(Position::Tail(true)), } .into(), ], @@ -371,8 +392,8 @@ async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { core.shutdown().await; } -/// The start offset is deliberately not compared, even when the command names -/// one itself. The comparison would only be sound for the run's first subscribe +/// The start is deliberately not compared, even when the command names an +/// absolute offset. The comparison would only be sound for the 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 from history. Failing a run that /// did nothing wrong costs more than the drift the check would catch. @@ -393,6 +414,7 @@ async fn a_subscribe_reissued_with_a_different_offset_is_accepted() { SubscribeStream { stream_name_or_id: "s1".to_string(), start_offset: 0, + start_position: start(Position::Earliest(true)), } .into(), ], @@ -419,7 +441,8 @@ async fn a_repeat_subscribe_is_accepted() { let subscribe = SubscribeStream { stream_name_or_id: "s1".to_string(), - start_offset: 100, + start_offset: 0, + start_position: start(Position::Offset(100)), }; let task = core.poll_workflow_activation().await.unwrap(); core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( From 245d72945c88248e764e9fa94058b198cc382572 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Thu, 1 Oct 2026 17:44:54 -0700 Subject: [PATCH 10/10] Subscribed workflows to notification channels by command. --- .../workflow_activation.proto | 12 ++ .../workflow_commands/workflow_commands.proto | 11 ++ crates/protos/src/protos/mod.rs | 19 +++ crates/sdk-core/src/core_tests/channels.rs | 134 ++++++++++++++++ crates/sdk-core/src/core_tests/mod.rs | 1 + crates/sdk-core/src/replay/history_builder.rs | 9 ++ .../src/worker/workflow/machines/mod.rs | 3 + ...ribe_notification_channel_state_machine.rs | 143 ++++++++++++++++++ .../workflow/machines/transition_coverage.rs | 5 + .../workflow/machines/workflow_machines.rs | 10 ++ crates/sdk-core/src/worker/workflow/mod.rs | 4 + crates/sdk/src/workflow_future.rs | 6 + crates/workflow/src/runtime/instance.rs | 11 ++ 13 files changed, 368 insertions(+) create mode 100644 crates/sdk-core/src/core_tests/channels.rs create mode 100644 crates/sdk-core/src/worker/workflow/machines/subscribe_notification_channel_state_machine.rs diff --git a/crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto b/crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto index d764a937f..08d13701a 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto @@ -14,6 +14,7 @@ import "temporal/api/failure/v1/message.proto"; import "temporal/api/update/v1/message.proto"; import "temporal/api/common/v1/message.proto"; import "temporal/api/stream/v1/message.proto"; +import "temporal/api/notification/v1/message.proto"; import "temporal/api/enums/v1/failed_cause.proto"; import "temporal/api/enums/v1/workflow.proto"; import "temporal/sdk/core/activity_result/activity_result.proto"; @@ -147,6 +148,8 @@ message WorkflowActivationJob { // // A range of a stream the workflow subscribed to. DeliverStreamRecords deliver_stream_records = 21; + // Notifications from the channels the workflow subscribed to. + NotificationsReceived notifications_received = 22; // Remove the workflow identified by the [WorkflowActivation] containing this job from the // cache after performing the activation. It is guaranteed that this will be the only job // in the activation if present. @@ -173,6 +176,15 @@ message DeliverStreamRecords { repeated temporal.api.stream.v1.StreamRecord records = 4; } +// Hand a workflow the notifications the server folded for its channels. +// +// They come from the scheduled event of the Workflow Task this activation +// belongs to. History is the record, so a replay yields the same job with the +// same notifications at the same point. +message NotificationsReceived { + repeated temporal.api.notification.v1.Notification notifications = 1; +} + // Initialize a new workflow message InitializeWorkflow { // The identifier the lang-specific sdk uses to execute workflow code diff --git a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto index 81c45b64c..2040fff11 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto @@ -60,6 +60,7 @@ message WorkflowCommand { // with that family and must not be reused. SubscribeStream subscribe_stream = 29; AppendStreamRecords append_stream_records = 30; + SubscribeNotificationChannel subscribe_notification_channel = 31; } } @@ -102,6 +103,16 @@ message SubscribeStream { temporal.api.stream.v1.StreamStartPosition start_position = 3; } +// Subscribe this workflow to a notification channel, so the scheduled event of +// each later Workflow Task carries the notifications folded for it. +// +// The notifications live in History rather than arriving by a side channel, so +// a replay reads the same ones the live run saw. +message SubscribeNotificationChannel { + // Name of the channel, scoped to the namespace. + string channel = 1; +} + message StartTimer { // Lang's incremental sequence number, used as the operation identifier uint32 seq = 1; diff --git a/crates/protos/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index c9f934ddd..a041d274c 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1284,6 +1284,9 @@ pub mod coresdk { d.stream_id, d.from_offset, d.to_offset ) } + workflow_activation_job::Variant::NotificationsReceived(n) => { + write!(f, "NotificationsReceived({})", n.notifications.len()) + } } } } @@ -1499,6 +1502,12 @@ pub mod coresdk { } } + impl Display for SubscribeNotificationChannel { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "SubscribeNotificationChannel({})", self.channel) + } + } + impl Display for AppendStreamRecords { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!( @@ -1934,6 +1943,16 @@ pub mod temporal { } } + impl From for Attributes { + fn from(s: workflow_commands::SubscribeNotificationChannel) -> Self { + Self::SubscribeNotificationChannelCommandAttributes( + SubscribeNotificationChannelCommandAttributes { + channel: s.channel, + }, + ) + } + } + impl From for command::Attributes { fn from(s: workflow_commands::StartTimer) -> Self { Self::StartTimerCommandAttributes(StartTimerCommandAttributes { diff --git a/crates/sdk-core/src/core_tests/channels.rs b/crates/sdk-core/src/core_tests/channels.rs new file mode 100644 index 000000000..d0c6de5fd --- /dev/null +++ b/crates/sdk-core/src/core_tests/channels.rs @@ -0,0 +1,134 @@ +use crate::{ + replay::TestHistoryBuilder, + test_help::{MockPollCfg, ResponseType, WorkerTestHelpers, build_mock_pollers, mock_worker}, + worker::client::mocks::mock_worker_client, +}; +use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, +}; +use temporalio_common::protos::{ + coresdk::{ + workflow_commands::SubscribeNotificationChannel, + workflow_completion::WorkflowActivationCompletion, + }, + temporal::api::{ + command::v1::command, + enums::v1::{CommandType, EventType, WorkflowTaskFailedCause}, + workflowservice::v1::RespondWorkflowTaskCompletedResponse, + }, +}; + +fn subscribe(channel: &str) -> SubscribeNotificationChannel { + SubscribeNotificationChannel { + channel: channel.to_string(), + } +} + +/// The command has to reach the server naming the channel. Only a task that is +/// not being replayed sends commands, so this drives a single open task. +#[tokio::test] +async fn subscribe_channel_command_reaches_the_server() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .times(1) + .returning(|resp, _| { + let cmd = resp.commands.first().expect("a command was sent"); + assert_eq!( + cmd.command_type(), + CommandType::SubscribeNotificationChannel + ); + match cmd.attributes.as_ref().unwrap() { + command::Attributes::SubscribeNotificationChannelCommandAttributes(a) => { + assert_eq!(a.channel, "orders"); + } + other => panic!("wrong attributes: {other:?}"), + } + Ok(RespondWorkflowTaskCompletedResponse::default()) + }); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![subscribe("orders").into()], + )) + .await + .unwrap(); + core.shutdown().await; +} + +/// A subscribe reissued on replay has to match the event the original run +/// wrote. Core pops one queued command per command-generated event, so this is +/// what keeps every later command lined up with its own event. +#[tokio::test] +async fn subscribe_channel_command_round_trips_through_replay() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_notification_channel_subscribed("orders"); + t.add_full_wf_task(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .returning(|_, _| Ok(RespondWorkflowTaskCompletedResponse::default())); + mock_client + .expect_fail_workflow_task() + .returning(|_, _, f| panic!("core rejected the reissued subscribe: {f:?}")); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![subscribe("orders").into()], + )) + .await + .unwrap(); + core.shutdown().await; +} + +/// The channel is part of what the recorded event holds the reissued command +/// to. Core sends no commands while replaying, so this check is the only place +/// a subscribe to the wrong channel can be noticed. +#[tokio::test] +async fn a_subscribe_reissued_to_a_different_channel_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_notification_channel_subscribed("orders"); + t.add_full_wf_task(); + + let failures = Arc::new(AtomicUsize::new(0)); + let counted = failures.clone(); + let mut mock = + MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_worker_client()); + mock.num_expected_fails = 1; + mock.expect_fail_wft_matcher = Box::new(move |_, cause, _| { + counted.fetch_add(1, Ordering::Relaxed); + *cause == WorkflowTaskFailedCause::NonDeterministicError + }); + let mut mock = build_mock_pollers(mock); + mock.make_wft_stream_interminable(); + let core = mock_worker(mock); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![subscribe("invoices").into()], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} diff --git a/crates/sdk-core/src/core_tests/mod.rs b/crates/sdk-core/src/core_tests/mod.rs index 4779016bb..4fe4ab5e8 100644 --- a/crates/sdk-core/src/core_tests/mod.rs +++ b/crates/sdk-core/src/core_tests/mod.rs @@ -1,4 +1,5 @@ mod activity_tasks; +mod channels; mod event_groups; mod queries; mod replay_flag; diff --git a/crates/sdk-core/src/replay/history_builder.rs b/crates/sdk-core/src/replay/history_builder.rs index 2d88d1b7f..60ba8f9a1 100644 --- a/crates/sdk-core/src/replay/history_builder.rs +++ b/crates/sdk-core/src/replay/history_builder.rs @@ -139,6 +139,15 @@ impl TestHistoryBuilder { self.add(attrs) } + /// Add the event a subscribe-notification-channel command produces. + pub fn add_notification_channel_subscribed(&mut self, channel: &str) -> i64 { + let attrs = WorkflowNotificationChannelSubscribedEventAttributes { + workflow_task_completed_event_id: self.previous_task_completed_id, + channel: channel.to_string(), + }; + self.add(attrs) + } + /// Add the event an append-stream-records command produces. The range is /// half-open, as it is on the event. pub fn add_stream_records_appended( diff --git a/crates/sdk-core/src/worker/workflow/machines/mod.rs b/crates/sdk-core/src/worker/workflow/machines/mod.rs index 364897465..4f1044ae9 100644 --- a/crates/sdk-core/src/worker/workflow/machines/mod.rs +++ b/crates/sdk-core/src/worker/workflow/machines/mod.rs @@ -16,6 +16,7 @@ mod modify_workflow_properties_state_machine; mod nexus_operation_state_machine; mod patch_state_machine; mod signal_external_state_machine; +mod subscribe_notification_channel_state_machine; mod subscribe_stream_state_machine; mod timer_state_machine; mod update_state_machine; @@ -49,6 +50,7 @@ use std::{ convert::{TryFrom, TryInto}, fmt::{Debug, Display}, }; +use subscribe_notification_channel_state_machine::SubscribeNotificationChannelMachine; use subscribe_stream_state_machine::SubscribeStreamMachine; use temporalio_common::{ fsm_trait::{StateMachine, TransitionResult}, @@ -86,6 +88,7 @@ enum Machines { ModifyWorkflowPropertiesMachine, SubscribeStreamMachine, AppendStreamRecordsMachine, + SubscribeNotificationChannelMachine, UpdateMachine, NexusOperationMachine, } diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_notification_channel_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_notification_channel_state_machine.rs new file mode 100644 index 000000000..a3b2b8794 --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_notification_channel_state_machine.rs @@ -0,0 +1,143 @@ +use super::{ + NewMachineWithCommand, StateMachine, TransitionResult, fsm, workflow_machines::MachineResponse, +}; +use crate::worker::workflow::{ + WFMachinesError, fatal, + machines::{EventInfo, HistEventData, WFMachinesAdapter}, + nondeterminism, +}; +use temporalio_common::protos::{ + coresdk::workflow_commands::SubscribeNotificationChannel, + temporal::api::{ + enums::v1::{CommandType, EventType}, + history::v1::{WorkflowNotificationChannelSubscribedEventAttributes, history_event}, + }, +}; + +fsm! { + pub(super) name SubscribeNotificationChannelMachine; + command SubscribeNotificationChannelMachineCommand; + error WFMachinesError; + shared_state SharedState; + + Created --(CommandScheduled) --> CommandIssued; + CommandIssued --(CommandRecorded(WorkflowNotificationChannelSubscribedEventAttributes), + shared on_command_recorded) --> Done; +} + +/// The channel the command named, kept so the recorded event can be held against it. +#[derive(Default, Clone)] +pub(super) struct SharedState { + channel: String, +} + +/// Subscribe this workflow to a notification channel. The command carries only +/// the channel name; the notifications arrive later on the scheduled event of +/// each Workflow Task, so nothing is handed back here. +pub(super) fn subscribe_notification_channel( + lang_cmd: SubscribeNotificationChannel, +) -> NewMachineWithCommand { + let sm = SubscribeNotificationChannelMachine::from_parts( + Created {}.into(), + SharedState { + channel: lang_cmd.channel.clone(), + }, + ); + NewMachineWithCommand { + command: lang_cmd.into(), + machine: sm.into(), + } +} + +#[derive(Debug, derive_more::Display)] +pub(super) enum SubscribeNotificationChannelMachineCommand {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Created {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct CommandIssued {} + +impl CommandIssued { + pub(super) fn on_command_recorded( + self, + dat: &mut SharedState, + attrs: WorkflowNotificationChannelSubscribedEventAttributes, + ) -> SubscribeNotificationChannelMachineTransition { + if dat.channel == attrs.channel { + TransitionResult::default() + } else { + TransitionResult::Err(nondeterminism!( + "Recorded subscription to notification channel {:?} does not match the \ + reissued subscription to channel {:?}", + attrs.channel, + dat.channel + )) + } + } +} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Done {} + +impl WFMachinesAdapter for SubscribeNotificationChannelMachine { + fn adapt_response( + &self, + _my_command: Self::Command, + _event_info: Option, + ) -> Result, Self::Error> { + Err(fatal!( + "SubscribeNotificationChannel does not use state machine commands" + )) + } +} + +impl TryFrom for SubscribeNotificationChannelMachineEvents { + type Error = WFMachinesError; + + fn try_from(e: HistEventData) -> Result { + let e = e.event; + match e.event_type() { + EventType::WorkflowNotificationChannelSubscribed => { + if let Some( + history_event::Attributes::WorkflowNotificationChannelSubscribedEventAttributes( + attrs, + ), + ) = e.attributes + { + Ok(SubscribeNotificationChannelMachineEvents::CommandRecorded( + attrs, + )) + } else { + Err(fatal!( + "Notification channel subscribed attributes were unset: {e}" + )) + } + } + _ => Err(Self::Error::Nondeterminism(format!( + "SubscribeNotificationChannelMachine does not handle {e}" + ))), + } + } +} + +impl TryFrom for SubscribeNotificationChannelMachineEvents { + type Error = WFMachinesError; + + fn try_from(c: CommandType) -> Result { + match c { + CommandType::SubscribeNotificationChannel => { + Ok(SubscribeNotificationChannelMachineEvents::CommandScheduled) + } + _ => Err(Self::Error::Nondeterminism(format!( + "SubscribeNotificationChannelMachine does not handle command type {c:?}" + ))), + } + } +} + +impl From for CommandIssued { + fn from(_: Created) -> Self { + Self {} + } +} diff --git a/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs index 41c054450..46cb072e1 100644 --- a/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs +++ b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs @@ -77,6 +77,7 @@ mod machine_coverage_report { modify_workflow_properties_state_machine::ModifyWorkflowPropertiesMachine, nexus_operation_state_machine::NexusOperationMachine, patch_state_machine::PatchMachine, signal_external_state_machine::SignalExternalMachine, + subscribe_notification_channel_state_machine::SubscribeNotificationChannelMachine, subscribe_stream_state_machine::SubscribeStreamMachine, timer_state_machine::TimerMachine, update_state_machine::UpdateMachine, upsert_search_attributes_state_machine::UpsertSearchAttributesMachine, @@ -121,6 +122,7 @@ mod machine_coverage_report { let mut nexus = NexusOperationMachine::visualizer().to_owned(); let mut subscribe_stream = SubscribeStreamMachine::visualizer().to_owned(); let mut append_stream_records = AppendStreamRecordsMachine::visualizer().to_owned(); + let mut subscribe_channel = SubscribeNotificationChannelMachine::visualizer().to_owned(); // This isn't at all efficient but doesn't need to be. // Replace transitions in the vizzes with green color if they are covered. @@ -154,6 +156,9 @@ mod machine_coverage_report { m @ "AppendStreamRecordsMachine" => { cover_transitions(m, &mut append_stream_records, coverage) } + m @ "SubscribeNotificationChannelMachine" => { + cover_transitions(m, &mut subscribe_channel, coverage) + } m => panic!("Unknown machine {m}"), } } diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index d96ae0eaf..b1f1751fc 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -11,6 +11,7 @@ use super::{ local_activity_state_machine::new_local_activity, patch_state_machine::has_change, signal_external_state_machine::new_external_signal, + subscribe_notification_channel_state_machine::subscribe_notification_channel, subscribe_stream_state_machine::subscribe_stream, timer_state_machine::new_timer, upsert_search_attributes_state_machine::upsert_search_attrs, @@ -1556,6 +1557,15 @@ impl WorkflowMachines { CommandIdKind::NeverResolves, ); } + WFCommandVariant::SubscribeNotificationChannel(attrs) => { + // Never resolves: the notifications arrive on the scheduled + // event of later tasks, not as a reply to this command. + self.add_cmd_to_wf_task( + subscribe_notification_channel(attrs), + annotations, + CommandIdKind::NeverResolves, + ); + } WFCommandVariant::UpdateResponse(ur) => { let m_key = self.get_machine_by_msg(&ur.protocol_instance_id)?; let m = if let Machines::UpdateMachine(m) = self.machine_mut(m_key) { diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index 68db21e9f..bcc45038d 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -1575,6 +1575,7 @@ enum WFCommandVariant { RequestCancelNexusOperation(RequestCancelNexusOperation), SubscribeStream(SubscribeStream), AppendStreamRecords(AppendStreamRecords), + SubscribeNotificationChannel(SubscribeNotificationChannel), } impl TryFrom for WFCommand { @@ -1587,6 +1588,9 @@ impl TryFrom for WFCommand { workflow_command::Variant::AppendStreamRecords(s) => { WFCommandVariant::AppendStreamRecords(s) } + workflow_command::Variant::SubscribeNotificationChannel(s) => { + WFCommandVariant::SubscribeNotificationChannel(s) + } workflow_command::Variant::CancelTimer(s) => WFCommandVariant::CancelTimer(s), workflow_command::Variant::ScheduleActivity(s) => WFCommandVariant::AddActivity(s), workflow_command::Variant::RequestCancelActivity(s) => { diff --git a/crates/sdk/src/workflow_future.rs b/crates/sdk/src/workflow_future.rs index 1c745a38a..681fe6098 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -345,6 +345,12 @@ impl WorkflowFuture { slice.stream_id ); } + Variant::NotificationsReceived(_) => { + // No channel API in this SDK, so nothing here could have + // subscribed. Bailing rather than ignoring keeps a run that + // depends on them from going on as if none arrived. + bail!("received channel notifications, which this SDK cannot deliver"); + } Variant::RemoveFromCache(_) => { unreachable!("Cache removal should happen higher up"); } diff --git a/crates/workflow/src/runtime/instance.rs b/crates/workflow/src/runtime/instance.rs index 07ce1da68..2ce9025cc 100644 --- a/crates/workflow/src/runtime/instance.rs +++ b/crates/workflow/src/runtime/instance.rs @@ -1132,6 +1132,17 @@ where ..Default::default() })); } + Some(ActivationVariant::NotificationsReceived(_)) => { + // This runtime has no command to subscribe to a channel, so + // a run that gets notifications was driven by something it + // cannot model. Failing beats handing back a workflow that + // silently behaves as if none arrived. + return Err(Box::new(Failure { + message: "received channel notifications, which this SDK cannot deliver" + .to_string(), + ..Default::default() + })); + } Some(ActivationVariant::RemoveFromCache(_)) => ActivationJobResult::None, None => { return Err(Box::new(Failure {