From bd62d2b83eab1f5d052f84a62328fce4f260391d Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Sat, 3 Oct 2026 01:45:01 -0700 Subject: [PATCH 1/2] Vendored the stream api into Core's api tree. Core decodes the stream commands, events, slices and consumed ranges only when its vendored api carries them. The type maps in the protos crate name the two new commands and events, which are not ignorable. --- crates/common/build.rs | 5 ++ crates/protos/build.rs | 1 + .../temporal/api/command/v1/message.proto | 43 ++++++++++ .../temporal/api/enums/v1/command_type.proto | 2 + .../temporal/api/enums/v1/event_type.proto | 8 ++ .../temporal/api/enums/v1/failed_cause.proto | 8 ++ .../temporal/api/history/v1/message.proto | 44 +++++++++++ .../temporal/api/stream/v1/message.proto | 79 ++++++++++++++++++- .../workflowservice/v1/request_response.proto | 5 ++ crates/protos/src/protos/mod.rs | 19 +++++ 10 files changed, 211 insertions(+), 3 deletions(-) diff --git a/crates/common/build.rs b/crates/common/build.rs index 8647a610e..f8f1429f7 100644 --- a/crates/common/build.rs +++ b/crates/common/build.rs @@ -829,6 +829,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", // Notification metadata: the server bounds it with its own notification size limit, not // the blob limit, so the SDK has nothing to mirror. "temporal.api.notification.v1.Notification.metadata", diff --git a/crates/protos/build.rs b/crates/protos/build.rs index 2163d16fd..d829d8d76 100644 --- a/crates/protos/build.rs +++ b/crates/protos/build.rs @@ -41,6 +41,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/temporal/api/command/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto index 0af8ba569..aaec7e4b0 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"; @@ -328,6 +329,8 @@ message Command { subscribe_notification_channel_command_attributes = 20; UnsubscribeNotificationChannelCommandAttributes unsubscribe_notification_channel_command_attributes = 21; + AppendStreamRecordsCommandAttributes append_stream_records_command_attributes = 22; + SubscribeStreamCommandAttributes subscribe_stream_command_attributes = 23; } } @@ -347,3 +350,43 @@ message UnsubscribeNotificationChannelCommandAttributes { // The channel to stop listening on, as the writers name it. string channel = 1; } + +// 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 { + // Name of a stream this Workflow owns, scoped to the Workflow. Created on + // first use. Empty means the Workflow's default output stream. A Workflow + // cannot append to a stream in another execution, so this is never the id + // of a standalone stream. + string stream_name = 1; + // Stored in order. The server sets `producer_id` to empty on each record, + // because the owning Workflow is the producer here. + repeated temporal.api.stream.v1.StreamRecord records = 2; +} + +// 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, named either way round: a stream this Workflow owns + // by the name it appends under, a stream in another execution by its id. + // The server tries them in that order, so a Workflow that owns a stream + // under this name cannot reach a standalone stream with the same id. When + // neither exists the Workflow gets a stream of its own by that name, which + // is how a reader subscribes before the first record is written. + string stream_name_or_id = 1; + // Where to start, as an absolute offset. Read only when `start_position` + // is unset. A negative value is refused: the head of the stream is asked + // for with `start_position.tail`. + int64 start_offset = 2; + // Where to start. The server resolves it once, when it registers the + // subscription, and records the resolved absolute offset on the subscribed + // event, so replay does not resolve it again. Setting it together with a + // non-zero `start_offset` fails the command. + temporal.api.stream.v1.StreamStartPosition start_position = 3; +} 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 c3b3d8e9f..a46064e48 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 @@ -31,4 +31,6 @@ enum CommandType { COMMAND_TYPE_REQUEST_CANCEL_NEXUS_OPERATION = 18; COMMAND_TYPE_SUBSCRIBE_NOTIFICATION_CHANNEL = 19; COMMAND_TYPE_UNSUBSCRIBE_NOTIFICATION_CHANNEL = 20; + COMMAND_TYPE_APPEND_STREAM_RECORDS = 21; + COMMAND_TYPE_SUBSCRIBE_STREAM = 22; } 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 e7d7a31a2..6dd2f5a79 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 @@ -182,4 +182,12 @@ enum EventType { // Recorded for every UnsubscribeNotificationChannel command, including one // naming a channel the run was not subscribed to. EVENT_TYPE_WORKFLOW_NOTIFICATION_CHANNEL_UNSUBSCRIBED = 62; + // 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 = 63; + // 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 = 64; } 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 96baff0f8..2ccfc6900 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 @@ -95,6 +95,14 @@ enum WorkflowTaskFailedCause { WORKFLOW_TASK_FAILED_CAUSE_BAD_SUBSCRIBE_NOTIFICATION_CHANNEL_ATTRIBUTES = 41; // An UnsubscribeNotificationChannel command named an empty or too-long channel. WORKFLOW_TASK_FAILED_CAUSE_BAD_UNSUBSCRIBE_NOTIFICATION_CHANNEL_ATTRIBUTES = 42; + // A workflow task completed with an invalid AppendStreamRecords command. + WORKFLOW_TASK_FAILED_CAUSE_BAD_APPEND_STREAM_RECORDS_ATTRIBUTES = 43; + // A workflow task completed with an invalid SubscribeStream command. + WORKFLOW_TASK_FAILED_CAUSE_BAD_SUBSCRIBE_STREAM_ATTRIBUTES = 44; + // 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 = 45; } // Activity tasks can fail for various reasons. Note that some of these reasons can only originate 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 40f2b0a11..85b071a23 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/notification/v1/message.proto"; import "temporal/api/deployment/v1/message.proto"; import "temporal/api/failure/v1/message.proto"; @@ -384,6 +385,12 @@ 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. + repeated temporal.api.stream.v1.StreamRange consumed_stream_ranges = 14; } message WorkflowTaskTimedOutEventAttributes { @@ -980,6 +987,41 @@ message WorkflowNotificationChannelUnsubscribedEventAttributes { int64 subscribed_event_id = 3; } +message WorkflowStreamSubscribedEventAttributes { + // The WorkflowTaskCompleted event of the task whose command created this + // subscription. + int64 workflow_task_completed_event_id = 1; + // The stream the Workflow subscribed to, as the command addressed it: + // either the name of a stream this Workflow owns or the id of one in + // another execution. + string stream_id = 2; + // The offset the subscription actually starts from. Resolved by the server + // when the subscription is registered and recorded here, so replay reads + // 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; + // Name of the stream the Workflow appended to. + string stream_id = 2; + // Inclusive. Same range vocabulary as StreamRange and StreamSlice, so a + // reader does not have to remember which of the three counts and which + // bounds. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 from_offset = 3; + // Exclusive. With from_offset this names the range without carrying any of + // it, which is what keeps this event a fixed size no matter how large the + // batch or its payloads are. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 to_offset = 4; +} + message WorkflowExecutionUpdateAcceptedEventAttributes { // The instance ID of the update protocol that generated this event. string protocol_instance_id = 1; @@ -1307,6 +1349,8 @@ message HistoryEvent { workflow_notification_channel_subscribed_event_attributes = 66; WorkflowNotificationChannelUnsubscribedEventAttributes workflow_notification_channel_unsubscribed_event_attributes = 67; + WorkflowStreamSubscribedEventAttributes workflow_stream_subscribed_event_attributes = 68; + WorkflowStreamRecordsAppendedEventAttributes workflow_stream_records_appended_event_attributes = 69; } } diff --git a/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto index dd1905803..e6a8fab2d 100644 --- a/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto @@ -14,9 +14,13 @@ import "temporal/api/common/v1/message.proto"; // One entry in a stream. The record is the wire format: stores keep it // serialized as is and readers in every language decode the same bytes. message StreamRecord { - // The value the producer published, stored as sent. A stream provider - // writes and reads it as part of the record, so applying a payload codec - // to it is the provider's job. + // The value the producer published, stored as sent. + // + // A payload codec applies on the paths this API owns: the append command on + // RespondWorkflowTaskCompleted, and the slices on PollWorkflowTaskQueue. + // Records a producer writes or reads through the stream service take a + // different path, whose messages are not part of this API yet and so are + // outside what a codec-applying proxy walks. temporal.api.common.v1.Payload body = 1; // Producer-supplied provenance, stored as sent. map metadata = 2; @@ -35,6 +39,75 @@ message StreamRecord { 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 { + // The stream, as the subscribing command addressed it: either the name of + // a stream the consuming Workflow owns or the id of one in another + // execution. + string stream_id = 1; + // Run id of the execution that owns the stream. Set on both a slice for the + // task being started and a re-supplied one. + string run_id = 2; + // Inclusive. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 from_offset = 3; + // Exclusive. Equal to from_offset when the subscription observed nothing, + // which is a fact replay has to reproduce rather than an absence of one. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 to_offset = 4; + repeated StreamRecord records = 5; + // The WorkflowTaskCompleted event whose consumed_stream_ranges recorded + // 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 { + // The stream, as the subscribing command addressed it: either the name of + // a stream the consuming Workflow owns or the id of one in another + // execution. + string stream_id = 1; + // Inclusive. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 from_offset = 2; + // Exclusive. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "from" and "to" name a half-open offset range. --) + int64 to_offset = 3; +} + +// Where a new subscription or read begins. The server resolves it against the +// stream as it stands in the same transaction that registers the reader, so +// the result does not race with appends or truncation, and records the +// resolved absolute offset. +message StreamStartPosition { + oneof position { + // Absolute and inclusive. Refused when below the stream's floor. + int64 offset = 1; + // The last N records the stream holds, or all of them when it holds + // fewer. Counts records of every kind. Must be positive. + int64 last_n = 2; + // The oldest record the stream still holds. Must be true. + bool earliest = 3; + // Only records appended after registration: the stream's head offset. + // Must be true. + bool tail = 4; + } +} + // What a record means to a reader. Kept on the record itself so every store // and every language reads it the same way without a private envelope. enum StreamRecordKind { diff --git a/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto b/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto index 47f6ebe24..48485ac79 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/notification/v1/message.proto"; import "temporal/api/history/v1/message.proto"; import "temporal/api/workflow/v1/message.proto"; @@ -385,6 +386,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/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index c6709adab..133c7b81b 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1840,6 +1840,12 @@ pub mod temporal { CommandType::ScheduleActivityTask } Attributes::StartTimerCommandAttributes(_) => CommandType::StartTimer, + Attributes::SubscribeStreamCommandAttributes(_) => { + CommandType::SubscribeStream + } + Attributes::AppendStreamRecordsCommandAttributes(_) => { + CommandType::AppendStreamRecords + } Attributes::SubscribeNotificationChannelCommandAttributes(_) => { CommandType::SubscribeNotificationChannel } @@ -2378,6 +2384,8 @@ pub mod temporal { | EventType::TimerStarted | EventType::UpsertWorkflowSearchAttributes | EventType::WorkflowPropertiesModified + | EventType::WorkflowStreamSubscribed + | EventType::WorkflowStreamRecordsAppended | EventType::WorkflowNotificationChannelSubscribed | EventType::WorkflowNotificationChannelUnsubscribed | EventType::NexusOperationScheduled @@ -2479,6 +2487,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::WorkflowNotificationChannelSubscribedEventAttributes(_) => { false } @@ -2572,6 +2584,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::WorkflowNotificationChannelSubscribedEventAttributes(_) => { EventType::WorkflowNotificationChannelSubscribed } Attributes::WorkflowNotificationChannelUnsubscribedEventAttributes(_) => { EventType::WorkflowNotificationChannelUnsubscribed } Attributes::WorkflowExecutionStartedEventAttributes(_) => { EventType::WorkflowExecutionStarted } @@ -2696,6 +2710,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"); From e917870a4423e6506147e0b056fc149d9dd9b624 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Sat, 3 Oct 2026 01:45:01 -0700 Subject: [PATCH 2/2] Added the DeliverStreamRecords activation job. A consumed stream range reaches lang as its own job. The Rust SDKs have no stream API yet, so they fail the activation rather than drop records the server will not send again. --- .../workflow_activation.proto | 24 ++++++++++++++++++- crates/protos/src/protos/mod.rs | 7 ++++++ crates/sdk/src/workflow_future.rs | 9 +++++++ crates/workflow/src/runtime/instance.rs | 14 +++++++++++ 4 files changed, 53 insertions(+), 1 deletion(-) 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 118c3a0c5..3d3998f6d 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/notification/v1/message.proto"; import "temporal/api/enums/v1/failed_cause.proto"; import "temporal/api/enums/v1/workflow.proto"; @@ -144,11 +145,13 @@ message WorkflowActivationJob { // 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 + // alongside this one and share this message. The numbers below are fixed // with that family and must not be reused. // // Notifications from the channels the workflow subscribed to. NotificationsReceived notifications_received = 21; + // A range of a stream the workflow subscribed to. + DeliverStreamRecords deliver_stream_records = 22; // Remove the workflow identified by the [WorkflowActivation] containing this job from the // cache after performing the activation. It is guaranteed that this will be the only job // in the activation if present. @@ -165,6 +168,25 @@ message NotificationsReceived { repeated temporal.api.notification.v1.Notification notifications = 1; } +// 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/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index 133c7b81b..2d1895330 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 + ) + } workflow_activation_job::Variant::NotificationsReceived(n) => { write!(f, "NotificationsReceived({})", n.notifications.len()) } diff --git a/crates/sdk/src/workflow_future.rs b/crates/sdk/src/workflow_future.rs index a62e2a0db..68085a5a5 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -336,6 +336,15 @@ impl WorkflowFuture { .context("Nexus operation must have result")?; push_polled_context!(ActivationJobContext::Passive); } + Variant::DeliverStreamRecords(slice) => { + // No stream API in this SDK. Bailing rather than ignoring: + // the server has recorded this range as consumed and will + // not send it again, so dropping it loses data silently. + bail!( + "received stream records for {}, which this SDK cannot deliver", + slice.stream_id + ); + } Variant::NotificationsReceived(_) => { // No channel API in this SDK, so nothing here subscribed, and a // notification carries no data a workflow could lose by this. diff --git a/crates/workflow/src/runtime/instance.rs b/crates/workflow/src/runtime/instance.rs index 2a06cf73c..4dfb4159e 100644 --- a/crates/workflow/src/runtime/instance.rs +++ b/crates/workflow/src/runtime/instance.rs @@ -1118,6 +1118,20 @@ where self.apply_resolution(resolution); ActivationJobResult::None } + Some(ActivationVariant::DeliverStreamRecords(slice)) => { + // The Rust workflow runtime has no stream API yet. Failing + // is the only safe answer: the server has already recorded + // this range as consumed, so dropping it would leave the + // workflow permanently behind data it will never be sent + // again. + return Err(Box::new(Failure { + message: format!( + "received stream records for {}, which this SDK cannot deliver", + slice.stream_id + ), + ..Default::default() + })); + } Some(ActivationVariant::NotificationsReceived(_)) => { // This runtime cannot subscribe to a channel, and a notification // carries no data the workflow could lose by ignoring it.