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/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 c6709adab..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()) } @@ -1840,6 +1847,12 @@ pub mod temporal { CommandType::ScheduleActivityTask } Attributes::StartTimerCommandAttributes(_) => CommandType::StartTimer, + Attributes::SubscribeStreamCommandAttributes(_) => { + CommandType::SubscribeStream + } + Attributes::AppendStreamRecordsCommandAttributes(_) => { + CommandType::AppendStreamRecords + } Attributes::SubscribeNotificationChannelCommandAttributes(_) => { CommandType::SubscribeNotificationChannel } @@ -2378,6 +2391,8 @@ pub mod temporal { | EventType::TimerStarted | EventType::UpsertWorkflowSearchAttributes | EventType::WorkflowPropertiesModified + | EventType::WorkflowStreamSubscribed + | EventType::WorkflowStreamRecordsAppended | EventType::WorkflowNotificationChannelSubscribed | EventType::WorkflowNotificationChannelUnsubscribed | EventType::NexusOperationScheduled @@ -2479,6 +2494,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 +2591,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 +2717,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/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.