From aec5383fd25127c166e322f368c416037721d0d4 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Sat, 3 Oct 2026 02:00:42 -0700 Subject: [PATCH 1/3] Delivered consumed stream ranges to workflows, live and on replay. The poll response carries the records for the task about to run and re-supplies what earlier tasks consumed, keyed by the completion that recorded each range. Core hands them over in a fixed order, gives a consuming task its own activation on replay, and fails a task whose recorded range arrived without its records. --- .../workflow_activation.proto | 2 +- crates/sdk-core/src/protosext/mod.rs | 7 + crates/sdk-core/src/replay/mod.rs | 19 ++ .../src/worker/workflow/history_update.rs | 62 ++++- .../workflow/machines/workflow_machines.rs | 248 +++++++++++++++++- .../src/worker/workflow/managed_run.rs | 13 +- crates/sdk-core/src/worker/workflow/mod.rs | 13 + 7 files changed, 351 insertions(+), 13 deletions(-) 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 3d3998f6d..3e57a6547 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 @@ -42,7 +42,7 @@ import "temporal/sdk/core/nexus/nexus.proto"; // * Signal and update handlers should be invoked before workflow routines are iterated. That is to // say before the users' main workflow function and anything spawned by it is allowed to continue. // * Channel notifications are input from outside the workflow, like signals, so they go with -// them. +// them and ahead of the stream ranges among the other jobs. // * Local activities resolutions go after other normal jobs because while *not* replaying, they // will always take longer than anything else that produces an immediate job (which is // effectively instant). When *replaying* we need to scan ahead for LA markers so that we can diff --git a/crates/sdk-core/src/protosext/mod.rs b/crates/sdk-core/src/protosext/mod.rs index 506b953c6..456cad5ce 100644 --- a/crates/sdk-core/src/protosext/mod.rs +++ b/crates/sdk-core/src/protosext/mod.rs @@ -39,6 +39,7 @@ use temporalio_common::protos::{ history::v1::{History, HistoryEvent, MarkerRecordedEventAttributes, history_event}, query::v1::WorkflowQuery, sdk::v1::{EventGroupMarker, UserMetadata}, + stream::v1::StreamSlice, workflowservice::v1::PollWorkflowTaskQueueResponse, }, utilities::TryIntoOrNone, @@ -64,6 +65,10 @@ pub(crate) struct ValidPollWFTQResponse { pub(crate) query_requests: Vec, /// Protocol messages pub(crate) messages: Vec, + /// Ranges of streams this workflow subscribed to. A slice tagged with a + /// completed-event id is re-supplying what an earlier task consumed; an + /// untagged one belongs to the task about to run. + pub(crate) stream_slices: Vec, /// Zero-size field to prevent explicit construction _cant_construct_me: (), @@ -109,6 +114,7 @@ impl TryFrom for ValidPollWFTQResponse { query, queries, messages, + stream_slices, .. } => { if task_token.is_empty() { @@ -133,6 +139,7 @@ impl TryFrom for ValidPollWFTQResponse { legacy_query: query, query_requests, messages, + stream_slices, _cant_construct_me: (), }) } diff --git a/crates/sdk-core/src/replay/mod.rs b/crates/sdk-core/src/replay/mod.rs index 6e4cadc50..692bc18b3 100644 --- a/crates/sdk-core/src/replay/mod.rs +++ b/crates/sdk-core/src/replay/mod.rs @@ -33,6 +33,7 @@ use temporalio_common::{ temporal::api::{ common::v1::WorkflowExecution, history::v1::History, + stream::v1::StreamSlice, workflowservice::v1::{ DescribeNamespaceResponse, RespondWorkflowTaskCompletedResponse, RespondWorkflowTaskFailedResponse, @@ -118,6 +119,7 @@ where workflow_id: history.workflow_id, run_id: hist_info.orig_run_id().to_string(), }); + resp.stream_slices = history.stream_slices; Ok(resp) } else { if let Some(wc) = hlock.worker_closer.get() { @@ -174,6 +176,7 @@ mod tests { pub struct HistoryForReplay { hist: History, workflow_id: String, + stream_slices: Vec, } impl HistoryForReplay { /// Create a new history from replay from something that looks like a history and a workflow id. @@ -181,8 +184,24 @@ impl HistoryForReplay { Self { hist: history.into(), workflow_id: workflow_id.into(), + stream_slices: Vec::new(), } } + + /// Attach the stream records the history's completed tasks consumed. + /// + /// History records the offsets a task consumed and never the payloads, so a workflow that + /// read a stream cannot be replayed from its history alone. Each slice carries the records + /// for one recorded range, tagged with the `WorkflowTaskCompleted` event that recorded it, + /// the same shape the server puts on a poll response when it re-supplies them. Replay hands + /// each range to the activation of the task that consumed it. A slice that disagrees with + /// the recorded range, and a recorded range with content that has no slice, both fail the + /// task as the worker's failure rather than the workflow's: the history and the slices are + /// both given to replay, so neither says the workflow diverged. + pub fn with_stream_slices(mut self, slices: impl IntoIterator) -> Self { + self.stream_slices = slices.into_iter().collect(); + self + } } #[cfg(any(feature = "test-utilities", test))] impl From for HistoryForReplay { diff --git a/crates/sdk-core/src/worker/workflow/history_update.rs b/crates/sdk-core/src/worker/workflow/history_update.rs index 0914c8457..972c5a8f7 100644 --- a/crates/sdk-core/src/worker/workflow/history_update.rs +++ b/crates/sdk-core/src/worker/workflow/history_update.rs @@ -162,6 +162,7 @@ impl HistoryPaginator { query_requests: wft.query_requests, update, messages: wft.messages, + stream_slices: wft.stream_slices, }; Ok((paginator, prepared)) } @@ -282,7 +283,7 @@ impl HistoryPaginator { // We only *really* have the last WFT if the events go all the way up to at least the // WFT started event id. Otherwise we somehow still have partial history. let no_more = matches!(self.next_page_token, NextPageToken::Done) && seen_enough_events; - let (update, extra) = HistoryUpdate::from_events( + let (mut update, extra) = HistoryUpdate::from_events( current_events, self.previous_wft_started_id, self.wft_started_event_id, @@ -309,8 +310,33 @@ impl HistoryPaginator { // There was not a meaningful WFT in the whole page. We must fetch more. continue; } + // The machines read the completion that closes an update's last task while that + // task is the one being replayed: it records the stream range the task consumed, + // and the workflow has to be handed that range in the same activation. A page that + // ends exactly on a WFT started event leaves that completion on the next page, so + // fetch it before handing the update over. + // + // The fetch costs a page for any run whose page boundary lands here, stream or not. + // There is no telling the two apart from here: a subscription made through the + // stream service records no event, so a run can consume ranges with nothing earlier + // in its history to say it would. + if !no_more && self.event_queue.is_empty() { + self.event_queue.extend(update.events); + continue; + } self.id_of_last_event_in_last_extracted_update = update.events.last().map(|e| e.event_id); + // Otherwise the completion is the first retained event. Carry a copy across the + // split so it can be peeked at. The original stays queued and is what the next + // update is built from, and a lone completion yields nothing to take, so no event + // is applied twice. + if let Some(completion) = self + .event_queue + .front() + .filter(|e| e.event_type() == EventType::WorkflowTaskCompleted) + { + update.events.push(completion.clone()); + } #[cfg(debug_assertions)] update.assert_contiguous(); return Ok(update); @@ -640,17 +666,21 @@ impl HistoryUpdate { true } - /// Returns the next WFT completed event attributes, if any, starting at (inclusive) the - /// `from_id` + /// Returns the next WFT completed event, if any, starting at (inclusive) the + /// `from_id`, as its event id and attributes. + /// + /// The id matters to callers that need to key something on the event rather + /// than only read its contents, such as the stream range a task consumed, + /// which is recorded on the completion that closes that task. pub(crate) fn peek_next_wft_completed( &self, from_id: i64, - ) -> Option<&WorkflowTaskCompletedEventAttributes> { + ) -> Option<(i64, &WorkflowTaskCompletedEventAttributes)> { self.events .iter() .skip_while(|e| e.event_id < from_id) .find_map(|e| match &e.attributes { - Some(Attributes::WorkflowTaskCompletedEventAttributes(a)) => Some(a), + Some(Attributes::WorkflowTaskCompletedEventAttributes(a)) => Some((e.event_id, a)), _ => None, }) } @@ -730,6 +760,10 @@ fn find_end_index_of_next_wft_seq( } if e.event_type() == EventType::WorkflowTaskStarted { + // Scoped to this started event. What its own completion consumed says nothing + // about the events the scan already passed, so it must not join the flags that + // carry across them. + let mut completion_consumed_a_range = false; wft_started_event_id_to_index.push((e.event_id, ix)); if let Some(next_event) = events.get(ix + 1) { let next_event_type = next_event.event_type(); @@ -748,8 +782,19 @@ fn find_end_index_of_next_wft_seq( wft_started_event_id_to_index.pop(); continue; } else if next_event_type == EventType::WorkflowTaskCompleted { + // A task that consumed a stream range issued nothing the machines match, + // but it was an activation of its own: the workflow was handed that range + // and ran on it. Replay has to give it its own activation too, rather than + // fold it into a heartbeat chain and hand several ranges over at once. + if let Some(Attributes::WorkflowTaskCompletedEventAttributes(ref attrs)) = + next_event.attributes + && !attrs.consumed_stream_ranges.is_empty() + { + completion_consumed_a_range = true; + } if let Some(next_next_event) = events.get(ix + 2) { if !saw_command + && !completion_consumed_a_range && next_next_event.event_type() == EventType::WorkflowTaskScheduled && !scheduled_with_notifications(next_next_event) { @@ -792,7 +837,10 @@ fn find_end_index_of_next_wft_seq( } return NextWFTSeqEndIndex::Complete(ix); } - } else if !has_last_wft && !saw_command_or_started { + } else if !has_last_wft + && !saw_command_or_started + && !completion_consumed_a_range + { // Don't have enough events to look ahead of the WorkflowTaskCompleted. Need // to fetch more. continue; @@ -803,7 +851,7 @@ fn find_end_index_of_next_wft_seq( // more. continue; } - if saw_command_or_started { + if saw_command_or_started || completion_consumed_a_range { return NextWFTSeqEndIndex::Complete(ix); } } 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 b127f1aee..ee50f0fc7 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -78,6 +78,7 @@ use temporalio_common::{ notification::v1::Notification, protocol::v1::{Message as ProtocolMessage, message::SequencingId}, sdk::v1::WorkflowTaskCompletedMetadata, + stream::v1::{StreamRange, StreamSlice}, }, }, worker::WorkerDeploymentVersion, @@ -94,6 +95,29 @@ pub(crate) struct WorkflowMachines { /// kept because the lang side polls & completes for every workflow task, but we do not need /// to poll the server that often during replay. last_history_from_server: HistoryUpdate, + /// Stream ranges an earlier task consumed, keyed by the WorkflowTaskCompleted + /// event that recorded each one. Only offsets are in History, so replay is + /// served by the server reading the stream again and sending the bytes back + /// tagged with that event. + stream_slices_by_event: HashMap>, + /// Ranges for the task about to run, which no event records yet. + current_stream_slices: Vec, + /// Highest completion event whose recorded range was delivered by looking + /// ahead, or 0. + /// + /// A task's consumed range is recorded on the completion that closes it, so + /// it is found one batch early. That same completion is then seen again as + /// an ordinary event in the next batch, where its slices are already spent + /// and must not be asked for twice. + stream_slices_lookahead_through: i64, + /// Workflow Task Started id of the last task this instance ran live, or 0. + /// + /// A task that ran here was handed its stream records as it ran, and the + /// completion event recording what it consumed arrives in the next task's + /// history. That event needs no bytes from the server, and `replaying` does + /// not say so: it is true for a cache hit as well, because the history + /// begins after an earlier task. + stream_slices_delivered_through: i64, /// Channel notifications from the scheduled events of the task being applied, folded per /// channel. The server clears what it put on a scheduled event, so a retry's event carries /// only later arrivals, and a failed task's notifications reach lang only here, together @@ -283,7 +307,7 @@ impl WorkflowMachines { basics.sdk_version.to_owned(), ); // Peek ahead to determine used flags in the first WFT. - if let Some(attrs) = basics.history.peek_next_wft_completed(0) { + if let Some((_, attrs)) = basics.history.peek_next_wft_completed(0) { observed_internal_flags.add_from_complete(attrs); }; Self { @@ -293,6 +317,10 @@ impl WorkflowMachines { workflow_type: basics.workflow_type, run_id: basics.run_id, drive_me: driven_wf, + stream_slices_by_event: Default::default(), + current_stream_slices: Default::default(), + stream_slices_delivered_through: 0, + stream_slices_lookahead_through: 0, pending_notifications: vec![], replaying, metrics: basics.metrics, @@ -344,11 +372,31 @@ impl WorkflowMachines { &mut self, update: HistoryUpdate, protocol_messages: Vec, + stream_slices: Vec, ) -> Result<()> { if !self.protocol_msgs.is_empty() { dbg_panic!("There are unprocessed protocol messages while receiving new work"); } self.protocol_msgs = protocol_messages; + // An untagged slice belongs to the task about to run; a tagged one is + // re-supplying a range some earlier task already consumed. + self.current_stream_slices.clear(); + self.stream_slices_by_event.clear(); + for slice in stream_slices { + if slice.workflow_task_completed_event_id == 0 { + self.current_stream_slices.push(slice); + } else { + self.stream_slices_by_event + .entry(slice.workflow_task_completed_event_id) + .or_default() + .push(slice); + } + } + // The order ranges arrive in is something the workflow can branch on, and + // the server's order is whatever its iteration happened to produce. Fixing + // it here is what makes it the same live and on replay. + self.current_stream_slices + .sort_by(|a, b| stream_order(&a.stream_id, a.from_offset, &b.stream_id, b.from_offset)); self.new_history_from_server(update)?; Ok(()) } @@ -636,19 +684,40 @@ impl WorkflowMachines { } }}; } + let mut replayed_slice_events: Vec<(i64, Vec)> = vec![]; + // Kept apart from the in-batch ones. An in-batch event whose bytes are + // missing is a real inconsistency; a looked-ahead one may simply not + // have been sent yet, and demanding it would turn an early delivery + // into a new way to fail. + let mut lookahead_slice_event: Option<(i64, Vec)> = None; let mut peeked_events = events.iter().peekable(); while let Some(event) = peeked_events.next() { if let Some(history_event::Attributes::WorkflowTaskCompletedEventAttributes(ref wtc)) = event.attributes { apply_wft_complete_data!(self, wtc); + // The event records which offsets that task consumed; the server + // sends the bytes back separately, keyed by this event. + if !wtc.consumed_stream_ranges.is_empty() { + replayed_slice_events + .push((event.event_id, wtc.consumed_stream_ranges.clone())); + } } if peeked_events.peek().is_none() - && let Some(wtc) = self + && let Some((wtc_id, wtc)) = self .last_history_from_server .peek_next_wft_completed(event.event_id) { apply_wft_complete_data!(self, wtc); + // The range this batch's task consumed is recorded on the + // completion that closes it, which lands in the *next* batch. + // Collecting it only when it appears would hand the workflow + // its input one activation after the commands that input + // caused, so a read-then-publish task replays with nothing to + // decide from and reissues no command. Look ahead for it here. + if !wtc.consumed_stream_ranges.is_empty() { + lookahead_slice_event = Some((wtc_id, wtc.consumed_stream_ranges.clone())); + } } } @@ -850,6 +919,91 @@ impl WorkflowMachines { } } + // Ranges an earlier task consumed, in the order those tasks ran, so a + // replaying workflow observes them exactly as it did the first time. + // + // Only while replaying. A cache hit is handed the previous task's + // completion in its history as well, and that event carries the range + // that task consumed, but that task ran here and its messages were + // delivered live at the time. There is nothing to re-supply, the server + // sends nothing, and re-supplying would hand the workflow the same + // messages twice. + replayed_slice_events.sort_unstable_by_key(|(event_id, _)| *event_id); + for (event_id, cursors) in replayed_slice_events { + // Delivered live to this instance when the task ran, so there is + // nothing to re-supply and the server sent nothing. + if event_id <= self.stream_slices_delivered_through + 1 { + self.stream_slices_by_event.remove(&event_id); + continue; + } + // Already handed over when this task was looked ahead to. + if event_id <= self.stream_slices_lookahead_through { + self.stream_slices_by_event.remove(&event_id); + continue; + } + let slices = self.stream_slices_by_event.remove(&event_id); + if slices.is_none() + && let Some(cursor) = cursors.iter().find(|c| c.from_offset < c.to_offset) + { + // History says a task consumed a range and the server sent no + // bytes for it. Replaying with less data than the original run + // had produces different commands, and the mismatch would + // surface later as an unrelated nondeterminism error. Reaching + // this from here also means the lookahead did not find the + // completion, so the range was due one activation before this. + return Err(WFMachinesError::MissingRecords(format!( + "Event {event_id} records that stream {} was consumed from offset {} to {}, \ + but the server sent no records for it. The workflow was owed that range \ + one activation earlier.", + cursor.stream_id, cursor.from_offset, cursor.to_offset + ))); + } + for job in resupplied_deliveries(event_id, cursors, slices.unwrap_or_default())? { + self.drive_me.send_job(job); + } + } + // The task about to be replayed, whose range is only visible by looking + // ahead to the completion that closes it. + if let Some((event_id, cursors)) = lookahead_slice_event + && event_id > self.stream_slices_delivered_through + 1 + && event_id > self.stream_slices_lookahead_through + { + let slices = self.stream_slices_by_event.remove(&event_id); + if slices.is_none() + && let Some(cursor) = cursors.iter().find(|c| c.from_offset < c.to_offset) + { + // The bytes only ever travel on the response that carried this + // task, so absent bytes for a range with content mean this + // worker was handed a task it cannot replay: a sticky task for a + // run it no longer holds, whose history it fetched itself. + // Failing here, before the workflow runs on less input than it + // had, keeps a legacy query from being answered from the wrong + // state, and the server's retry on the normal task queue + // carries the records for both task kinds. + return Err(WFMachinesError::MissingRecords(format!( + "Event {event_id} records that stream {} was consumed from offset {} to {}, \ + but the server sent no records for it. A task dispatched with a partial \ + history to a worker that no longer holds the run cannot replay it; the \ + retry on the normal task queue carries the records.", + cursor.stream_id, cursor.from_offset, cursor.to_offset + ))); + } + for job in resupplied_deliveries(event_id, cursors, slices.unwrap_or_default())? { + self.drive_me.send_job(job); + } + self.stream_slices_lookahead_through = event_id; + } + // Then the range for the task about to run, which is only meaningful + // once we have caught up to it. + if !self.replaying { + for slice in std::mem::take(&mut self.current_stream_slices) { + self.drive_me.send_job(deliver_stream_records_job(slice)); + } + // This task is running here, so whatever it consumes is already in + // hand and its completion event will not need re-supplying. + self.stream_slices_delivered_through = self.next_started_event_id; + } + // Only record replay latency if we actually did replay work. This avoids recording // near-zero latencies for the first workflow task (which has no history to replay) or // when there were no events to process. @@ -1904,3 +2058,93 @@ fn fold_notification(folded: &mut Vec, n: Notification) { None => folded.push(n), } } + +/// Turn a slice the server supplied into the job lang sees. +fn deliver_stream_records_job(slice: StreamSlice) -> OutgoingJob { + workflow_activation::DeliverStreamRecords { + stream_id: slice.stream_id, + from_offset: slice.from_offset, + to_offset: slice.to_offset, + records: slice.records, + } + .into() +} + +/// The one order ranges are handed to a workflow in, live and on replay. +fn stream_order( + stream_a: &str, + from_offset_a: i64, + stream_b: &str, + from_offset_b: i64, +) -> std::cmp::Ordering { + (stream_a, from_offset_a).cmp(&(stream_b, from_offset_b)) +} + +/// Pair the ranges a completion event recorded with the slices the server sent +/// back for it, in the order the workflow is handed them. +/// +/// The event is the record of what the task saw, so the slices have to match +/// it rather than the other way round. A range that observed nothing needs no +/// bytes and is rebuilt from the cursor alone; a range with content has to +/// arrive, and what arrives has to cover exactly the recorded offsets. Anything +/// else would replay the task with different input than it ran on. +/// +/// Every disagreement here is between two things the server produced, History +/// on one side and the poll response on the other. The workflow's own commands +/// reach none of it, so none of these is the workflow's fault and none of them +/// is nondeterminism. +fn resupplied_deliveries( + event_id: i64, + mut cursors: Vec, + mut slices: Vec, +) -> Result> { + cursors.sort_by(|a, b| stream_order(&a.stream_id, a.from_offset, &b.stream_id, b.from_offset)); + let mut jobs = Vec::with_capacity(cursors.len()); + for cursor in cursors { + let slice = slices + .iter() + .position(|s| s.stream_id == cursor.stream_id) + .map(|ix| slices.swap_remove(ix)); + let slice = match slice { + Some(slice) + if slice.from_offset == cursor.from_offset + && slice.to_offset == cursor.to_offset => + { + slice + } + Some(slice) => { + return Err(WFMachinesError::MissingRecords(format!( + "Event {event_id} records that stream {} was consumed from offset {} to {}, \ + but the server sent offsets {} to {} for it", + cursor.stream_id, + cursor.from_offset, + cursor.to_offset, + slice.from_offset, + slice.to_offset + ))); + } + None if cursor.from_offset == cursor.to_offset => StreamSlice { + stream_id: cursor.stream_id, + from_offset: cursor.from_offset, + to_offset: cursor.to_offset, + ..Default::default() + }, + None => { + return Err(WFMachinesError::MissingRecords(format!( + "Event {event_id} records that stream {} was consumed from offset {} to {}, \ + but the server sent no messages for it", + cursor.stream_id, cursor.from_offset, cursor.to_offset + ))); + } + }; + jobs.push(deliver_stream_records_job(slice)); + } + if let Some(extra) = slices.first() { + return Err(WFMachinesError::MissingRecords(format!( + "The server sent stream {} from offset {} to {} for event {event_id}, which records \ + no such range", + extra.stream_id, extra.from_offset, extra.to_offset + ))); + } + Ok(jobs) +} diff --git a/crates/sdk-core/src/worker/workflow/managed_run.rs b/crates/sdk-core/src/worker/workflow/managed_run.rs index e9bcb6272..a5c881cd1 100644 --- a/crates/sdk-core/src/worker/workflow/managed_run.rs +++ b/crates/sdk-core/src/worker/workflow/managed_run.rs @@ -42,6 +42,7 @@ use temporalio_common::protos::{ temporal::api::{ enums::v1::{VersioningBehavior, WorkflowTaskFailedCause}, failure::v1::Failure, + stream::v1::StreamSlice, }, }; use tokio::sync::oneshot; @@ -247,7 +248,8 @@ impl ManagedRun { if is_incremental { self.metrics.sticky_cache_hit(); } - self.wfm.new_work_from_server(work.update, work.messages)? + self.wfm + .new_work_from_server(work.update, work.messages, work.stream_slices)? } else { let r = self.wfm.get_next_activation()?; if r.jobs.is_empty() { @@ -626,7 +628,10 @@ impl ManagedRun { EvictionReason::Unspecified | EvictionReason::PaginationOrHistoryFetch ); - let rur = if is_no_report_query_fail { + // An unreported query failure leaves an intact run in the cache for the retry to use. A + // run whose machines broke while it was being brought up to the query is given up + // instead, since it can produce nothing more, so the retry starts from history. + let rur = if is_no_report_query_fail && !self.am_broken { None } else { // Blow up any cached data associated with the workflow @@ -1446,8 +1451,10 @@ impl WorkflowManager { &mut self, update: HistoryUpdate, messages: Vec, + stream_slices: Vec, ) -> Result { - self.machines.new_work_from_server(update, messages)?; + self.machines + .new_work_from_server(update, messages, stream_slices)?; self.get_next_activation() } diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index 35f44f95b..d02714723 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -87,6 +87,7 @@ use temporalio_common::{ protocol::v1::Message as ProtocolMessage, query::v1::WorkflowQuery, sdk::v1::{EventGroupMarker, UserMetadata, WorkflowTaskCompletedMetadata}, + stream::v1::StreamSlice, taskqueue::v1::StickyExecutionAttributes, workflowservice::v1::{PollActivityTaskQueueResponse, get_system_info_response}, }, @@ -1016,6 +1017,7 @@ struct PreparedWFT { query_requests: Vec, update: HistoryUpdate, messages: Vec, + stream_slices: Vec, } impl PreparedWFT { @@ -1723,6 +1725,15 @@ pub(crate) enum WFMachinesError { Nondeterminism(String), #[error("Fatal error in workflow machines: {0}")] Fatal(String), + /// History records that a task consumed stream records and the response that carried the + /// task did not bring them, so this worker cannot replay the run. Covers a response that + /// brought none of them and one whose records do not cover what History says the task read. + /// Not the workflow's fault either way: both sides of that comparison come from the server, + /// and a worker handed a sticky task for a run it no longer holds has no way to fetch the + /// records. Treated like a failed history fetch, so a legacy query goes unanswered and the + /// server retries it where the records travel. + #[error("Workflow task cannot be replayed on this worker: {0}")] + MissingRecords(String), } /// Helper macro to create Nondeterminism errors with automatic assertion @@ -1802,6 +1813,7 @@ impl WFMachinesError { match self { WFMachinesError::Nondeterminism(_) => EvictionReason::Nondeterminism, WFMachinesError::Fatal(_) => EvictionReason::Fatal, + WFMachinesError::MissingRecords(_) => EvictionReason::PaginationOrHistoryFetch, } } @@ -1907,6 +1919,7 @@ fn prepare_to_ship_activation(wfa: &mut WorkflowActivation) { workflow_activation_job::Variant::UpdateRandomSeed(_) => 2, workflow_activation_job::Variant::SignalWorkflow(_) => 3, workflow_activation_job::Variant::DoUpdate(_) => 3, + // Ahead of the stream ranges, which fall in the default bucket below. workflow_activation_job::Variant::NotificationsReceived(_) => 3, workflow_activation_job::Variant::ResolveActivity(ra) if ra.is_local => 5, // In principle we should never actually need to sort these with the others, since From ed6840b7c8eba500a522bb532dcce1c6df43e8d6 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Sat, 3 Oct 2026 02:00:42 -0700 Subject: [PATCH 2/3] Tested stream delivery and replay re-supply. The cases cover live delivery, replay across history pages, the lookahead to the closing completion, missing and mismatched records, and notifications ordered ahead of a stream range. --- crates/sdk-core/src/core_tests/channels.rs | 33 +- crates/sdk-core/src/core_tests/streams.rs | 1081 ++++++++++++++++- crates/sdk-core/src/replay/history_builder.rs | 17 + .../sdk-core/src/test_help/integ_helpers.rs | 40 +- 4 files changed, 1162 insertions(+), 9 deletions(-) diff --git a/crates/sdk-core/src/core_tests/channels.rs b/crates/sdk-core/src/core_tests/channels.rs index 147e7e37e..4cb2fee9c 100644 --- a/crates/sdk-core/src/core_tests/channels.rs +++ b/crates/sdk-core/src/core_tests/channels.rs @@ -2,8 +2,8 @@ use crate::{ init_replay_worker, replay::{HistoryFeeder, HistoryForReplay, ReplayWorkerInput, TestHistoryBuilder}, test_help::{ - MockPollCfg, ResponseType, WorkerTestHelpers, build_mock_pollers, mock_worker, - test_worker_cfg, + MockPollCfg, PollWFTRespExt, ResponseType, WorkerTestHelpers, build_mock_pollers, + hist_to_poll_resp, mock_worker, test_worker_cfg, }, worker::client::mocks::mock_worker_client, }; @@ -351,6 +351,7 @@ fn job_kinds(task: &WorkflowActivation) -> Vec<&'static str> { .map(|j| match j.variant.as_ref().unwrap() { workflow_activation_job::Variant::InitializeWorkflow(_) => "init", workflow_activation_job::Variant::NotificationsReceived(_) => "notifications", + workflow_activation_job::Variant::DeliverStreamRecords(_) => "stream", _ => "other", }) .collect() @@ -368,6 +369,34 @@ fn one_task_with_notifications() -> TestHistoryBuilder { t } +/// The notifications on the scheduled event reach the activation of that task, +/// as one job, ahead of the stream ranges the same task was handed. +#[tokio::test] +async fn notifications_on_the_scheduled_event_reach_the_live_activation() { + let t = one_task_with_notifications(); + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("s1", 0, 0, &["alpha"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + assert!(!task.is_replaying); + assert_eq!(job_kinds(&task), vec!["init", "notifications", "stream"]); + assert_eq!( + received(&task), + vec![vec![notification("orders", 3), notification("invoices", 7)]] + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); +} + /// Replaying the same history yields the same job. Nothing is re-supplied: the /// notifications are on the event, so History alone carries them. #[tokio::test] diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs index fa6514a75..8b1423679 100644 --- a/crates/sdk-core/src/core_tests/streams.rs +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -1,7 +1,18 @@ +//! Delivery of server-side stream ranges to a workflow. +//! +//! History records the offsets a task consumed and never the payloads, so the +//! server sends the bytes on the poll response: untagged for the task about to +//! run, and tagged with a WorkflowTaskCompleted event id when it is +//! re-supplying what an earlier task consumed. These tests pin that both +//! arrive, in the right order, and that an empty range is still delivered. + use crate::{ - Worker, - replay::TestHistoryBuilder, - test_help::{MockPollCfg, ResponseType, WorkerTestHelpers, build_mock_pollers, mock_worker}, + Worker, init_replay_worker, + replay::{HistoryFeeder, HistoryForReplay, ReplayWorkerInput, TestHistoryBuilder}, + test_help::{ + MockPollCfg, PollWFTRespExt, ResponseType, WorkerTestHelpers, build_mock_pollers, + hist_to_poll_resp, mock_worker, test_worker_cfg, + }, worker::client::mocks::mock_worker_client, }; use std::sync::{ @@ -10,17 +21,68 @@ use std::sync::{ }; use temporalio_common::protos::{ coresdk::{ - workflow_commands::{AppendStreamRecords, SubscribeStream}, + workflow_activation::{ + RemoveFromCache, WorkflowActivation, WorkflowActivationJob, + remove_from_cache::EvictionReason, workflow_activation_job, + }, + workflow_commands::{AppendStreamRecords, CompleteWorkflowExecution, SubscribeStream}, workflow_completion::WorkflowActivationCompletion, }, temporal::api::{ command::v1::command, enums::v1::{CommandType, EventType, WorkflowTaskFailedCause}, - stream::v1::{StreamRecord, StreamStartPosition, stream_start_position::Position}, - workflowservice::v1::RespondWorkflowTaskCompletedResponse, + history::v1::{History, HistoryEvent}, + query::v1::WorkflowQuery, + stream::v1::{ + StreamRange, StreamRecord, StreamSlice, StreamStartPosition, + stream_start_position::Position, + }, + workflowservice::v1::{ + GetWorkflowExecutionHistoryResponse, RespondWorkflowTaskCompletedResponse, + }, }, }; +fn cursor(stream_id: &str, from: i64, to: i64) -> StreamRange { + StreamRange { + stream_id: stream_id.to_string(), + from_offset: from, + to_offset: to, + } +} + +fn delivered(job: &WorkflowActivationJob) -> (&str, i64, i64, Vec<&[u8]>) { + match job.variant.as_ref().unwrap() { + workflow_activation_job::Variant::DeliverStreamRecords(d) => ( + d.stream_id.as_str(), + d.from_offset, + d.to_offset, + d.records + .iter() + .map(|m| m.body.as_ref().unwrap().data.as_slice()) + .collect(), + ), + other => panic!("expected a stream delivery, got {other:?}"), + } +} + +/// The stream deliveries in an activation, as (stream, from, to). +fn delivered_ranges(task: &WorkflowActivation) -> Vec<(String, i64, i64)> { + task.jobs + .iter() + .filter(|j| { + matches!( + j.variant, + Some(workflow_activation_job::Variant::DeliverStreamRecords(_)) + ) + }) + .map(|j| { + let (stream, from, to, _) = delivered(j); + (stream.to_string(), from, to) + }) + .collect() +} + /// 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 @@ -65,6 +127,165 @@ fn worker_rejecting_any_failure(t: TestHistoryBuilder) -> Worker { let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); mock_worker(build_mock_pollers(mock)) } + +fn history_page( + events: &[HistoryEvent], + next_page_token: Vec, +) -> GetWorkflowExecutionHistoryResponse { + GetWorkflowExecutionHistoryResponse { + history: Some(History { + events: events.to_vec(), + }), + next_page_token, + ..Default::default() + } +} + +#[tokio::test] +async fn delivers_the_range_for_the_current_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("s1", 0, 0, &["alpha", "beta"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + let stream_jobs: Vec<_> = task + .jobs + .iter() + .filter(|j| { + matches!( + j.variant, + Some(workflow_activation_job::Variant::DeliverStreamRecords(_)) + ) + }) + .collect(); + assert_eq!(stream_jobs.len(), 1); + assert_eq!( + delivered(stream_jobs[0]), + ("s1", 0, 2, vec![b"alpha".as_slice(), b"beta".as_slice()]) + ); + + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); +} + +// A task where the subscription saw nothing is a fact replay has to reproduce, +// so the range still has to arrive rather than being dropped as uninteresting. +#[tokio::test] +async fn delivers_an_empty_range() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("s1", 0, 4, &[]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + let job = task + .jobs + .iter() + .find(|j| { + matches!( + j.variant, + Some(workflow_activation_job::Variant::DeliverStreamRecords(_)) + ) + }) + .expect("an empty range is still delivered"); + assert_eq!(delivered(job), ("s1", 4, 4, vec![])); + + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); +} + +// On a cache miss every prior task replays, so ranges those tasks consumed have +// to be handed back in the order they were consumed, before the range for the +// task about to run. +#[tokio::test] +async fn replays_recorded_ranges_in_order_before_the_current_one() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let first_completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 0, 2)]); + t.add_workflow_task_scheduled_and_started(); + let second_completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 2, 3)]); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + // Deliberately out of order, to prove the ordering comes from the events + // rather than from however the server happened to lay them out. + poll_resp.add_stream_slice("s1", second_completed, 2, &["gamma"]); + poll_resp.add_stream_slice("s1", 0, 3, &["delta"]); + poll_resp.add_stream_slice("s1", first_completed, 0, &["alpha", "beta"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let mut mock = build_mock_pollers(mock); + mock.worker_cfg(|wc| wc.max_cached_workflows = 1); + let core = mock_worker(mock); + + // Each entry is (activation, from, to, bodies). The activation matters as + // much as the order: a range that came back one activation late would + // still be in order. + let mut seen = vec![]; + for activation in 1..=3 { + let task = core.poll_workflow_activation().await.unwrap(); + for job in &task.jobs { + if matches!( + job.variant, + Some(workflow_activation_job::Variant::DeliverStreamRecords(_)) + ) { + let (_, from, to, bodies) = delivered(job); + seen.push(( + activation, + from, + to, + bodies.iter().map(|b| b.to_vec()).collect::>(), + )); + } + } + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + } + + assert_eq!( + seen, + vec![ + (1, 0, 2, vec![b"alpha".to_vec(), b"beta".to_vec()]), + (2, 2, 3, vec![b"gamma".to_vec()]), + (3, 3, 4, vec![b"delta".to_vec()]), + ], + "each recorded range comes back in the activation of the task that consumed it, \ + then the live one" + ); +} + // 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 @@ -275,6 +496,296 @@ fn publish_two(stream_name: &str) -> AppendStreamRecords { ], } } + +/// A task that reads a range and publishes because of what it read. +/// +/// This is the shape the product requires: workflow code observes stream input, +/// decides, and writes. Replaying it means the workflow has to be handed the +/// input again *before* core matches the command that input caused, otherwise +/// lang has nothing to decide from and reissues nothing. +/// +/// The recorded range lives on the WorkflowTaskCompleted that closes the task, +/// which is the event *after* the one that started it. So the range for the +/// task about to be replayed is only visible by looking ahead, and a lookahead +/// that reads the completion for its flags but not for its cursors delivers the +/// input one activation too late. +#[tokio::test] +async fn read_then_publish_replays_when_the_range_is_only_visible_by_lookahead() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + // The task that reads [0,1) and publishes because of it. + let read_completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("in", 0, 1)]); + t.add_stream_records_appended("out", 0, 1); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("in", read_completed, 0, &["go"]); + + 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 read-caused publish: {f:?}")); + + let mock = + MockPollCfg::from_resp_batches("wfid", t, [ResponseType::Raw(poll_resp.resp)], mock_client); + let mut mock = build_mock_pollers(mock); + // Cold: nothing cached, so this is reconstruction from History. + mock.worker_cfg(|wc| wc.max_cached_workflows = 0); + let core = mock_worker(mock); + + let task = core.poll_workflow_activation().await.unwrap(); + let got_input = task.jobs.iter().any(|j| { + matches!( + j.variant, + Some(workflow_activation_job::Variant::DeliverStreamRecords(_)) + ) + }); + assert!( + got_input, + "the replayed task must receive the range it consumed before its \ + resulting publish is matched; jobs were {:?}", + task.jobs + ); + + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + AppendStreamRecords { + stream_name: "out".to_string(), + records: vec![StreamRecord { + body: Some(b"accept".to_vec().into()), + ..Default::default() + }], + } + .into(), + ], + )) + .await + .unwrap(); + + core.shutdown().await; +} + +/// The real shape: a subscribe task with no consumed range, then a task that +/// reads and publishes, then a third task replaying both. +/// +/// The first completion carries no cursors, so the lookahead that finds the +/// range has to keep looking past it rather than stopping at the first +/// completion it sees. +#[tokio::test] +async fn read_then_publish_replays_after_a_task_that_consumed_nothing() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + // Task 1 subscribes and consumes nothing. + t.add_workflow_task_completed(); + t.add_stream_subscribed("in", 0); + t.add_workflow_task_scheduled_and_started(); + // Task 2 reads [0,3) and publishes because of it. + let read_completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("in", 0, 3)]); + t.add_stream_records_appended("out", 0, 1); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("in", read_completed, 0, &["a", "b", "c"]); + + 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 replayed read-caused publish: {f:?}")); + + let mock = + MockPollCfg::from_resp_batches("wfid", t, [ResponseType::Raw(poll_resp.resp)], mock_client); + let mut mock = build_mock_pollers(mock); + mock.worker_cfg(|wc| wc.max_cached_workflows = 0); + let core = mock_worker(mock); + + // Task 1: subscribe, no input yet. + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_name_or_id: "in".to_string(), + start_offset: 0, + start_position: start(Position::Earliest(true)), + } + .into(), + ], + )) + .await + .unwrap(); + + // Task 2: the recorded range has to arrive before its publish is matched. + let task = core.poll_workflow_activation().await.unwrap(); + let bodies: Vec> = task + .jobs + .iter() + .filter(|j| { + matches!( + j.variant, + Some(workflow_activation_job::Variant::DeliverStreamRecords(_)) + ) + }) + .flat_map(|j| delivered(j).3.into_iter().map(|b| b.to_vec())) + .collect(); + assert_eq!( + bodies, + vec![b"a".to_vec(), b"b".to_vec(), b"c".to_vec()], + "the replayed task must be handed the range it consumed; jobs were {:?}", + task.jobs + ); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + AppendStreamRecords { + stream_name: "out".to_string(), + records: vec![StreamRecord { + body: Some(b"accept".to_vec().into()), + ..Default::default() + }], + } + .into(), + ], + )) + .await + .unwrap(); + + core.shutdown().await; +} + +/// The reading task's range has to arrive in its own activation when History +/// comes in pages and a page boundary falls at that task. +/// +/// The range is recorded on the completion that closes the task, and the +/// paginator hands the machines updates cut at WFT started events. Whether the +/// boundary lands right after the reading task's started event or right after +/// its completion, the completion is outside the update the machines are +/// replaying from, and a lookahead reading only that update would find nothing. +#[rstest::rstest] +#[tokio::test] +async fn read_then_publish_replays_across_a_page_boundary( + #[values(12, 13)] second_page_end: usize, +) { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); // 3 + t.add_we_signaled("go", vec![]); + t.add_full_wf_task(); // 7 + // Task 2 subscribes. Its command event is what lets the paginator tell the + // reading task's sequence is complete once it sees the started event. + t.add_stream_subscribed("in", 0); + t.add_we_signaled("go", vec![]); + t.add_workflow_task_scheduled_and_started(); // 12 + // Task 3 reads [0,1) and publishes because of it. + let read_completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("in", 0, 1)]); + t.add_stream_records_appended("out", 0, 1); + t.add_workflow_task_scheduled_and_started(); // 16 + + let events = t.get_full_history_info().unwrap().into_events(); + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.history.as_mut().unwrap().events.truncate(3); + poll_resp.next_page_token = vec![1]; + // Two complete tasks past the previous started id fit in the first two + // pages, so the paginator would hand them over before fetching the third. + poll_resp.previous_started_event_id = 3; + poll_resp.add_stream_slice("in", read_completed, 0, &["go"]); + poll_resp.add_stream_slice("in", 0, 1, &["next"]); + + let second_page = history_page(&events[3..second_page_end], vec![2]); + let third_page = history_page(&events[second_page_end..], vec![]); + let mut mock_client = mock_worker_client(); + mock_client + .expect_get_workflow_execution_history() + .returning(move |_, _, token| match token.as_slice() { + [1] => Ok(second_page.clone()), + [2] => Ok(third_page.clone()), + other => panic!("unexpected page token {other:?}"), + }); + mock_client + .expect_fail_workflow_task() + .returning(|_, _, f| panic!("core rejected the replayed read-caused publish: {f:?}")); + + let mock = + MockPollCfg::from_resp_batches("wfid", t, [ResponseType::Raw(poll_resp.resp)], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![]); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![]); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_name_or_id: "in".to_string(), + start_offset: 0, + start_position: start(Position::Earliest(true)), + } + .into(), + ], + )) + .await + .unwrap(); + + // The reading task. Its signal and its range belong to the same activation. + let task = core.poll_workflow_activation().await.unwrap(); + assert!(task.is_replaying); + assert!( + task.jobs.iter().any(|j| matches!( + j.variant, + Some(workflow_activation_job::Variant::SignalWorkflow(_)) + )), + "expected the reading task's signal; jobs were {:?}", + task.jobs + ); + assert_eq!( + delivered_ranges(&task), + vec![("in".to_string(), 0, 1)], + "the replayed task must be handed the range it consumed; jobs were {:?}", + task.jobs + ); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + AppendStreamRecords { + stream_name: "out".to_string(), + records: vec![StreamRecord { + body: Some(b"accept".to_vec().into()), + ..Default::default() + }], + } + .into(), + ], + )) + .await + .unwrap(); + + // The live task gets only its own range. + let task = core.poll_workflow_activation().await.unwrap(); + assert!(!task.is_replaying); + assert_eq!(delivered_ranges(&task), vec![("in".to_string(), 1, 2)]); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + + core.shutdown().await; +} + /// 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. @@ -453,3 +964,561 @@ async fn a_repeat_subscribe_is_accepted() { .unwrap(); core.shutdown().await; } + +/// A task that consumes a range and issues no command still ran as its own +/// activation, so replay hands each such range over in its own activation +/// rather than collapsing the run of them into one, the way it does for +/// heartbeats that did nothing. +#[tokio::test] +async fn data_only_tasks_replay_one_range_per_activation() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let first = t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 0, 1)]); + t.add_workflow_task_scheduled_and_started(); + let second = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 1, 2)]); + t.add_workflow_task_scheduled_and_started(); + let third = t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 2, 3)]); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("s1", first, 0, &["a"]); + poll_resp.add_stream_slice("s1", second, 1, &["b"]); + poll_resp.add_stream_slice("s1", third, 2, &["c"]); + poll_resp.add_stream_slice("s1", 0, 3, &["d"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let core = mock_worker(build_mock_pollers(mock)); + + let mut per_activation = vec![]; + for _ in 0..4 { + let task = core.poll_workflow_activation().await.unwrap(); + per_activation.push(delivered_ranges(&task)); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + } + assert_eq!( + per_activation, + vec![ + vec![("s1".to_string(), 0, 1)], + vec![("s1".to_string(), 1, 2)], + vec![("s1".to_string(), 2, 3)], + vec![("s1".to_string(), 3, 4)], + ], + "each consumed range replays in the activation of the task that consumed it" + ); + core.shutdown().await; +} + +/// Ranges for several streams on one task arrive ordered by stream, however +/// the server laid them out. A workflow that waits on two streams with a +/// first-completed pattern would otherwise take whichever branch the server's +/// iteration order happened to pick, live and again differently on replay. +#[tokio::test] +async fn ranges_for_several_streams_arrive_in_stream_order() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let completed = t.add_workflow_task_completed_with_consumed_stream_ranges(vec![ + cursor("s2", 0, 1), + cursor("s1", 0, 1), + ]); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("s2", completed, 0, &["two"]); + poll_resp.add_stream_slice("s1", completed, 0, &["one"]); + poll_resp.add_stream_slice("s2", 0, 1, &["four"]); + poll_resp.add_stream_slice("s1", 0, 1, &["three"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!( + delivered_ranges(&task), + vec![("s1".to_string(), 0, 1), ("s2".to_string(), 0, 1)], + "re-supplied ranges follow stream order" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!( + delivered_ranges(&task), + vec![("s1".to_string(), 1, 2), ("s2".to_string(), 1, 2)], + "live ranges follow stream order" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + core.shutdown().await; +} + +/// A recorded range that observed nothing is rebuilt from the cursor alone. The +/// event already says everything the workflow needs, so replay does not depend +/// on the server sending an empty slice back for it. +#[tokio::test] +async fn an_empty_recorded_range_replays_without_a_slice_from_the_server() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 4, 4)]); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.add_stream_slice("s1", 0, 4, &["e"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![("s1".to_string(), 4, 4)]); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![("s1".to_string(), 4, 5)]); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + core.shutdown().await; +} + +/// A re-supplied slice is checked against the cursor it claims to satisfy. The +/// event is the record of what the task saw, so bytes covering other offsets +/// would replay the task on different input than it ran on. Both sides of that +/// comparison come from the server, so the task fails as the worker's failure. +#[tokio::test] +async fn a_resupplied_slice_that_disagrees_with_the_record_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 0, 2)]); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + // One record where the event says two. + poll_resp.add_stream_slice("s1", completed, 0, &["a"]); + + let (core, failures) = worker_expecting_one_failure( + t, + ResponseType::Raw(poll_resp.resp), + WorkflowTaskFailedCause::WorkflowWorkerUnhandledFailure, + ); + + // The mismatch is found while the poll response is applied, so the first + // activation is already the eviction. + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// A run that stays cached was handed its range as the task ran. When the next +/// task arrives, the completion recording that range is in its history, and the +/// server may re-supply the range tagged with it, as it would for a worker that +/// lost the run. This worker did not, so it must not process the same records +/// twice. +#[tokio::test] +async fn a_cached_run_is_not_handed_its_own_range_again() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let completed = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 0, 2)]); + t.add_workflow_task_scheduled_and_started(); + + let mut first_poll = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::ToTaskNum(1)); + first_poll.add_stream_slice("s1", 0, 0, &["a", "b"]); + let mut second_poll = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::OneTask(2)); + second_poll.add_stream_slice("s1", completed, 0, &["a", "b"]); + second_poll.add_stream_slice("s1", 0, 2, &["c"]); + + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ + ResponseType::Raw(first_poll.resp), + ResponseType::Raw(second_poll.resp), + ], + mock_worker_client(), + ); + let mut mock = build_mock_pollers(mock); + mock.worker_cfg(|wc| wc.max_cached_workflows = 1); + let core = mock_worker(mock); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![("s1".to_string(), 0, 2)]); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + + let task = core.poll_workflow_activation().await.unwrap(); + assert!(!task.is_replaying); + assert_eq!( + delivered_ranges(&task), + vec![("s1".to_string(), 2, 3)], + "only the new range; the first one was delivered live to this worker" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + core.shutdown().await; +} + +/// History says a task consumed a range with content and the server sent no +/// bytes for it. The bytes only travel on the response that carried the task, +/// so the lookahead fails the task as soon as it sees the range, before the +/// workflow runs on less input than it ran on. A sticky task handed to a worker +/// that no longer holds the run is the case that reaches this, and the server's +/// retry on the normal queue carries the records. The task is failed as the +/// worker's failure, not the workflow's: nothing the workflow did was wrong. +#[tokio::test] +async fn a_missing_resupply_for_a_consumed_range_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("s1", 0, 2)]); + t.add_workflow_task_scheduled_and_started(); + + let (core, failures) = worker_expecting_one_failure( + t, + ResponseType::AllHistory, + WorkflowTaskFailedCause::WorkflowWorkerUnhandledFailure, + ); + + // Found while the poll response is applied, so the first activation is + // already the eviction. + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// A task that read two streams and a response that re-supplies only one of +/// them. The recorded range is the whole of what the task ran on, so a partial +/// re-supply leaves the workflow as short of input as none at all, and it is +/// the same worker failure: the history and the response both come from the +/// server, and the retry on the normal task queue carries every stream. +#[tokio::test] +async fn a_partial_resupply_for_a_consumed_task_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let completed = t.add_workflow_task_completed_with_consumed_stream_ranges(vec![ + cursor("s1", 0, 1), + cursor("s2", 0, 1), + ]); + t.add_workflow_task_scheduled_and_started(); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + // Only one of the two subscribed streams came back. + poll_resp.add_stream_slice("s1", completed, 0, &["one"]); + + let (core, failures) = worker_expecting_one_failure( + t, + ResponseType::Raw(poll_resp.resp), + WorkflowTaskFailedCause::WorkflowWorkerUnhandledFailure, + ); + + // Found while the poll response is applied, so the first activation is + // already the eviction. + // The reason travels in the message, since an eviction for a failure found + // before any activation ran reports none of its own. + let task = core.poll_workflow_activation().await.unwrap(); + let evict = eviction(&task); + assert!( + evict.message.contains( + "stream s2 was consumed from offset 0 to 1, but the server sent no \ + messages for it" + ), + "eviction did not name the stream that went missing: {evict:?}" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// A legacy query for a run this worker no longer holds arrives on the sticky +/// queue with partial history and no records, and the history the worker +/// fetches itself carries none either. Answering it from a replay on less +/// input would be wrong, and failing it would end the query: the server +/// retries a query it hears nothing about on the normal queue, where the +/// records travel with it. So the query goes unanswered, no task is failed, +/// and the run is given up so the retry starts from history. +#[tokio::test] +async fn a_legacy_query_owed_records_it_was_not_sent_goes_unanswered() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + // Task 1 subscribes and consumes nothing, so the missing range is only + // reached once the workflow has run its first activation. + t.add_workflow_task_completed(); + t.add_stream_subscribed("in", 0); + t.add_workflow_task_scheduled_and_started(); + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("in", 0, 2)]); + t.add_stream_records_appended("out", 0, 2); + + let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory); + poll_resp.resp.query = Some(WorkflowQuery { + query_type: "trace".to_string(), + query_args: None, + header: None, + }); + // No slices: the mock plays the sticky queue, and the defaults of zero + // expected task failures and zero legacy query responses are the assertion. + let mock = MockPollCfg::from_resp_batches( + "wfid", + t, + [ResponseType::Raw(poll_resp.resp)], + mock_worker_client(), + ); + let mut mock = build_mock_pollers(mock); + mock.worker_cfg(|wc| { + wc.max_cached_workflows = 10; + wc.ignore_evicts_on_shutdown = false; + }); + let core = mock_worker(mock); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![]); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_name_or_id: "in".to_string(), + start_offset: 0, + start_position: start(Position::Earliest(true)), + } + .into(), + ], + )) + .await + .unwrap(); + + // The missing range is found while the next task is applied. The run is + // evicted as a fetch failure would evict it, and nothing is reported. + let task = core.poll_workflow_activation().await.unwrap(); + let evict = eviction(&task); + assert_eq!(evict.reason(), EvictionReason::PaginationOrHistoryFetch); + assert!( + evict + .message + .contains("but the server sent no records for it"), + "eviction did not name the missing range: {evict:?}" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + core.shutdown().await; +} + +/// A slice as the server re-supplies it: the records of one recorded range, +/// tagged with the completion that recorded it. +fn replay_slice( + stream_id: &str, + completed_event_id: i64, + from: i64, + bodies: &[&str], +) -> StreamSlice { + StreamSlice { + stream_id: stream_id.to_string(), + from_offset: from, + to_offset: from + bodies.len() as i64, + records: bodies + .iter() + .map(|b| StreamRecord { + body: Some(b.as_bytes().to_vec().into()), + ..Default::default() + }) + .collect(), + workflow_task_completed_event_id: completed_event_id, + ..Default::default() + } +} + +/// A publish of one record per body. +fn publish(stream_name: &str, bodies: &[&str]) -> AppendStreamRecords { + AppendStreamRecords { + stream_name: stream_name.to_string(), + records: bodies + .iter() + .map(|b| StreamRecord { + body: Some(b.as_bytes().to_vec().into()), + ..Default::default() + }) + .collect(), + } +} + +fn eviction(task: &WorkflowActivation) -> &RemoveFromCache { + match task.jobs.as_slice() { + [ + WorkflowActivationJob { + variant: Some(workflow_activation_job::Variant::RemoveFromCache(evict)), + }, + ] => evict, + other => panic!("expected an eviction, got {other:?}"), + } +} + +/// A recorded read-then-publish workflow: each task consumes one record and +/// publishes because of it, the second one also completes the run. +fn read_then_publish_history() -> (TestHistoryBuilder, i64, i64) { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let first = t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("in", 0, 1)]); + t.add_stream_records_appended("out", 0, 1); + t.add_workflow_task_scheduled_and_started(); + let second = + t.add_workflow_task_completed_with_consumed_stream_ranges(vec![cursor("in", 1, 2)]); + t.add_stream_records_appended("out", 1, 2); + t.add_workflow_execution_completed(); + (t, first, second) +} + +/// A replay worker fed one history. The feeder is handed back so the history +/// stream stays open until the test drops it, as a language replayer keeps it +/// open: once the stream ends the worker closes, and an eviction still owed +/// for the last history would be lost to the shutdown. +async fn replay_worker(history: HistoryForReplay) -> (Worker, HistoryFeeder) { + let (feeder, stream) = HistoryFeeder::new(1); + feeder.feed(history).await.unwrap(); + let core = init_replay_worker(ReplayWorkerInput::new( + test_worker_cfg().build().unwrap(), + stream, + )) + .unwrap(); + (core, feeder) +} + +/// A history pushed for replay can carry the ranges its tasks consumed, in the +/// shape the server re-supplies them. The replay worker puts them on its +/// synthetic poll response, so the ordinary delivery path runs: each range +/// reaches the activation of the task that consumed it, and the publish that +/// task reissues is matched against its event. +#[tokio::test] +async fn a_pushed_history_replays_with_the_slices_it_carries() { + let (t, first, second) = read_then_publish_history(); + let history = HistoryForReplay::new(t.get_full_history_info().unwrap(), "wfid") + .with_stream_slices([ + replay_slice("in", first, 0, &["go"]), + replay_slice("in", second, 1, &["stop"]), + ]); + let (core, feeder) = replay_worker(history).await; + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![("in".to_string(), 0, 1)]); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish("out", &["accept"]).into()], + )) + .await + .unwrap(); + + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(delivered_ranges(&task), vec![("in".to_string(), 1, 2)]); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + publish("out", &["done"]).into(), + CompleteWorkflowExecution { result: None }.into(), + ], + )) + .await + .unwrap(); + + // Replay is over and the worker lets the run go; a nondeterminism eviction + // would say the reissued publishes did not match their events. + let task = core.poll_workflow_activation().await.unwrap(); + assert_eq!(eviction(&task).reason(), EvictionReason::LangRequested); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + drop(feeder); + core.shutdown().await; +} + +/// A slice attached to a pushed history is held to the same check as one the +/// server sends: offsets other than the recorded range fail the task rather +/// than replay it on different input. +#[tokio::test] +async fn a_pushed_history_with_a_wrong_slice_fails_as_a_worker_failure() { + let (t, first, second) = read_then_publish_history(); + let history = HistoryForReplay::new(t.get_full_history_info().unwrap(), "wfid") + .with_stream_slices([ + // Two records where the event says one. + replay_slice("in", first, 0, &["go", "extra"]), + replay_slice("in", second, 1, &["stop"]), + ]); + let (core, feeder) = replay_worker(history).await; + + // The mismatch is found while the poll response is applied, so the first + // activation is already the eviction. The task is failed as the worker's + // failure (the mock-client test above checks the cause); the eviction + // carries the reason in its message, since an eviction for a failure found + // before any activation ran reports no reason of its own. + let task = core.poll_workflow_activation().await.unwrap(); + let evict = eviction(&task); + assert!( + evict.message.contains( + "Event 4 records that stream in was consumed from offset 0 to 1, but the server \ + sent offsets 0 to 2 for it" + ), + "eviction did not name the mismatch: {evict:?}" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + drop(feeder); + core.shutdown().await; +} + +/// A recorded range with content and no slice for it is the same failure. A +/// language replayer that has no store to fetch from cannot replay a consuming +/// workflow, and the task says so before the workflow runs on less input. +#[tokio::test] +async fn a_pushed_history_without_slices_for_a_consumed_range_fails_as_a_worker_failure() { + let (t, _, _) = read_then_publish_history(); + let history = HistoryForReplay::new(t.get_full_history_info().unwrap(), "wfid"); + let (core, feeder) = replay_worker(history).await; + + // Found while the poll response is applied, so the first activation is + // already the eviction, which carries the reason in its message. + let task = core.poll_workflow_activation().await.unwrap(); + let evict = eviction(&task); + assert!( + evict.message.contains( + "Event 4 records that stream in was consumed from offset 0 to 1, but the server \ + sent no records for it" + ), + "eviction did not name the missing range: {evict:?}" + ); + core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id)) + .await + .unwrap(); + drop(feeder); + core.shutdown().await; +} diff --git a/crates/sdk-core/src/replay/history_builder.rs b/crates/sdk-core/src/replay/history_builder.rs index d5cd59df4..4e827af97 100644 --- a/crates/sdk-core/src/replay/history_builder.rs +++ b/crates/sdk-core/src/replay/history_builder.rs @@ -24,6 +24,7 @@ use temporalio_common::protos::{ failure::v1::{CanceledFailureInfo, Failure, failure}, history::v1::{history_event::Attributes, *}, notification::v1::Notification, + stream::v1::StreamRange, taskqueue::v1::TaskQueue, update, update::v1::outcome, @@ -142,6 +143,22 @@ impl TestHistoryBuilder { self.previous_task_completed_id = id; } + /// Add a workflow task completed event recording the stream offsets that + /// task consumed. Only the range is in History; the payloads come back from + /// the server on the poll response. + pub fn add_workflow_task_completed_with_consumed_stream_ranges( + &mut self, + cursors: Vec, + ) -> i64 { + let id = self.add(WorkflowTaskCompletedEventAttributes { + scheduled_event_id: self.workflow_task_scheduled_event_id, + consumed_stream_ranges: cursors, + ..Default::default() + }); + self.previous_task_completed_id = 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 { diff --git a/crates/sdk-core/src/test_help/integ_helpers.rs b/crates/sdk-core/src/test_help/integ_helpers.rs index 0e446be95..cdfcc13da 100644 --- a/crates/sdk-core/src/test_help/integ_helpers.rs +++ b/crates/sdk-core/src/test_help/integ_helpers.rs @@ -50,10 +50,11 @@ use temporalio_common::{ workflow_completion::WorkflowActivationCompletion, }, temporal::api::{ - common::v1::WorkflowExecution, + common::v1::{Payload, WorkflowExecution}, enums::v1::WorkflowTaskFailedCause, failure::v1::Failure, protocol::{self, v1::message}, + stream::v1::{StreamRecord, StreamSlice}, update, workflowservice::v1::{ DescribeNamespaceResponse, PollActivityTaskQueueResponse, @@ -913,6 +914,17 @@ pub trait PollWFTRespExt { update_id: impl ToString, after_event_id: i64, ) -> update::v1::Request; + + /// Attach a range of a stream. Passing a `completed_event_id` of zero makes + /// it the range for the task about to run; anything else re-supplies what + /// the task closed by that event consumed. + fn add_stream_slice( + &mut self, + stream_id: impl ToString, + completed_event_id: i64, + from_offset: i64, + bodies: &[&str], + ); } impl PollWFTRespExt for PollWorkflowTaskQueueResponse { @@ -947,6 +959,32 @@ impl PollWFTRespExt for PollWorkflowTaskQueueResponse { }); upd_req_body } + + fn add_stream_slice( + &mut self, + stream_id: impl ToString, + completed_event_id: i64, + from_offset: i64, + bodies: &[&str], + ) { + self.stream_slices.push(StreamSlice { + stream_id: stream_id.to_string(), + from_offset, + to_offset: from_offset + bodies.len() as i64, + records: bodies + .iter() + .map(|b| StreamRecord { + body: Some(Payload { + data: b.as_bytes().to_vec(), + ..Default::default() + }), + ..Default::default() + }) + .collect(), + workflow_task_completed_event_id: completed_event_id, + ..Default::default() + }); + } } pub fn hist_to_poll_resp( From bb5444ed2331b091d319a69587626b4cab050935 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Sat, 3 Oct 2026 02:00:42 -0700 Subject: [PATCH 3/3] Added changelog entries for native streams. --- crates/sdk-core/CHANGELOG.md | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/crates/sdk-core/CHANGELOG.md b/crates/sdk-core/CHANGELOG.md index 62647f78f..f54aa2004 100644 --- a/crates/sdk-core/CHANGELOG.md +++ b/crates/sdk-core/CHANGELOG.md @@ -38,6 +38,17 @@ relevant information. command and end a subscription with `UnsubscribeNotificationChannel`. The notifications the server folds for a Workflow Task arrive as one `NotificationsReceived` activation job, taken from the task's scheduled event, so replay yields the same job at the same point. +* Workflows can subscribe to server-side streams and append batches of records to them with the + `SubscribeStream` and `AppendStreamRecords` commands. Consumed ranges reach the workflow as + `DeliverStreamRecords` activation jobs, and replay hands each recorded range back in the + activation of the task that consumed it. +* A history fed to a replay worker can carry the stream records its tasks consumed + (`HistoryForReplay::with_stream_slices`), so a language replayer that fetched them from the + stream service can replay a consuming workflow. History alone holds only the offsets. +* A task whose history records a consumed range with content that the response carried no + records for fails before the workflow runs, rather than after it ran on less input. A legacy + query dispatched that way to a worker that no longer holds the run goes unanswered, so the + server retries it on the normal task queue, where the records travel with it. ### Fixed * Task-poll targets no longer decrease after cancelled or timed-out polls. Affected pollers still