diff --git a/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto index c8c6859ea..aadc8f89d 100644 --- a/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto @@ -339,8 +339,11 @@ message Command { // carrying the offset range and none of the payload; it schedules no further // work. message AppendStreamRecordsCommandAttributes { - // 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; // Stored in order. The server sets `producer_id` to empty on each record, // because the owning Workflow is the producer here. repeated temporal.api.stream.v1.StreamRecord records = 2; @@ -353,16 +356,22 @@ message AppendStreamRecordsCommandAttributes { // 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 SubscribeStreamCommandAttributes { - // Stream to consume. A stream in another execution is addressed by its id; - // one this Workflow owns is addressed by the name it was published under. - // The server resolves an owned name first and falls back to a standalone - // id, so a Workflow that owns a stream under this name cannot reach a - // standalone stream with the same id. - string stream_id = 1; - // Where to start. Negative means from wherever the stream is when the - // subscription is registered, which the server resolves and records so - // replay does not resolve it again. + // 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. Read only when `start_position` + // is unset. A negative value is refused: the head of the stream is asked + // for with `start_position.tail`. int64 start_offset = 2; + // Where to start. The server resolves it once, when it registers the + // subscription, and records the resolved absolute offset on the subscribed + // event, so replay does not resolve it again. Setting it together with a + // non-zero `start_offset` fails the command. + temporal.api.stream.v1.StreamStartPosition start_position = 3; } // Makes the Workflow a listener of a notification channel for this run. The diff --git a/crates/protos/protos/api_upstream/temporal/api/history/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/history/v1/message.proto index 31636d401..553683100 100644 --- a/crates/protos/protos/api_upstream/temporal/api/history/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/history/v1/message.proto @@ -390,8 +390,11 @@ message WorkflowTaskCompletedEventAttributes { // Recorded on every task where a subscription is active, including when it // observed nothing: an empty range is a fact replay must reproduce, and // omitting it would let replay deliver records the Workflow did not have. - // Numbered 20 to leave 14 through 19 free for fields added on the main line. repeated temporal.api.stream.v1.StreamRange consumed_stream_ranges = 20; + + // Held for fields added on the main line, so a rebase does not land one of + // them on a number this fork already writes. + reserved 14 to 19; } message WorkflowTaskTimedOutEventAttributes { @@ -972,7 +975,9 @@ message WorkflowStreamSubscribedEventAttributes { // The WorkflowTaskCompleted event of the task whose command created this // subscription. int64 workflow_task_completed_event_id = 1; - // Stream the Workflow subscribed to. + // The stream the Workflow subscribed to, as the command addressed it: + // either the name of a stream this Workflow owns or the id of one in + // another execution. string stream_id = 2; // The offset the subscription actually starts from. Resolved by the server // when the subscription is registered and recorded here, so replay reads @@ -1005,14 +1010,20 @@ message WorkflowStreamRecordsAppendedEventAttributes { // The WorkflowTaskCompleted event of the task whose command appended this // batch. int64 workflow_task_completed_event_id = 1; - // Stream the Workflow appended to. + // Name of the stream the Workflow appended to. string stream_id = 2; - // Offset the first record of the batch landed at. - int64 first_offset = 3; - // How many records the batch held. With first_offset this names the range - // without carrying any of it, which is what keeps this event a fixed size - // no matter how large the batch or its payloads are. - int64 record_count = 4; + // Inclusive. Same range vocabulary as StreamRange and StreamSlice, so a + // reader does not have to remember which of the three counts and which + // bounds. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 from_offset = 3; + // Exclusive. With from_offset this names the range without carrying any of + // it, which is what keeps this event a fixed size no matter how large the + // batch or its payloads are. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 to_offset = 4; } message WorkflowExecutionUpdateAcceptedEventAttributes { diff --git a/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto index 9f0c88dfb..e6a8fab2d 100644 --- a/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto @@ -14,8 +14,13 @@ import "temporal/api/common/v1/message.proto"; // One entry in a stream. The record is the wire format: stores keep it // serialized as is and readers in every language decode the same bytes. message StreamRecord { - // The value the producer published, stored as sent. A payload codec - // applies here as it does to any other payload. + // The value the producer published, stored as sent. + // + // A payload codec applies on the paths this API owns: the append command on + // RespondWorkflowTaskCompleted, and the slices on PollWorkflowTaskQueue. + // Records a producer writes or reads through the stream service take a + // different path, whose messages are not part of this API yet and so are + // outside what a codec-applying proxy walks. temporal.api.common.v1.Payload body = 1; // Producer-supplied provenance, stored as sent. map metadata = 2; @@ -28,8 +33,9 @@ message StreamRecord { // The producer's attempt. Readers treat a later attempt by the same // producer as superseding what the earlier one wrote. int64 attempt = 6; - // The producer's position within its attempt, or -1 when unnumbered. - // Stored as sent; the server does not assign, validate or order by it. + // The producer's position within its attempt, zero when it does not number + // its records. Stored as sent; the server does not assign, validate or + // order by it, and the stream's own offsets are what order a read. int64 sequence = 7; } @@ -37,14 +43,21 @@ message StreamRecord { // offsets it covers. The offsets are what History records; the records // themselves are never written to History. message StreamSlice { + // The stream, as the subscribing command addressed it: either the name of + // a stream the consuming Workflow owns or the id of one in another + // execution. string stream_id = 1; // Run id of the execution that owns the stream. Set on both a slice for the // task being started and a re-supplied one. string run_id = 2; // Inclusive. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) int64 from_offset = 3; // Exclusive. Equal to from_offset when the subscription observed nothing, // which is a fact replay has to reproduce rather than an absence of one. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) int64 to_offset = 4; repeated StreamRecord records = 5; // The WorkflowTaskCompleted event whose consumed_stream_ranges recorded @@ -62,13 +75,39 @@ message StreamSlice { // WorkflowTaskCompleted so History grows with Workflow Tasks rather than with // records. message StreamRange { + // The stream, as the subscribing command addressed it: either the name of + // a stream the consuming Workflow owns or the id of one in another + // execution. string stream_id = 1; // Inclusive. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) int64 from_offset = 2; // Exclusive. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) int64 to_offset = 3; } +// Where a new subscription or read begins. The server resolves it against the +// stream as it stands in the same transaction that registers the reader, so +// the result does not race with appends or truncation, and records the +// resolved absolute offset. +message StreamStartPosition { + oneof position { + // Absolute and inclusive. Refused when below the stream's floor. + int64 offset = 1; + // The last N records the stream holds, or all of them when it holds + // fewer. Counts records of every kind. Must be positive. + int64 last_n = 2; + // The oldest record the stream still holds. Must be true. + bool earliest = 3; + // Only records appended after registration: the stream's head offset. + // Must be true. + bool tail = 4; + } +} + // What a record means to a reader. Kept on the record itself so every store // and every language reads it the same way without a private envelope. enum StreamRecordKind { 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 03235adff..e0126739e 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 @@ -13,6 +13,7 @@ import "google/protobuf/empty.proto"; 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/enums/v1/failed_cause.proto"; import "temporal/api/enums/v1/workflow.proto"; import "temporal/api/notification/v1/message.proto"; @@ -156,6 +157,11 @@ message WorkflowActivationJob { // Runtime-internal: encode the terminal boundary for a marker Core is about to write. // Runs no user code. FinalizeExternalStreams finalize_external_streams = 20; + // The number below is shared with the native stream tree, which leaves 17 + // to 20 to the jobs above, and must not be reused. + // + // 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 @@ -231,6 +237,25 @@ message FinalizeExternalStreams { coresdk.external_data.ParkReason reason = 3; } +// Hand a workflow the next range of a stream it subscribed to. +// +// The range is delivered once, on the task the server decided it belongs to, +// and the offsets it covered are recorded in History rather than the payloads. +// On replay the server re-supplies the same range by reading the stream again, +// so this job appears at the same point with the same contents both times. +// +// An empty range is still delivered: a task where the subscription saw nothing +// is a fact replay has to reproduce, not an absence of one. +message DeliverStreamRecords { + // Id of the stream this range came from. + string stream_id = 1; + // Inclusive. + int64 from_offset = 2; + // Exclusive. Equal to from_offset when the subscription saw nothing. + int64 to_offset = 3; + repeated temporal.api.stream.v1.StreamRecord records = 4; +} + // 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 d73cdfa1c..b94e8618f 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"; @@ -61,6 +62,10 @@ message WorkflowCommand { ExternalStreamFinalized external_stream_finalized = 26; WorkflowOutputStreamCommit workflow_output_stream_commit = 27; WorkflowOutputStreamBuffered workflow_output_stream_buffered = 28; + // The two numbers below are shared with the native stream tree, which + // leaves 23 to 28 to the commands above, and must not be reused. + SubscribeStream subscribe_stream = 29; + AppendStreamRecords append_stream_records = 30; SubscribeNotificationChannel subscribe_notification_channel = 31; UnsubscribeNotificationChannel unsubscribe_notification_channel = 32; WorkflowStreamChannels workflow_stream_channels = 33; @@ -186,6 +191,45 @@ message WorkflowOutputStreamBuffered { google.protobuf.Duration max_publish_latency = 1; } +// 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 { + // 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; +} + +// 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 { + // 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. 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 { // 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 e0ff9a7f1..450e8d448 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1374,6 +1374,13 @@ pub mod coresdk { fin.reason() ) } + workflow_activation_job::Variant::DeliverStreamRecords(d) => { + write!( + f, + "DeliverStreamRecords({}, {}..{})", + d.stream_id, d.from_offset, d.to_offset + ) + } workflow_activation_job::Variant::NotificationsReceived(n) => { write!(f, "NotificationsReceived({})", n.notifications.len()) } @@ -1586,6 +1593,23 @@ pub mod coresdk { } } + impl Display for SubscribeStream { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "SubscribeStream({})", self.stream_name_or_id) + } + } + + impl Display for AppendStreamRecords { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!( + f, + "AppendStreamRecords({}, {} records)", + self.stream_name, + self.records.len() + ) + } + } + impl Display for StartTimer { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "StartTimer({})", self.seq) @@ -2094,6 +2118,29 @@ pub mod temporal { } } + impl From for command::Attributes { + fn from(s: workflow_commands::AppendStreamRecords) -> Self { + Self::AppendStreamRecordsCommandAttributes( + AppendStreamRecordsCommandAttributes { + stream_name: s.stream_name, + records: s.records, + }, + ) + } + } + + impl From for command::Attributes { + fn from(s: workflow_commands::SubscribeStream) -> Self { + Self::SubscribeStreamCommandAttributes( + SubscribeStreamCommandAttributes { + stream_name_or_id: s.stream_name_or_id, + start_offset: s.start_offset, + start_position: s.start_position, + }, + ) + } + } + impl From for Attributes { fn from(s: workflow_commands::SubscribeNotificationChannel) -> Self { Self::SubscribeNotificationChannelCommandAttributes( diff --git a/crates/sdk-core/CHANGELOG.md b/crates/sdk-core/CHANGELOG.md index a4557e377..782186e8f 100644 --- a/crates/sdk-core/CHANGELOG.md +++ b/crates/sdk-core/CHANGELOG.md @@ -33,6 +33,19 @@ relevant information. ## Unreleased +### Added +* 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 retain their slot during backoff, while resource-exhaustion errors still reduce the target. @@ -71,7 +84,6 @@ relevant information. metrics now carry a `failure_reason` attribute. Each is now split into one time series per reason, which may affect existing dashboards. * Workflow task completions larger than the gRPC request size limit are now paginated automatically when the namespace supports it. Paginated workflow task completions require Temporal Server 1.32.0 or later. - ### Breaking Changes :boom: * The following types are now non-exhaustive: `Priority`, `WorkerDeploymentVersion`, `WorkerCallbacks`, `WorkflowExecutionInfo`, `ActivityCloseTimeouts`, diff --git a/crates/sdk-core/src/core_tests/mod.rs b/crates/sdk-core/src/core_tests/mod.rs index 0c3b4174f..66d2f5979 100644 --- a/crates/sdk-core/src/core_tests/mod.rs +++ b/crates/sdk-core/src/core_tests/mod.rs @@ -3,6 +3,7 @@ mod event_groups; mod external_streams; 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..8b1423679 --- /dev/null +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -0,0 +1,1524 @@ +//! 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, 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::{ + Arc, + atomic::{AtomicUsize, Ordering}, +}; +use temporalio_common::protos::{ + coresdk::{ + 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}, + 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 +/// 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 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)) +} + +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 +// 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_name_or_id: "s1".to_string(), + start_offset: 0, + start_position: start(Position::Tail(true)), + } + .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_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 + .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. 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_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(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"); + assert_eq!(a.start_offset, 0); + assert_eq!(a.start_position, expected); + } + 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_name_or_id: "s1".to_string(), + start_offset: 0, + start_position, + } + .into(), + ], + )) + .await + .unwrap(); + 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(), + records: vec![ + StreamRecord { + body: Some(b"one".to_vec().into()), + ..Default::default() + }, + StreamRecord { + body: Some(b"two".to_vec().into()), + ..Default::default() + }, + ], + } +} + +/// 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. +#[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 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(); + 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_name: "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 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(); + 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_name_or_id: "s2".to_string(), + start_offset: 0, + start_position: start(Position::Tail(true)), + } + .into(), + ], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// 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. +#[tokio::test] +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(); + 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_name_or_id: "s1".to_string(), + start_offset: 0, + start_position: start(Position::Earliest(true)), + } + .into(), + ], + )) + .await + .unwrap(); + core.shutdown().await; +} + +/// 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_accepted() { + 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_name_or_id: "s1".to_string(), + start_offset: 0, + start_position: start(Position::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; +} + +/// 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/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/history_builder.rs b/crates/sdk-core/src/replay/history_builder.rs index 616eaf68d..fb5e719fd 100644 --- a/crates/sdk-core/src/replay/history_builder.rs +++ b/crates/sdk-core/src/replay/history_builder.rs @@ -26,6 +26,7 @@ use temporalio_common::protos::{ enums::v1::{EventType, TaskQueueKind, WorkflowTaskFailedCause}, failure::v1::{CanceledFailureInfo, Failure, failure}, history::v1::{history_event::Attributes, *}, + stream::v1::StreamRange, taskqueue::v1::TaskQueue, update, update::v1::outcome, @@ -169,6 +170,49 @@ 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 { + 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. The range is + /// half-open, as it is on the event. + pub fn add_stream_records_appended( + &mut self, + stream_id: &str, + 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(), + from_offset, + to_offset, + }; + self.add(attrs) + } + /// Add a workflow task timed out event. pub fn add_workflow_task_timed_out(&mut self) { let attrs = WorkflowTaskTimedOutEventAttributes { 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/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( diff --git a/crates/sdk-core/src/worker/workflow/history_update.rs b/crates/sdk-core/src/worker/workflow/history_update.rs index 109edb4a3..4e7ca5de9 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, }) } @@ -720,6 +750,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(); @@ -738,8 +772,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 && !carries_notifications(next_next_event) { @@ -785,7 +830,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; @@ -796,7 +844,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/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..e7c2ad302 --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs @@ -0,0 +1,170 @@ +use super::{ + NewMachineWithCommand, StateMachine, TransitionResult, fsm, workflow_machines::MachineResponse, +}; +use crate::worker::workflow::{ + WFMachinesError, fatal, + machines::{EventInfo, HistEventData, WFMachinesAdapter}, + nondeterminism, +}; +use std::{cell::RefCell, rc::Rc}; +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. +/// +/// 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_name: String, + record_count: i64, + 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 DefaultStreamNameRef = 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, + default_stream_name: DefaultStreamNameRef, +) -> NewMachineWithCommand { + let sm = AppendStreamRecordsMachine::from_parts( + Created {}.into(), + SharedState { + stream_name: lang_cmd.stream_name.clone(), + record_count: lang_cmd.records.len() as i64, + default_stream_name, + }, + ); + 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 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_name.clone() + }; + // 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 {:?}", + recorded_count, + attrs.stream_id, + dat.record_count, + expected + )) + } + } +} + +#[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(fatal!( + "AppendStreamRecords does not use state machine commands" + )) + } +} + +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 4698cabb3..466be535a 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; @@ -17,6 +18,7 @@ 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 unsubscribe_notification_channel_state_machine; mod update_state_machine; @@ -34,6 +36,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; @@ -51,6 +54,7 @@ use std::{ fmt::{Debug, Display}, }; use subscribe_notification_channel_state_machine::SubscribeNotificationChannelMachine; +use subscribe_stream_state_machine::SubscribeStreamMachine; use temporalio_common::{ fsm_trait::{StateMachine, TransitionResult}, protos::temporal::api::{ @@ -87,6 +91,8 @@ enum Machines { WorkflowTaskMachine, UpsertSearchAttributesMachine, ModifyWorkflowPropertiesMachine, + SubscribeStreamMachine, + AppendStreamRecordsMachine, UpdateMachine, NexusOperationMachine, SubscribeNotificationChannelMachine, 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..2e547d0d3 --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -0,0 +1,141 @@ +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 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_name_or_id: String, +} + +/// 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. +pub(super) fn subscribe_stream(lang_cmd: SubscribeStream) -> NewMachineWithCommand { + let sm = SubscribeStreamMachine::from_parts( + Created {}.into(), + SharedState { + stream_name_or_id: lang_cmd.stream_name_or_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_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_name_or_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(fatal!( + "SubscribeStream does not use state machine commands" + )) + } +} + +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/transition_coverage.rs b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs index db12b886f..da382c11b 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, @@ -78,7 +79,7 @@ mod machine_coverage_report { nexus_operation_state_machine::NexusOperationMachine, patch_state_machine::PatchMachine, signal_external_state_machine::SignalExternalMachine, subscribe_notification_channel_state_machine::SubscribeNotificationChannelMachine, - timer_state_machine::TimerMachine, + subscribe_stream_state_machine::SubscribeStreamMachine, timer_state_machine::TimerMachine, unsubscribe_notification_channel_state_machine::UnsubscribeNotificationChannelMachine, update_state_machine::UpdateMachine, upsert_search_attributes_state_machine::UpsertSearchAttributesMachine, @@ -122,6 +123,8 @@ mod machine_coverage_report { let mut update = UpdateMachine::visualizer().to_owned(); let mut nexus = NexusOperationMachine::visualizer().to_owned(); let mut external_stream = ExternalStreamMachine::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(); let mut unsubscribe_channel = UnsubscribeNotificationChannelMachine::visualizer().to_owned(); @@ -153,6 +156,12 @@ mod machine_coverage_report { m @ "UpdateMachine" => cover_transitions(m, &mut update, coverage), m @ "NexusOperationMachine" => cover_transitions(m, &mut nexus, coverage), m @ "ExternalStreamMachine" => cover_transitions(m, &mut external_stream, coverage), + m @ "SubscribeStreamMachine" => { + cover_transitions(m, &mut subscribe_stream, coverage) + } + m @ "AppendStreamRecordsMachine" => { + cover_transitions(m, &mut append_stream_records, coverage) + } m @ "SubscribeNotificationChannelMachine" => { cover_transitions(m, &mut subscribe_channel, coverage) } 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 44e94ad58..b80f0a95a 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,17 @@ mod local_acts; use super::{ Machines, NewMachineWithCommand, TemporalStateMachine, + 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, 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, + 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_notification_channel_state_machine::subscribe_notification_channel, + subscribe_stream_state_machine::subscribe_stream, timer_state_machine::new_timer, unsubscribe_notification_channel_state_machine::unsubscribe_notification_channel, upsert_search_attributes_state_machine::upsert_search_attrs, @@ -82,6 +86,7 @@ use temporalio_common::{ notification::v1::Notification, protocol::v1::{Message as ProtocolMessage, message::SequencingId}, sdk::v1::WorkflowTaskCompletedMetadata, + stream::v1::{StreamRange, StreamSlice}, }, }, worker::WorkerDeploymentVersion, @@ -98,6 +103,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, /// Protocol messages that have yet to be processed for the current WFT. protocol_msgs: Vec, /// Reserved external stream wake Signals seen in history, decoded and suppressed from user @@ -197,6 +225,10 @@ pub(crate) struct WorkflowMachines { /// Contains extra local-activity related data local_activity_data: LocalActivityData, + /// 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_name: DefaultStreamNameRef, + /// The workflow that is being driven by this instance of the machines drive_me: DrivenWorkflow, @@ -300,7 +332,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 { @@ -310,6 +342,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, replaying, metrics: basics.metrics, // In an ideal world one could say ..Default::default() here and it'd still work. @@ -341,6 +377,7 @@ impl WorkflowMachines { message_outbox: Default::default(), encountered_patch_markers: Default::default(), local_activity_data: LocalActivityData::default(), + default_stream_name: Default::default(), have_seen_terminal_event: false, worker_config: basics.worker_config, } @@ -364,11 +401,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(()) } @@ -835,19 +892,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())); + } } } @@ -1056,6 +1134,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. @@ -1811,6 +1974,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, self.default_stream_name.clone()), + 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::SubscribeNotificationChannel(attrs) => { // Never resolves: the event it produces records the subscription and hands // nothing back. The notifications arrive later on a scheduled event. @@ -2160,6 +2343,96 @@ fn decode_wake_signal( Some(wake) } +/// 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) +} + /// Adds a scheduled event's notifications to the ones already read for this task. /// /// Uses the server's own fold rule, one per channel with the highest counter kept, so a channel diff --git a/crates/sdk-core/src/worker/workflow/managed_run.rs b/crates/sdk-core/src/worker/workflow/managed_run.rs index 649a14780..44423346b 100644 --- a/crates/sdk-core/src/worker/workflow/managed_run.rs +++ b/crates/sdk-core/src/worker/workflow/managed_run.rs @@ -301,9 +301,11 @@ impl ManagedRun { if is_incremental { self.metrics.sticky_cache_hit(); } - self.wfm - .machines - .new_work_from_server(work.update, work.messages)?; + self.wfm.machines.new_work_from_server( + work.update, + work.messages, + work.stream_slices, + )?; } // A wake Signal reaches Core as a history event, so it can only be classified once that @@ -789,7 +791,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 diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index 0043be83a..bf367fca8 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -95,6 +95,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}, }, @@ -1169,6 +1170,7 @@ struct PreparedWFT { query_requests: Vec, update: HistoryUpdate, messages: Vec, + stream_slices: Vec, } impl PreparedWFT { @@ -1860,6 +1862,8 @@ enum WFCommandVariant { ExternalOutputStreamBuffered(WorkflowOutputStreamBuffered), /// The complete set of notification channels the run listens on. Never implies retention. ExternalStreamChannels(WorkflowStreamChannels), + SubscribeStream(SubscribeStream), + AppendStreamRecords(AppendStreamRecords), SubscribeNotificationChannel(SubscribeNotificationChannel), UnsubscribeNotificationChannel(UnsubscribeNotificationChannel), } @@ -1870,6 +1874,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) => { @@ -2028,6 +2036,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 @@ -2107,6 +2124,7 @@ impl WFMachinesError { match self { WFMachinesError::Nondeterminism(_) => EvictionReason::Nondeterminism, WFMachinesError::Fatal(_) => EvictionReason::Fatal, + WFMachinesError::MissingRecords(_) => EvictionReason::PaginationOrHistoryFetch, } } diff --git a/crates/sdk/src/workflow_future.rs b/crates/sdk/src/workflow_future.rs index cbb960762..e1226bcd1 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -336,6 +336,15 @@ impl WorkflowFuture { .context("Nexus operation must have result")?; push_polled_context!(ActivationJobContext::Passive); } + Variant::DeliverStreamRecords(slice) => { + // No stream API in this SDK. Bailing rather than ignoring: + // the server has recorded this range as consumed and will + // not send it again, so dropping it loses data silently. + bail!( + "received stream records for {}, which this SDK cannot deliver", + slice.stream_id + ); + } 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 d5d1a443e..ba8a89430 100644 --- a/crates/workflow/src/runtime/instance.rs +++ b/crates/workflow/src/runtime/instance.rs @@ -1118,6 +1118,20 @@ where self.apply_resolution(resolution); ActivationJobResult::None } + Some(ActivationVariant::DeliverStreamRecords(slice)) => { + // The Rust workflow runtime has no stream API yet. Failing + // is the only safe answer: the server has already recorded + // this range as consumed, so dropping it would leave the + // workflow permanently behind data it will never be sent + // again. + return Err(Box::new(Failure { + message: format!( + "received stream records for {}, which this SDK cannot deliver", + slice.stream_id + ), + ..Default::default() + })); + } Some(ActivationVariant::RemoveFromCache(_)) => ActivationJobResult::None, // This runtime has no way to subscribe to a channel. A notification carries no // data, only a hint to go read a source, so dropping one loses nothing.