diff --git a/.github/workflows/per-pr.yml b/.github/workflows/per-pr.yml index 11325303a..89126c4cc 100644 --- a/.github/workflows/per-pr.yml +++ b/.github/workflows/per-pr.yml @@ -123,7 +123,7 @@ jobs: integ-tests: name: Integ tests - timeout-minutes: ${{ github.ref == 'refs/heads/main' && 30 || 25 }} + timeout-minutes: ${{ matrix.timeoutMinutes || (github.ref == 'refs/heads/main' && 30 || 25) }} strategy: fail-fast: false matrix: diff --git a/crates/common/build.rs b/crates/common/build.rs index 2e4a378de..1baae9de8 100644 --- a/crates/common/build.rs +++ b/crates/common/build.rs @@ -820,6 +820,11 @@ const NOT_VALIDATED_FIELDS: &[&str] = &[ "temporal.api.workflowservice.v1.StartWorkflowExecutionRequest.continued_failure", "temporal.api.workflowservice.v1.StartWorkflowExecutionRequest.last_completion_result", "temporal.api.workflowservice.v1.TerminateWorkflowExecutionRequest.details", + // Stream records: the blob limit is a per-event limit, and these bodies never reach an + // event. The server bounds the batch by message count (MaxMessagesPerBatch) instead, which + // is not a payload size the SDK can mirror. + "temporal.api.stream.v1.StreamRecord.body", + "temporal.api.stream.v1.StreamRecord.metadata", // Dedicated, non-fetchable limits (not blob/memo, not in DescribeNamespace): UserMetadata // (nexus-start only); Nexus EndpointSpec.description (maxDescriptionSize; cloud variant cloud-only). "temporal.api.sdk.v1.UserMetadata.details", diff --git a/crates/protos/build.rs b/crates/protos/build.rs index 509ba9259..7365df524 100644 --- a/crates/protos/build.rs +++ b/crates/protos/build.rs @@ -36,6 +36,7 @@ const SERDE_DERIVE_PREFIXES: &[&str] = &[ ".temporal.api.rules", ".temporal.api.schedule", ".temporal.api.sdk", + ".temporal.api.stream", ".temporal.api.taskqueue", ".temporal.api.testservice", ".temporal.api.update", diff --git a/crates/protos/protos/api_upstream/nexus/deps/nexus-temporal-types/model.wit b/crates/protos/protos/api_upstream/nexus/deps/nexus-temporal-types/model.wit index 91de3e838..b6e6e5e8f 100644 --- a/crates/protos/protos/api_upstream/nexus/deps/nexus-temporal-types/model.wit +++ b/crates/protos/protos/api_upstream/nexus/deps/nexus-temporal-types/model.wit @@ -147,12 +147,7 @@ interface model { /// typescript="common.WorkflowIdReusePolicy" /// dotnet="Temporalio.Api.Enums.V1.WorkflowIdReusePolicy" /// typescript-import="@temporalio/common" - enum workflow-id-reuse-policy { - allow-duplicate, - allow-duplicate-failed-only, - reject-duplicate, - terminate-if-running, - } + type workflow-id-reuse-policy = placeholder; /// @nexus.proto "temporal.api.enums.v1.WorkflowIdConflictPolicy" typescript-import="@temporalio/proto" /// @nexus.type @@ -160,11 +155,7 @@ interface model { /// typescript="common.WorkflowIdConflictPolicy" /// dotnet="Temporalio.Api.Enums.V1.WorkflowIdConflictPolicy" /// typescript-import="@temporalio/common" - enum workflow-id-conflict-policy { - fail, - use-existing, - terminate-existing, - } + type workflow-id-conflict-policy = placeholder; /// @nexus.proto "temporal.api.sdk.v1.UserMetadata" typescript-import="@temporalio/proto" /// @nexus.flatten-in-api 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 ee839115b..7d10e441b 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 @@ -14,6 +14,7 @@ import "google/protobuf/duration.proto"; import "temporal/api/enums/v1/workflow.proto"; import "temporal/api/enums/v1/command_type.proto"; import "temporal/api/common/v1/message.proto"; +import "temporal/api/stream/v1/message.proto"; import "temporal/api/failure/v1/message.proto"; import "temporal/api/taskqueue/v1/message.proto"; import "temporal/api/workflow/v1/message.proto"; @@ -324,5 +325,38 @@ message Command { ScheduleNexusOperationCommandAttributes schedule_nexus_operation_command_attributes = 18; RequestCancelNexusOperationCommandAttributes request_cancel_nexus_operation_command_attributes = 19; + AppendStreamRecordsCommandAttributes append_stream_records_command_attributes = 20; + SubscribeStreamCommandAttributes subscribe_stream_command_attributes = 21; } } + +// Appends records to a stream the Workflow owns. Applied inside the Workflow +// Task's own commit. Produces one `WorkflowStreamRecordsAppended` event +// 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; + // 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; +} + +// 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 rather than supplied here. +// 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. + int64 start_offset = 2; +} diff --git a/crates/protos/protos/api_upstream/temporal/api/enums/v1/command_type.proto b/crates/protos/protos/api_upstream/temporal/api/enums/v1/command_type.proto index 067d95391..967169b19 100644 --- a/crates/protos/protos/api_upstream/temporal/api/enums/v1/command_type.proto +++ b/crates/protos/protos/api_upstream/temporal/api/enums/v1/command_type.proto @@ -29,4 +29,6 @@ enum CommandType { COMMAND_TYPE_MODIFY_WORKFLOW_PROPERTIES = 16; COMMAND_TYPE_SCHEDULE_NEXUS_OPERATION = 17; COMMAND_TYPE_REQUEST_CANCEL_NEXUS_OPERATION = 18; + COMMAND_TYPE_APPEND_STREAM_RECORDS = 19; + COMMAND_TYPE_SUBSCRIBE_STREAM = 20; } diff --git a/crates/protos/protos/api_upstream/temporal/api/enums/v1/event_type.proto b/crates/protos/protos/api_upstream/temporal/api/enums/v1/event_type.proto index b879f51e8..a386b1db7 100644 --- a/crates/protos/protos/api_upstream/temporal/api/enums/v1/event_type.proto +++ b/crates/protos/protos/api_upstream/temporal/api/enums/v1/event_type.proto @@ -175,4 +175,12 @@ enum EventType { EVENT_TYPE_WORKFLOW_EXECUTION_UNPAUSED = 59; // An event that indicates time skipping advanced time or was disabled automatically after a bound was reached. EVENT_TYPE_WORKFLOW_EXECUTION_TIME_SKIPPING_TRANSITIONED = 60; + // A Workflow subscribed to a stream. Recorded once per subscription, not + // per record: the offsets a task consumed ride WorkflowTaskCompleted and + // the payloads never enter History at all. + EVENT_TYPE_WORKFLOW_STREAM_SUBSCRIBED = 61; + // A Workflow appended a batch of records to a stream. Recorded per + // batch, and carrying only the offset range it landed at: the bodies go to + // the stream's own log, never into History. + EVENT_TYPE_WORKFLOW_STREAM_RECORDS_APPENDED = 62; } diff --git a/crates/protos/protos/api_upstream/temporal/api/enums/v1/failed_cause.proto b/crates/protos/protos/api_upstream/temporal/api/enums/v1/failed_cause.proto index 81cbde73e..274e533d5 100644 --- a/crates/protos/protos/api_upstream/temporal/api/enums/v1/failed_cause.proto +++ b/crates/protos/protos/api_upstream/temporal/api/enums/v1/failed_cause.proto @@ -90,6 +90,14 @@ enum WorkflowTaskFailedCause { WORKFLOW_TASK_FAILED_CAUSE_WORKFLOW_PAUSE_REQUESTED_BEFORE_TASK_STARTED = 39; // A workflow task failed because the request exceeded a size limit. WORKFLOW_TASK_FAILED_CAUSE_REQUEST_TOO_LARGE = 40; + // A workflow task completed with an invalid AppendStreamRecords command. + WORKFLOW_TASK_FAILED_CAUSE_BAD_APPEND_STREAM_RECORDS_ATTRIBUTES = 41; + // A workflow task completed with an invalid SubscribeStream command. + WORKFLOW_TASK_FAILED_CAUSE_BAD_SUBSCRIBE_STREAM_ATTRIBUTES = 42; + // A workflow task could not be started because a stream range it consumed and recorded in + // History can no longer be served, for example after truncation or because it exceeds the + // replay bound. Check the workflow task failure message for more information. + WORKFLOW_TASK_FAILED_CAUSE_STREAM_RANGE_UNAVAILABLE = 43; } enum StartChildWorkflowExecutionFailedCause { 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 0211c6f55..8deca8d15 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 @@ -17,6 +17,7 @@ import "temporal/api/enums/v1/failed_cause.proto"; import "temporal/api/enums/v1/update.proto"; import "temporal/api/enums/v1/workflow.proto"; import "temporal/api/common/v1/message.proto"; +import "temporal/api/stream/v1/message.proto"; import "temporal/api/deployment/v1/message.proto"; import "temporal/api/failure/v1/message.proto"; import "temporal/api/taskqueue/v1/message.proto"; @@ -379,6 +380,13 @@ message WorkflowTaskCompletedEventAttributes { // The Worker Deployment Version that completed this task. Must be set if `versioning_behavior` // is set. This value updates workflow execution's `versioning_info.deployment_version`. temporal.api.deployment.v1.WorkerDeploymentVersion deployment_version = 11; + + // Offset ranges this Workflow Task consumed from streams it subscribes to. + // 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; } message WorkflowTaskTimedOutEventAttributes { @@ -953,6 +961,33 @@ message ActivityPropertiesModifiedExternallyEventAttributes { temporal.api.common.v1.RetryPolicy new_retry_policy = 2; } +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. + 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 + // the resolved value rather than resolving it again against a stream that + // has since moved. + int64 start_offset = 3; +} + +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. + 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; +} + message WorkflowExecutionUpdateAcceptedEventAttributes { // The instance ID of the update protocol that generated this event. string protocol_instance_id = 1; @@ -1276,6 +1311,8 @@ message HistoryEvent { WorkflowExecutionPausedEventAttributes workflow_execution_paused_event_attributes = 63; WorkflowExecutionUnpausedEventAttributes workflow_execution_unpaused_event_attributes = 64; WorkflowExecutionTimeSkippingTransitionedEventAttributes workflow_execution_time_skipping_transitioned_event_attributes = 65; + WorkflowStreamSubscribedEventAttributes workflow_stream_subscribed_event_attributes = 66; + WorkflowStreamRecordsAppendedEventAttributes workflow_stream_records_appended_event_attributes = 67; } } 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 new file mode 100644 index 000000000..9f0c88dfb --- /dev/null +++ b/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto @@ -0,0 +1,82 @@ +syntax = "proto3"; + +package temporal.api.stream.v1; + +option go_package = "go.temporal.io/api/stream/v1;stream"; +option java_package = "io.temporal.api.stream.v1"; +option java_multiple_files = true; +option java_outer_classname = "MessageProto"; +option ruby_package = "Temporalio::Api::Stream::V1"; +option csharp_namespace = "Temporalio.Api.Stream.V1"; + +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. + temporal.api.common.v1.Payload body = 1; + // Producer-supplied provenance, stored as sent. + map metadata = 2; + // Producer-supplied grouping label, stored as sent. + string topic = 3; + // How to read this record. Unspecified is read as DATA. + StreamRecordKind kind = 4; + // Who wrote the record. Empty when the owning Workflow did. + string producer_id = 5; + // 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. + int64 sequence = 7; +} + +// A contiguous range of a stream delivered to a Workflow Task, along with the +// offsets it covers. The offsets are what History records; the records +// themselves are never written to History. +message StreamSlice { + 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. + 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. + int64 to_offset = 4; + repeated StreamRecord records = 5; + // The WorkflowTaskCompleted event whose consumed_stream_ranges recorded + // this range. Set only when the server is re-supplying a range for a task + // being replayed; a slice for the task now being started leaves it unset, + // because the event closing that task does not exist yet. + // + // Replay needs this because a Workflow Task response carries one slice set + // while a cache miss replays every prior task, so the ranges have to be + // matched to the events that recorded them rather than to the response. + int64 workflow_task_completed_event_id = 6; +} + +// The offsets a Workflow Task consumed, without the payloads. Recorded on +// WorkflowTaskCompleted so History grows with Workflow Tasks rather than with +// records. +message StreamRange { + string stream_id = 1; + // Inclusive. + int64 from_offset = 2; + // Exclusive. + int64 to_offset = 3; +} + +// 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 { + // Read as DATA. + STREAM_RECORD_KIND_UNSPECIFIED = 0; + // A value the producer published; `body` carries it. + STREAM_RECORD_KIND_DATA = 1; + // The producer named by `producer_id` writes nothing more on `topic`. + // Says nothing about that producer's outcome and does not end the stream. + STREAM_RECORD_KIND_FINISH = 2; +} diff --git a/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto b/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto index c3dd95769..b396de5b4 100644 --- a/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto +++ b/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto @@ -24,6 +24,7 @@ import "temporal/api/enums/v1/activity.proto"; import "temporal/api/enums/v1/nexus.proto"; import "temporal/api/activity/v1/message.proto"; import "temporal/api/common/v1/message.proto"; +import "temporal/api/stream/v1/message.proto"; import "temporal/api/history/v1/message.proto"; import "temporal/api/workflow/v1/message.proto"; import "temporal/api/command/v1/message.proto"; @@ -383,6 +384,10 @@ message PollWorkflowTaskQueueResponse { // 3. If every group has some pending polls, assign the next poll to a group randomly // according to the weights. temporal.api.taskqueue.v1.PollerGroupsInfo poller_groups_info = 19; + + // Stream data attached to this task. Delivered out of band so the payloads + // never enter History; only the offset ranges are recorded there. + repeated temporal.api.stream.v1.StreamSlice stream_slices = 20; } message RespondWorkflowTaskCompletedRequest { 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 097994ac6..8106089e7 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/workflow.proto"; import "temporal/sdk/core/activity_result/activity_result.proto"; import "temporal/sdk/core/child_workflow/child_workflow.proto"; @@ -139,6 +140,12 @@ message WorkflowActivationJob { ResolveNexusOperationStart resolve_nexus_operation_start = 15; // A nexus operation resolved. ResolveNexusOperation resolve_nexus_operation = 16; + // 17 to 20 are taken by the external stream jobs, which are developed + // alongside this one and share this message. The number below is fixed + // with that family and must not be reused. + // + // A range of a stream the workflow subscribed to. + DeliverStreamRecords deliver_stream_records = 21; // Remove the workflow identified by the [WorkflowActivation] containing this job from the // cache after performing the activation. It is guaranteed that this will be the only job // in the activation if present. @@ -146,6 +153,25 @@ message WorkflowActivationJob { } } +// 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 03e172216..456055d6c 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_commands/workflow_commands.proto @@ -15,6 +15,7 @@ import "temporal/api/common/v1/message.proto"; import "temporal/api/enums/v1/workflow.proto"; import "temporal/api/failure/v1/message.proto"; import "temporal/api/sdk/v1/user_metadata.proto"; +import "temporal/api/stream/v1/message.proto"; import "temporal/api/sdk/v1/event_group_marker.proto"; import "temporal/sdk/core/child_workflow/child_workflow.proto"; import "temporal/sdk/core/nexus/nexus.proto"; @@ -54,9 +55,40 @@ message WorkflowCommand { UpdateResponse update_response = 20; ScheduleNexusOperation schedule_nexus_operation = 21; RequestCancelNexusOperation request_cancel_nexus_operation = 22; + // 23 to 28 are taken by the external stream commands, which are developed + // alongside these and share this message. The two numbers below are fixed + // with that family and must not be reused. + SubscribeStream subscribe_stream = 29; + AppendStreamRecords append_stream_records = 30; } } +// Append a batch of records to a stream this workflow owns. +// +// The bodies go to the stream's own log rather than into History, which gets +// one fixed-size event naming the offset range. That is what makes the batch +// size free: a thousand records cost the same in History as one. The server +// stores each record with an empty producer id, because the workflow is the +// producer here. +message AppendStreamRecords { + // Empty means the workflow's default output stream. + string stream_id = 1; + repeated temporal.api.stream.v1.StreamRecord records = 2; +} + +// Subscribe this workflow to a stream, so later Workflow Tasks carry the +// ranges it has not consumed yet. +// +// The stream's addressing is resolved by the server. A workflow cannot look it +// up without doing I/O, and a value it carried would be a reading rather than a +// fact, so it could differ on replay. +message SubscribeStream { + string stream_id = 1; + // Negative means from wherever the stream is when the subscription is + // registered. The server resolves that once and records it. + int64 start_offset = 2; +} + message StartTimer { // Lang's incremental sequence number, used as the operation identifier uint32 seq = 1; diff --git a/crates/protos/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index 4bbfa5ac5..3cf466a9a 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1277,6 +1277,13 @@ pub mod coresdk { workflow_activation_job::Variant::ResolveNexusOperation(_) => { write!(f, "ResolveNexusOperation") } + workflow_activation_job::Variant::DeliverStreamRecords(d) => { + write!( + f, + "DeliverStreamRecords({}, {}..{})", + d.stream_id, d.from_offset, d.to_offset + ) + } } } } @@ -1486,6 +1493,23 @@ pub mod coresdk { } } + impl Display for SubscribeStream { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "SubscribeStream({})", self.stream_id) + } + } + + impl Display for AppendStreamRecords { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!( + f, + "AppendStreamRecords({}, {} records)", + self.stream_id, + self.records.len() + ) + } + } + impl Display for StartTimer { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "StartTimer({})", self.seq) @@ -1825,6 +1849,12 @@ pub mod temporal { CommandType::ScheduleActivityTask } Attributes::StartTimerCommandAttributes(_) => CommandType::StartTimer, + Attributes::SubscribeStreamCommandAttributes(_) => { + CommandType::SubscribeStream + } + Attributes::AppendStreamRecordsCommandAttributes(_) => { + CommandType::AppendStreamRecords + } Attributes::CompleteWorkflowExecutionCommandAttributes(_) => { CommandType::CompleteWorkflowExecution } @@ -1878,6 +1908,28 @@ pub mod temporal { } } + impl From for command::Attributes { + fn from(s: workflow_commands::AppendStreamRecords) -> Self { + Self::AppendStreamRecordsCommandAttributes( + AppendStreamRecordsCommandAttributes { + stream_id: s.stream_id, + records: s.records, + }, + ) + } + } + + impl From for command::Attributes { + fn from(s: workflow_commands::SubscribeStream) -> Self { + Self::SubscribeStreamCommandAttributes( + SubscribeStreamCommandAttributes { + stream_id: s.stream_id, + start_offset: s.start_offset, + }, + ) + } + } + impl From for command::Attributes { fn from(s: workflow_commands::StartTimer) -> Self { Self::StartTimerCommandAttributes(StartTimerCommandAttributes { @@ -2337,6 +2389,8 @@ pub mod temporal { | EventType::TimerStarted | EventType::UpsertWorkflowSearchAttributes | EventType::WorkflowPropertiesModified + | EventType::WorkflowStreamSubscribed + | EventType::WorkflowStreamRecordsAppended | EventType::NexusOperationScheduled | EventType::NexusOperationCancelRequested | EventType::WorkflowExecutionCanceled @@ -2436,6 +2490,10 @@ pub mod temporal { // mark any new event types as ignorable or not. if let Some(a) = self.attributes.as_ref() { match a { + Attributes::WorkflowStreamSubscribedEventAttributes(_) => false, + Attributes::WorkflowStreamRecordsAppendedEventAttributes(_) => { + false + } Attributes::WorkflowExecutionStartedEventAttributes(_) => false, Attributes::WorkflowExecutionCompletedEventAttributes(_) => false, Attributes::WorkflowExecutionFailedEventAttributes(_) => false, @@ -2523,6 +2581,8 @@ pub mod temporal { pub fn event_type(&self) -> EventType { // I just absolutely _love_ this match self { + Attributes::WorkflowStreamSubscribedEventAttributes(_) => { EventType::WorkflowStreamSubscribed } + Attributes::WorkflowStreamRecordsAppendedEventAttributes(_) => { EventType::WorkflowStreamRecordsAppended } Attributes::WorkflowExecutionStartedEventAttributes(_) => { EventType::WorkflowExecutionStarted } Attributes::WorkflowExecutionCompletedEventAttributes(_) => { EventType::WorkflowExecutionCompleted } Attributes::WorkflowExecutionFailedEventAttributes(_) => { EventType::WorkflowExecutionFailed } @@ -2640,6 +2700,11 @@ pub mod temporal { tonic::include_proto!("temporal.api.sdk.v1"); } } + pub mod stream { + pub mod v1 { + tonic::include_proto!("temporal.api.stream.v1"); + } + } pub mod taskqueue { pub mod v1 { tonic::include_proto!("temporal.api.taskqueue.v1"); diff --git a/crates/sdk-core/CHANGELOG.md b/crates/sdk-core/CHANGELOG.md index bfb8e04f9..4a50406c0 100644 --- a/crates/sdk-core/CHANGELOG.md +++ b/crates/sdk-core/CHANGELOG.md @@ -50,6 +50,17 @@ 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. +* 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. ### Breaking Changes :boom: * The following types are now non-exhaustive: `Priority`, `WorkerDeploymentVersion`, diff --git a/crates/sdk-core/src/core_tests/mod.rs b/crates/sdk-core/src/core_tests/mod.rs index 6dd4479bb..4779016bb 100644 --- a/crates/sdk-core/src/core_tests/mod.rs +++ b/crates/sdk-core/src/core_tests/mod.rs @@ -2,6 +2,7 @@ mod activity_tasks; mod event_groups; mod queries; mod replay_flag; +mod streams; mod updates; mod workers; mod workflow_cancels; diff --git a/crates/sdk-core/src/core_tests/streams.rs b/crates/sdk-core/src/core_tests/streams.rs new file mode 100644 index 000000000..4e1a55a44 --- /dev/null +++ b/crates/sdk-core/src/core_tests/streams.rs @@ -0,0 +1,1343 @@ +//! 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}, + 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) +} + +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_id: "s1".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); +} + +/// The command has to reach the server carrying the bodies. +/// +/// Only a task that is not being replayed sends commands, so this drives a +/// single open Workflow Task rather than a recorded history. +#[tokio::test] +async fn publish_command_reaches_the_server_with_its_payloads() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .times(1) + .returning(|resp, _| { + let cmd = resp.commands.first().expect("a command was sent"); + assert_eq!(cmd.command_type(), CommandType::AppendStreamRecords); + match cmd.attributes.as_ref().unwrap() { + command::Attributes::AppendStreamRecordsCommandAttributes(a) => { + assert_eq!(a.stream_id, "s1"); + // The bodies are the half of the batch History never sees, + // so the command is the only thing that can carry them. + let bodies: Vec<_> = a + .records + .iter() + .map(|m| m.body.as_ref().unwrap().data.clone()) + .collect(); + assert_eq!(bodies, vec![b"one".to_vec(), b"two".to_vec()]); + } + other => panic!("wrong attributes: {other:?}"), + } + Ok(RespondWorkflowTaskCompletedResponse::default()) + }); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish_two("s1").into()], + )) + .await + .unwrap(); + core.shutdown().await; +} + +/// A publish reissued on replay has to match the event the original run wrote. +/// +/// This is the property the event exists for. Core pops one queued command per +/// command-generated event, so a publish producing none would leave every later +/// command matched against the wrong event. Core signals the mismatch by +/// failing the workflow task, which the mock turns into a panic. +#[tokio::test] +async fn publish_command_round_trips_through_replay() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_records_appended("s1", 0, 2); + t.add_full_wf_task(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .returning(|_, _| Ok(RespondWorkflowTaskCompletedResponse::default())); + mock_client + .expect_fail_workflow_task() + .returning(|_, _, f| panic!("core rejected the reissued publish: {f:?}")); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + // First activation replays the recorded publish: lang reissues it, and core + // has to match it to that event rather than calling it nondeterministic. + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![publish_two("s1").into()], + )) + .await + .unwrap(); + + core.shutdown().await; +} + +/// The subscribe command has to reach the server, which only a task that is +/// not being replayed will send. +#[tokio::test] +async fn subscribe_command_reaches_the_server() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + + let mut mock_client = mock_worker_client(); + mock_client + .expect_complete_workflow_task() + .times(1) + .returning(|resp, _| { + let cmd = resp.commands.first().expect("a command was sent"); + assert_eq!(cmd.command_type(), CommandType::SubscribeStream); + match cmd.attributes.as_ref().unwrap() { + command::Attributes::SubscribeStreamCommandAttributes(a) => { + assert_eq!(a.stream_id, "s1"); + // Passed through unresolved: the server turns it into a + // real offset and records that. + assert_eq!(a.start_offset, -1); + } + other => panic!("wrong attributes: {other:?}"), + } + Ok(RespondWorkflowTaskCompletedResponse::default()) + }); + + let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client); + let core = mock_worker(build_mock_pollers(mock)); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s1".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); + core.shutdown().await; +} + +fn publish_two(stream_id: &str) -> AppendStreamRecords { + AppendStreamRecords { + stream_id: stream_id.to_string(), + records: vec![ + StreamRecord { + body: Some(b"one".to_vec().into()), + ..Default::default() + }, + StreamRecord { + body: Some(b"two".to_vec().into()), + ..Default::default() + }, + ], + } +} + +/// A 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_id: "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_id: "in".to_string(), + start_offset: 0, + } + .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_id: "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_id: "in".to_string(), + start_offset: 0, + } + .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_id: "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 how many records +/// landed, so a replay that publishes fewer has diverged from the original run. +#[tokio::test] +async fn a_publish_reissued_with_a_different_batch_size_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_records_appended("s1", 0, 2); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + AppendStreamRecords { + stream_id: "s1".to_string(), + records: vec![StreamRecord { + body: Some(b"one".to_vec().into()), + ..Default::default() + }], + } + .into(), + ], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// A subscription is checked on the stream alone. The recorded start offset is +/// the server's resolution of what the command asked for, so it is not the +/// command's to reproduce. +#[tokio::test] +async fn a_subscribe_reissued_to_a_different_stream_fails_the_task() { + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_full_wf_task(); + t.add_stream_subscribed("s1", 4); + t.add_full_wf_task(); + + let (core, failures) = worker_expecting_one_nondeterminism_failure(t, ResponseType::AllHistory); + + let task = core.poll_workflow_activation().await.unwrap(); + core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds( + task.run_id, + vec![ + SubscribeStream { + stream_id: "s2".to_string(), + start_offset: -1, + } + .into(), + ], + )) + .await + .unwrap(); + core.handle_eviction().await; + assert_eq!(failures.load(Ordering::Relaxed), 1); + core.shutdown().await; +} + +/// 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. +#[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_nondeterminism_failure(t, ResponseType::Raw(poll_resp.resp)); + + // 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 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_id: "in".to_string(), + start_offset: 0, + } + .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_id: &str, bodies: &[&str]) -> AppendStreamRecords { + AppendStreamRecords { + stream_id: stream_id.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, 1); + 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 as +/// nondeterministic rather than replay it on different input. +#[tokio::test] +async fn a_pushed_history_with_a_wrong_slice_fails_as_nondeterministic() { + 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 + // nondeterministic (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_nondeterministic() { + 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 9f0b77348..9af592890 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 418600403..43723d243 100644 --- a/crates/sdk-core/src/replay/history_builder.rs +++ b/crates/sdk-core/src/replay/history_builder.rs @@ -23,6 +23,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, @@ -129,6 +130,48 @@ 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. + pub fn add_stream_records_appended( + &mut self, + stream_id: &str, + first_offset: i64, + record_count: i64, + ) -> i64 { + let attrs = WorkflowStreamRecordsAppendedEventAttributes { + workflow_task_completed_event_id: self.previous_task_completed_id, + stream_id: stream_id.to_string(), + first_offset, + record_count, + }; + self.add(attrs) + } + /// Add a workflow task timed out event. pub fn add_workflow_task_timed_out(&mut self) { let attrs = WorkflowTaskTimedOutEventAttributes { diff --git a/crates/sdk-core/src/replay/mod.rs b/crates/sdk-core/src/replay/mod.rs index 6e4cadc50..34bae9043 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,23 @@ 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 and fails the task as + /// nondeterministic when a slice disagrees with the recorded range or a recorded range with + /// content has no slice. + 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 856589946..4100b6e2e 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,28 @@ 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. + 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 +661,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, }) } @@ -738,6 +763,17 @@ 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() + { + saw_command = true; + saw_command_or_started = true; + } if let Some(next_next_event) = events.get(ix + 2) { if !saw_command && next_next_event.event_type() == EventType::WorkflowTaskScheduled diff --git a/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs new file mode 100644 index 000000000..5acdc7eda --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/append_stream_records_state_machine.rs @@ -0,0 +1,144 @@ +use super::{ + NewMachineWithCommand, StateMachine, TransitionResult, fsm, workflow_machines::MachineResponse, +}; +use crate::worker::workflow::{ + WFMachinesError, fatal, + machines::{EventInfo, HistEventData, WFMachinesAdapter}, + nondeterminism, +}; +use temporalio_common::protos::{ + coresdk::workflow_commands::AppendStreamRecords, + temporal::api::{ + enums::v1::{CommandType, EventType}, + history::v1::{WorkflowStreamRecordsAppendedEventAttributes, history_event}, + }, +}; + +fsm! { + pub(super) name AppendStreamRecordsMachine; + command AppendStreamRecordsMachineCommand; + error WFMachinesError; + shared_state SharedState; + + Created --(CommandScheduled) --> CommandIssued; + CommandIssued --(CommandRecorded(WorkflowStreamRecordsAppendedEventAttributes), + shared on_command_recorded) --> Done; +} + +/// What the command claimed, kept so the recorded event can be held against it. +#[derive(Default, Clone)] +pub(super) struct SharedState { + stream_id: String, + record_count: i64, +} + +/// Append a batch of records to a stream this workflow owns. +/// +/// The bodies go to the stream's own log, and History gets one event naming the +/// offset range the batch landed at. The offsets are assigned by the server, so +/// nothing here predicts them. +pub(super) fn append_stream_records(lang_cmd: AppendStreamRecords) -> NewMachineWithCommand { + let sm = AppendStreamRecordsMachine::from_parts( + Created {}.into(), + SharedState { + stream_id: lang_cmd.stream_id.clone(), + record_count: lang_cmd.records.len() as i64, + }, + ); + NewMachineWithCommand { + command: lang_cmd.into(), + machine: sm.into(), + } +} + +#[derive(Debug, derive_more::Display)] +pub(super) enum AppendStreamRecordsMachineCommand {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Created {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct CommandIssued {} + +impl CommandIssued { + pub(super) fn on_command_recorded( + self, + dat: &mut SharedState, + attrs: WorkflowStreamRecordsAppendedEventAttributes, + ) -> AppendStreamRecordsMachineTransition { + // An empty id names the workflow's default stream, and the server is the + // one that resolves that name, so only a named stream can be compared. + let same_stream = dat.stream_id.is_empty() || dat.stream_id == attrs.stream_id; + if same_stream && dat.record_count == attrs.record_count { + TransitionResult::default() + } else { + TransitionResult::Err(nondeterminism!( + "Recorded append of {} records to stream {:?} does not match the reissued \ + append of {} records to stream {:?}", + attrs.record_count, + attrs.stream_id, + dat.record_count, + dat.stream_id + )) + } + } +} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Done {} + +impl WFMachinesAdapter for AppendStreamRecordsMachine { + fn adapt_response( + &self, + _my_command: Self::Command, + _event_info: Option, + ) -> Result, Self::Error> { + Err(Self::Error::Nondeterminism( + "AppendStreamRecords does not use state machine commands".to_string(), + )) + } +} + +impl TryFrom for AppendStreamRecordsMachineEvents { + type Error = WFMachinesError; + + fn try_from(e: HistEventData) -> Result { + let e = e.event; + match e.event_type() { + EventType::WorkflowStreamRecordsAppended => { + if let Some( + history_event::Attributes::WorkflowStreamRecordsAppendedEventAttributes(attrs), + ) = e.attributes + { + Ok(AppendStreamRecordsMachineEvents::CommandRecorded(attrs)) + } else { + Err(fatal!("Stream records appended attributes were unset: {e}")) + } + } + _ => Err(Self::Error::Nondeterminism(format!( + "AppendStreamRecordsMachine does not handle {e}" + ))), + } + } +} + +impl TryFrom for AppendStreamRecordsMachineEvents { + type Error = WFMachinesError; + + fn try_from(c: CommandType) -> Result { + match c { + CommandType::AppendStreamRecords => { + Ok(AppendStreamRecordsMachineEvents::CommandScheduled) + } + _ => Err(Self::Error::Nondeterminism(format!( + "AppendStreamRecordsMachine does not handle command type {c:?}" + ))), + } + } +} + +impl From for CommandIssued { + fn from(_: Created) -> Self { + Self {} + } +} diff --git a/crates/sdk-core/src/worker/workflow/machines/mod.rs b/crates/sdk-core/src/worker/workflow/machines/mod.rs index bc5e3fcb9..364897465 100644 --- a/crates/sdk-core/src/worker/workflow/machines/mod.rs +++ b/crates/sdk-core/src/worker/workflow/machines/mod.rs @@ -1,3 +1,4 @@ +mod append_stream_records_state_machine; mod workflow_machines; mod activity_state_machine; @@ -15,6 +16,7 @@ mod modify_workflow_properties_state_machine; mod nexus_operation_state_machine; mod patch_state_machine; mod signal_external_state_machine; +mod subscribe_stream_state_machine; mod timer_state_machine; mod update_state_machine; mod upsert_search_attributes_state_machine; @@ -31,6 +33,7 @@ use crate::{ worker::workflow::{WFMachinesError, fatal, nondeterminism}, }; use activity_state_machine::ActivityMachine; +use append_stream_records_state_machine::AppendStreamRecordsMachine; use cancel_external_state_machine::CancelExternalMachine; use cancel_workflow_state_machine::CancelWorkflowMachine; use child_workflow_state_machine::ChildWorkflowMachine; @@ -46,6 +49,7 @@ use std::{ convert::{TryFrom, TryInto}, fmt::{Debug, Display}, }; +use subscribe_stream_state_machine::SubscribeStreamMachine; use temporalio_common::{ fsm_trait::{StateMachine, TransitionResult}, protos::temporal::api::{ @@ -80,6 +84,8 @@ enum Machines { WorkflowTaskMachine, UpsertSearchAttributesMachine, ModifyWorkflowPropertiesMachine, + SubscribeStreamMachine, + AppendStreamRecordsMachine, UpdateMachine, NexusOperationMachine, } diff --git a/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs new file mode 100644 index 000000000..3fc08e768 --- /dev/null +++ b/crates/sdk-core/src/worker/workflow/machines/subscribe_stream_state_machine.rs @@ -0,0 +1,136 @@ +use super::{ + NewMachineWithCommand, StateMachine, TransitionResult, fsm, workflow_machines::MachineResponse, +}; +use crate::worker::workflow::{ + WFMachinesError, fatal, + machines::{EventInfo, HistEventData, WFMachinesAdapter}, + nondeterminism, +}; +use temporalio_common::protos::{ + coresdk::workflow_commands::SubscribeStream, + temporal::api::{ + enums::v1::{CommandType, EventType}, + history::v1::{WorkflowStreamSubscribedEventAttributes, history_event}, + }, +}; + +fsm! { + pub(super) name SubscribeStreamMachine; + command SubscribeStreamMachineCommand; + error WFMachinesError; + shared_state SharedState; + + Created --(CommandScheduled) --> CommandIssued; + CommandIssued --(CommandRecorded(WorkflowStreamSubscribedEventAttributes), + shared on_command_recorded) --> Done; +} + +/// The stream the command named, kept so the recorded event can be held against it. +/// The start offset is not kept: the server resolves it, so the recorded value is +/// its answer rather than what the command said. +#[derive(Default, Clone)] +pub(super) struct SharedState { + stream_id: String, +} + +/// Subscribe this workflow to a stream. The command carries only the stream id +/// and a start offset; the server resolves the addressing, because a workflow +/// cannot look it up without doing I/O and a value it carried would be a +/// reading rather than a fact. +pub(super) fn subscribe_stream(lang_cmd: SubscribeStream) -> NewMachineWithCommand { + let sm = SubscribeStreamMachine::from_parts( + Created {}.into(), + SharedState { + stream_id: lang_cmd.stream_id.clone(), + }, + ); + NewMachineWithCommand { + command: lang_cmd.into(), + machine: sm.into(), + } +} + +#[derive(Debug, derive_more::Display)] +pub(super) enum SubscribeStreamMachineCommand {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Created {} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct CommandIssued {} + +impl CommandIssued { + pub(super) fn on_command_recorded( + self, + dat: &mut SharedState, + attrs: WorkflowStreamSubscribedEventAttributes, + ) -> SubscribeStreamMachineTransition { + if dat.stream_id == attrs.stream_id { + TransitionResult::default() + } else { + TransitionResult::Err(nondeterminism!( + "Recorded subscription to stream {:?} does not match the reissued subscription \ + to stream {:?}", + attrs.stream_id, + dat.stream_id + )) + } + } +} + +#[derive(Debug, Default, Clone, derive_more::Display)] +pub(super) struct Done {} + +impl WFMachinesAdapter for SubscribeStreamMachine { + fn adapt_response( + &self, + _my_command: Self::Command, + _event_info: Option, + ) -> Result, Self::Error> { + Err(Self::Error::Nondeterminism( + "SubscribeStream does not use state machine commands".to_string(), + )) + } +} + +impl TryFrom for SubscribeStreamMachineEvents { + type Error = WFMachinesError; + + fn try_from(e: HistEventData) -> Result { + let e = e.event; + match e.event_type() { + EventType::WorkflowStreamSubscribed => { + if let Some(history_event::Attributes::WorkflowStreamSubscribedEventAttributes( + attrs, + )) = e.attributes + { + Ok(SubscribeStreamMachineEvents::CommandRecorded(attrs)) + } else { + Err(fatal!("Stream subscribed attributes were unset: {e}")) + } + } + _ => Err(Self::Error::Nondeterminism(format!( + "SubscribeStreamMachine does not handle {e}" + ))), + } + } +} + +impl TryFrom for SubscribeStreamMachineEvents { + type Error = WFMachinesError; + + fn try_from(c: CommandType) -> Result { + match c { + CommandType::SubscribeStream => Ok(SubscribeStreamMachineEvents::CommandScheduled), + _ => Err(Self::Error::Nondeterminism(format!( + "SubscribeStreamMachine does not handle command type {c:?}" + ))), + } + } +} + +impl From for CommandIssued { + fn from(_: Created) -> Self { + Self {} + } +} diff --git a/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs index b266137fa..41c054450 100644 --- a/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs +++ b/crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs @@ -66,6 +66,7 @@ mod machine_coverage_report { use super::*; use crate::worker::workflow::machines::{ StateMachine, activity_state_machine::ActivityMachine, + append_stream_records_state_machine::AppendStreamRecordsMachine, cancel_external_state_machine::CancelExternalMachine, cancel_workflow_state_machine::CancelWorkflowMachine, child_workflow_state_machine::ChildWorkflowMachine, @@ -75,7 +76,8 @@ mod machine_coverage_report { local_activity_state_machine::LocalActivityMachine, modify_workflow_properties_state_machine::ModifyWorkflowPropertiesMachine, nexus_operation_state_machine::NexusOperationMachine, patch_state_machine::PatchMachine, - signal_external_state_machine::SignalExternalMachine, timer_state_machine::TimerMachine, + signal_external_state_machine::SignalExternalMachine, + subscribe_stream_state_machine::SubscribeStreamMachine, timer_state_machine::TimerMachine, update_state_machine::UpdateMachine, upsert_search_attributes_state_machine::UpsertSearchAttributesMachine, workflow_task_state_machine::WorkflowTaskMachine, @@ -117,6 +119,8 @@ mod machine_coverage_report { let mut modify_wf_props = ModifyWorkflowPropertiesMachine::visualizer().to_owned(); let mut update = UpdateMachine::visualizer().to_owned(); let mut nexus = NexusOperationMachine::visualizer().to_owned(); + let mut subscribe_stream = SubscribeStreamMachine::visualizer().to_owned(); + let mut append_stream_records = AppendStreamRecordsMachine::visualizer().to_owned(); // This isn't at all efficient but doesn't need to be. // Replace transitions in the vizzes with green color if they are covered. @@ -144,6 +148,12 @@ mod machine_coverage_report { } m @ "UpdateMachine" => cover_transitions(m, &mut update, coverage), m @ "NexusOperationMachine" => cover_transitions(m, &mut nexus, coverage), + m @ "SubscribeStreamMachine" => { + cover_transitions(m, &mut subscribe_stream, coverage) + } + m @ "AppendStreamRecordsMachine" => { + cover_transitions(m, &mut append_stream_records, coverage) + } m => panic!("Unknown machine {m}"), } } diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index 8d4b8f54e..c742e3b8f 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -2,13 +2,15 @@ mod local_acts; use super::{ Machines, NewMachineWithCommand, TemporalStateMachine, + append_stream_records_state_machine::append_stream_records, cancel_external_state_machine::new_external_cancel, cancel_workflow_state_machine::cancel_workflow, complete_workflow_state_machine::complete_workflow, continue_as_new_workflow_state_machine::continue_as_new, fail_workflow_state_machine::fail_workflow, local_activity_state_machine::new_local_activity, patch_state_machine::has_change, signal_external_state_machine::new_external_signal, - timer_state_machine::new_timer, upsert_search_attributes_state_machine::upsert_search_attrs, + subscribe_stream_state_machine::subscribe_stream, timer_state_machine::new_timer, + upsert_search_attributes_state_machine::upsert_search_attrs, workflow_machines::local_acts::LocalActivityData, workflow_task_state_machine::WorkflowTaskMachine, }; @@ -70,6 +72,7 @@ use temporalio_common::{ history::v1::{HistoryEvent, history_event}, protocol::v1::{Message as ProtocolMessage, message::SequencingId}, sdk::v1::WorkflowTaskCompletedMetadata, + stream::v1::{StreamRange, StreamSlice}, }, }, worker::WorkerDeploymentVersion, @@ -86,6 +89,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, /// EventId of the last handled WorkflowTaskStarted event @@ -266,7 +292,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 { @@ -276,6 +302,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. @@ -325,11 +355,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(()) } @@ -617,19 +667,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())); + } } } @@ -824,6 +895,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. @@ -1526,6 +1682,26 @@ impl WorkflowMachines { CommandIdKind::NeverResolves, ); } + WFCommandVariant::AppendStreamRecords(attrs) => { + // Never resolves: the event names the offset range the + // server assigned and hands nothing back. A workflow that + // wants to know where its batch landed reads the stream. + self.add_cmd_to_wf_task( + append_stream_records(attrs), + annotations, + CommandIdKind::NeverResolves, + ); + } + WFCommandVariant::SubscribeStream(attrs) => { + // Never resolves: the event it produces records the + // subscription and hands nothing back to the workflow. The + // ranges arrive later as their own activation jobs. + self.add_cmd_to_wf_task( + subscribe_stream(attrs), + annotations, + CommandIdKind::NeverResolves, + ); + } WFCommandVariant::UpdateResponse(ur) => { let m_key = self.get_machine_by_msg(&ur.protocol_instance_id)?; let m = if let Machines::UpdateMachine(m) = self.machine_mut(m_key) { @@ -1817,3 +1993,92 @@ enum CommandIdKind { /// A command which is fire-and-forget (ex: Upsert search attribs) NeverResolves, } + +/// 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. +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(nondeterminism!( + "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(nondeterminism!( + "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(nondeterminism!( + "The server sent stream {} from offset {} to {} for event {event_id}, which records \ + no such range", + extra.stream_id, + extra.from_offset, + extra.to_offset + )); + } + Ok(jobs) +} diff --git a/crates/sdk-core/src/worker/workflow/managed_run.rs b/crates/sdk-core/src/worker/workflow/managed_run.rs index 603b38acc..9d6d23e39 100644 --- a/crates/sdk-core/src/worker/workflow/managed_run.rs +++ b/crates/sdk-core/src/worker/workflow/managed_run.rs @@ -42,6 +42,7 @@ use temporalio_common::protos::{ temporal::api::{ enums::v1::{VersioningBehavior, WorkflowTaskFailedCause}, failure::v1::Failure, + stream::v1::StreamSlice, }, }; use tokio::sync::oneshot; @@ -247,7 +248,8 @@ impl ManagedRun { if is_incremental { self.metrics.sticky_cache_hit(); } - self.wfm.new_work_from_server(work.update, work.messages)? + self.wfm + .new_work_from_server(work.update, work.messages, work.stream_slices)? } else { let r = self.wfm.get_next_activation()?; if r.jobs.is_empty() { @@ -623,7 +625,10 @@ impl ManagedRun { EvictionReason::Unspecified | EvictionReason::PaginationOrHistoryFetch ); - let (should_report, 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 (should_report, rur) = if is_no_report_query_fail && !self.am_broken { (false, None) } else { // Blow up any cached data associated with the workflow @@ -633,11 +638,14 @@ impl ManagedRun { reason, auto_reply_fail_tt: None, }); - let should_report = match &evict_req_outcome { - EvictionRequestResult::EvictionRequested(Some(attempt), _) - | EvictionRequestResult::EvictionAlreadyRequested(Some(attempt)) => *attempt <= 1, - _ => false, - }; + let should_report = !is_no_report_query_fail + && match &evict_req_outcome { + EvictionRequestResult::EvictionRequested(Some(attempt), _) + | EvictionRequestResult::EvictionAlreadyRequested(Some(attempt)) => { + *attempt <= 1 + } + _ => false, + }; let rur = evict_req_outcome.into_run_update_resp(); (should_report, rur) }; @@ -1448,8 +1456,10 @@ impl WorkflowManager { &mut self, update: HistoryUpdate, messages: Vec, + stream_slices: Vec, ) -> Result { - self.machines.new_work_from_server(update, messages)?; + self.machines + .new_work_from_server(update, messages, stream_slices)?; self.get_next_activation() } diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index b0da392a3..aa2cef0af 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -87,6 +87,7 @@ use temporalio_common::{ protocol::v1::Message as ProtocolMessage, query::v1::WorkflowQuery, sdk::v1::{EventGroupMarker, UserMetadata, WorkflowTaskCompletedMetadata}, + stream::v1::StreamSlice, taskqueue::v1::StickyExecutionAttributes, workflowservice::v1::{PollActivityTaskQueueResponse, get_system_info_response}, }, @@ -1017,6 +1018,7 @@ struct PreparedWFT { query_requests: Vec, update: HistoryUpdate, messages: Vec, + stream_slices: Vec, } impl PreparedWFT { @@ -1527,6 +1529,8 @@ enum WFCommandVariant { UpdateResponse(UpdateResponse), ScheduleNexusOperation(ScheduleNexusOperation), RequestCancelNexusOperation(RequestCancelNexusOperation), + SubscribeStream(SubscribeStream), + AppendStreamRecords(AppendStreamRecords), } impl TryFrom for WFCommand { @@ -1535,6 +1539,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) => { @@ -1663,6 +1671,13 @@ 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 brought none of them, so this worker cannot replay the run. Not the workflow's + /// fault: the records only travel with the task, and a worker handed a sticky task for a run + /// it no longer holds has no way to fetch them. 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 @@ -1742,6 +1757,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 d6920870d..6e7f421f0 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -340,6 +340,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 c3417ef31..8af164913 100644 --- a/crates/workflow/src/runtime/instance.rs +++ b/crates/workflow/src/runtime/instance.rs @@ -1093,6 +1093,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, None => { return Err(Box::new(Failure {