diff --git a/.github/workflows/per-pr.yml b/.github/workflows/per-pr.yml index 9c3bfa785..c088bb13a 100644 --- a/.github/workflows/per-pr.yml +++ b/.github/workflows/per-pr.yml @@ -141,7 +141,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/client/src/grpc.rs b/crates/client/src/grpc.rs index 1c3d5b0a6..9aa1c8306 100644 --- a/crates/client/src/grpc.rs +++ b/crates/client/src/grpc.rs @@ -955,6 +955,51 @@ proxier! { r.extensions_mut().insert(labels); } ); + ( + notify_channel, + NotifyChannelRequest, + NotifyChannelResponse, + |r| { + let labels = namespaced_request!(r); + r.extensions_mut().insert(labels); + } + ); + ( + register_channel_listener, + RegisterChannelListenerRequest, + RegisterChannelListenerResponse, + |r| { + let labels = namespaced_request!(r); + r.extensions_mut().insert(labels); + } + ); + ( + unregister_channel_listener, + UnregisterChannelListenerRequest, + UnregisterChannelListenerResponse, + |r| { + let labels = namespaced_request!(r); + r.extensions_mut().insert(labels); + } + ); + ( + poll_channel, + PollChannelRequest, + PollChannelResponse, + |r| { + let labels = namespaced_request!(r); + r.extensions_mut().insert(labels); + } + ); + ( + describe_channel, + DescribeChannelRequest, + DescribeChannelResponse, + |r| { + let labels = namespaced_request!(r); + r.extensions_mut().insert(labels); + } + ); ( signal_with_start_workflow_execution, SignalWithStartWorkflowExecutionRequest, diff --git a/crates/common/build.rs b/crates/common/build.rs index b3170c97e..f8f1429f7 100644 --- a/crates/common/build.rs +++ b/crates/common/build.rs @@ -829,6 +829,14 @@ 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", // 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 346d506c6..d829d8d76 100644 --- a/crates/protos/build.rs +++ b/crates/protos/build.rs @@ -33,6 +33,7 @@ const SERDE_DERIVE_PREFIXES: &[&str] = &[ ".temporal.api.nexus", ".temporal.api.nexusoperation", ".temporal.api.nexusservices", + ".temporal.api.notification", ".temporal.api.operatorservice", ".temporal.api.protocol", ".temporal.api.query", @@ -40,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 ee839115b..aadc8f89d 100644 --- a/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/command/v1/message.proto @@ -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,68 @@ 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; + SubscribeNotificationChannelCommandAttributes + subscribe_notification_channel_command_attributes = 22; + UnsubscribeNotificationChannelCommandAttributes + unsubscribe_notification_channel_command_attributes = 23; } } + +// 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; +} + +// Makes the Workflow a listener of a notification channel for this run. The +// next notifications on the channel arrive on the scheduled event of a Workflow +// Task. The subscription ends with the run, and a successor subscribes again. +message SubscribeNotificationChannelCommandAttributes { + // The channel to listen on, as the writers name it. + string channel = 1; +} + +// Ends the run's subscription to a notification channel. Notifications already +// recorded on a scheduled event still reach that Workflow Task; later ones do +// not. A command naming a channel the run is not subscribed to records its +// event and changes nothing, so replay matches every command to an event. +message UnsubscribeNotificationChannelCommandAttributes { + // The channel to stop listening on, as the writers name it. + string channel = 1; +} 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..629a26528 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,8 @@ 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; + COMMAND_TYPE_SUBSCRIBE_NOTIFICATION_CHANNEL = 21; + COMMAND_TYPE_UNSUBSCRIBE_NOTIFICATION_CHANNEL = 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 b879f51e8..815e1aebe 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,19 @@ 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; + // A Workflow became a listener of a notification channel for its run. + // The notifications themselves ride the WorkflowTaskScheduled event. + EVENT_TYPE_WORKFLOW_NOTIFICATION_CHANNEL_SUBSCRIBED = 63; + // A Workflow stopped listening on a notification channel for its run. + // Recorded for every UnsubscribeNotificationChannel command, including one + // naming a channel the run was not subscribed to. + EVENT_TYPE_WORKFLOW_NOTIFICATION_CHANNEL_UNSUBSCRIBED = 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 f902aa9a1..9ff4c9a53 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,19 @@ 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; + // A SubscribeNotificationChannel command named an empty or too-long channel, or hit a + // subscription or listener limit. + WORKFLOW_TASK_FAILED_CAUSE_BAD_SUBSCRIBE_NOTIFICATION_CHANNEL_ATTRIBUTES = 44; + // An UnsubscribeNotificationChannel command named an empty or too-long channel. + WORKFLOW_TASK_FAILED_CAUSE_BAD_UNSUBSCRIBE_NOTIFICATION_CHANNEL_ATTRIBUTES = 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 b40324d68..553683100 100644 --- a/crates/protos/protos/api_upstream/temporal/api/history/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/history/v1/message.proto @@ -17,6 +17,8 @@ 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"; import "temporal/api/taskqueue/v1/message.proto"; @@ -302,6 +304,10 @@ message WorkflowTaskScheduledEventAttributes { google.protobuf.Duration start_to_close_timeout = 2; // Starting at 1, how many attempts there have been to complete this task int32 attempt = 3; + // Notifications for channels this Workflow listens to, folded per channel + // since the last task was scheduled. In History so a Workflow may act on + // them deterministically and replay sees the same. + repeated temporal.api.notification.v1.Notification notifications = 4; } message WorkflowTaskStartedEventAttributes { @@ -379,6 +385,16 @@ 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 = 20; + + // Held for fields added on the main line, so a rebase does not land one of + // them on a number this fork already writes. + reserved 14 to 19; } message WorkflowTaskTimedOutEventAttributes { @@ -955,6 +971,61 @@ 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; + // 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 WorkflowNotificationChannelSubscribedEventAttributes { + // The WorkflowTaskCompleted event of the task whose command created this + // subscription. + int64 workflow_task_completed_event_id = 1; + // The channel the Workflow listens on for the rest of this run. + string channel = 2; +} + +message WorkflowNotificationChannelUnsubscribedEventAttributes { + // The WorkflowTaskCompleted event of the task whose command ended this + // subscription. + int64 workflow_task_completed_event_id = 1; + // The channel the Workflow stopped listening on. + string channel = 2; + // The WorkflowNotificationChannelSubscribed event that recorded the + // subscription this command ended. Zero when the run held no subscription + // for the channel. + int64 subscribed_event_id = 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; @@ -1278,6 +1349,12 @@ 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; + WorkflowNotificationChannelSubscribedEventAttributes + workflow_notification_channel_subscribed_event_attributes = 68; + WorkflowNotificationChannelUnsubscribedEventAttributes + workflow_notification_channel_unsubscribed_event_attributes = 69; } } diff --git a/crates/protos/protos/api_upstream/temporal/api/notification/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/notification/v1/message.proto new file mode 100644 index 000000000..584619582 --- /dev/null +++ b/crates/protos/protos/api_upstream/temporal/api/notification/v1/message.proto @@ -0,0 +1,75 @@ +syntax = "proto3"; + +package temporal.api.notification.v1; + +option go_package = "go.temporal.io/api/notification/v1;notification"; +option java_package = "io.temporal.api.notification.v1"; +option java_multiple_files = true; +option java_outer_classname = "MessageProto"; +option ruby_package = "Temporalio::Api::Notification::V1"; +option csharp_namespace = "Temporalio.Api.Notification.V1"; + +import "google/protobuf/timestamp.proto"; + +import "temporal/api/common/v1/message.proto"; + +// A notification tells the listeners of a channel that a source they consume +// has moved. It is not data: the listener reads the source itself. A channel +// is named by the writer and its listeners; for a stream, the provider formats +// the stream's identity into the name. Writers never learn who listens. The +// server folds notifications per listener while one is pending and no task has +// been scheduled for it, keeping the one with the highest counter. +message Notification { + // The channel the writer notified. Listeners register on the same name. + string channel = 1; + // Where the source stands after the write that caused this notification, + // in the writer's terms. Opaque to the server. + bytes position = 2; + // Orders notifications from one channel's writers. The writer derives it + // from the position, since only the source can order its positions. Among + // notifications folded together, the one with the highest counter is kept. + int64 counter = 3; + // Details for the listener, such as which topic moved. Bounded in size and + // carried as payloads, so a codec applies as to any payload. This is state, + // not a log: a fold keeps the latest notification only, so a writer puts + // here what is true at `position`, such as which topic moved or a close + // flag, never something a consumer must see once per write. + map metadata = 4; + // Set for a channel linked to an execution: the owner and the run that + // received the notification. Empty for an independent channel. A listener + // that holds both kinds routes the notification by it. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "to" names the owner the channel is linked to. --) + temporal.api.common.v1.Execution linked_to = 5; +} + +// A listener of a channel: a Workflow Execution woken with a Workflow Task, or +// a callback the server invokes with each notification. +message ChannelListener { + // Assigned by the server when the listener registers. + string listener_id = 1; + oneof listener { + WorkflowListener workflow = 2; + temporal.api.common.v1.Callback callback = 3; + } + google.protobuf.Timestamp registered_time = 4; +} + +// A Workflow Execution listening on a channel. +message WorkflowListener { + string workflow_id = 1; + // The run that subscribed. The server follows a continue-as-new to the + // chain's current run when it delivers. + string run_id = 2; +} + +// Where a channel lives, which decides how a call addresses it. +enum ChannelKind { + CHANNEL_KIND_UNSPECIFIED = 0; + // Its own execution, keyed by namespace and channel name. Any number of + // workflows and callbacks listen to it. + CHANNEL_KIND_INDEPENDENT = 1; + // Kept in one execution's state, keyed by namespace, execution and + // channel name. The owning execution is its listener by construction. + CHANNEL_KIND_LINKED = 2; +} 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..e6a8fab2d --- /dev/null +++ b/crates/protos/protos/api_upstream/temporal/api/stream/v1/message.proto @@ -0,0 +1,121 @@ +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 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; + // 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, zero when it does not number + // its records. Stored as sent; the server does not assign, validate or + // order by it, and the stream's own offsets are what order a read. + int64 sequence = 7; +} + +// 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 { + // 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/workflow/v1/message.proto b/crates/protos/protos/api_upstream/temporal/api/workflow/v1/message.proto index cf763aa12..1cb3e5eba 100644 --- a/crates/protos/protos/api_upstream/temporal/api/workflow/v1/message.proto +++ b/crates/protos/protos/api_upstream/temporal/api/workflow/v1/message.proto @@ -21,6 +21,7 @@ import "temporal/api/enums/v1/workflow.proto"; import "temporal/api/common/v1/message.proto"; import "temporal/api/deployment/v1/message.proto"; import "temporal/api/failure/v1/message.proto"; +import "temporal/api/notification/v1/message.proto"; import "temporal/api/taskqueue/v1/message.proto"; import "temporal/api/sdk/v1/user_metadata.proto"; @@ -586,6 +587,32 @@ message NexusOperationCancellationInfo { string blocked_reason = 7; } +// A workflow's standing on a notification channel, as reported by DescribeWorkflowExecution. +message ChannelSubscriptionInfo { + // Channel name. + string channel = 1; + // CHANNEL_KIND_INDEPENDENT for a channel the workflow subscribed to with a + // SubscribeNotificationChannel command. CHANNEL_KIND_LINKED for a channel linked to this + // workflow, which lists it once the channel holds any state. + temporal.api.notification.v1.ChannelKind kind = 2; + // Independent kind: id of the WorkflowNotificationChannelSubscribed event that recorded the + // subscription. Zero for the linked kind. + int64 subscribed_event_id = 3; + // Highest counter the workflow has accepted from the channel. Zero when none has arrived. + int64 last_counter = 4; + // The notification held for the workflow's next Workflow Task, when one is pending. + temporal.api.notification.v1.Notification pending_notification = 5; + // Counter carried by the scheduled event of a Workflow Task that has not started yet. Zero + // otherwise. + int64 scheduled_counter = 6; + // Linked kind: callback listeners registered on the channel. + int32 listener_count = 7; + // Linked kind: notifications retained for pollers. + int32 retained_count = 8; + // Linked kind: notifications the channel has accepted over its life. + int64 accepted_count = 9; +} + message WorkflowExecutionOptions { // If set, takes precedence over the Versioning Behavior sent by the SDK on Workflow Task completion. VersioningOverride versioning_override = 1; 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 1aae988d8..d69732a71 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,8 @@ 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"; import "temporal/api/command/v1/message.proto"; @@ -384,6 +386,13 @@ 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; + + // Used once by a repeated field this fork has since removed. + reserved 21; } message RespondWorkflowTaskCompletedRequest { @@ -880,6 +889,112 @@ message SignalWorkflowExecutionResponse { temporal.api.common.v1.Link link = 1; } +message NotifyChannelRequest { + string namespace = 1; + temporal.api.notification.v1.Notification notification = 2; + // The identity of the caller, for audit, metrics and logs. It is not copied + // into the notification. A writer that wants the consumer to see who wrote + // puts that in the notification's `metadata`. + string identity = 3; + // Used to de-dupe a retried notification. + string request_id = 4; + // When set, the call addresses the channel linked to this execution. + // `run_id` is optional and resolves to the current run of a workflow chain, + // as a Signal does. When unset, the call addresses the independent channel + // of that name. + temporal.api.common.v1.Execution execution = 5; +} + +message NotifyChannelResponse { + // Listeners registered when the notification was accepted. Zero means the + // notification was retained for pollers and woke nobody. + int32 listener_count = 1; +} + +message RegisterChannelListenerRequest { + string namespace = 1; + string channel = 2; + // Invoked with each notification on the channel. + temporal.api.common.v1.Callback callback = 3; + // Used to de-dupe a retried registration. + string request_id = 4; + // The identity of the caller, for metrics and logs. + string identity = 5; + // When set, the call addresses the channel linked to this execution. + // `run_id` is optional and resolves to the current run of a workflow chain, + // as a Signal does. When unset, the call addresses the independent channel + // of that name. + temporal.api.common.v1.Execution execution = 6; +} + +message RegisterChannelListenerResponse { + string listener_id = 1; +} + +message UnregisterChannelListenerRequest { + string namespace = 1; + string channel = 2; + string listener_id = 3; + // The identity of the caller, for metrics and logs. + string identity = 4; + // When set, the call addresses the channel linked to this execution. + // `run_id` is optional and resolves to the current run of a workflow chain, + // as a Signal does. When unset, the call addresses the independent channel + // of that name. + temporal.api.common.v1.Execution execution = 5; +} + +message UnregisterChannelListenerResponse { +} + +message PollChannelRequest { + string namespace = 1; + string channel = 2; + // Only notifications with a counter above this one are returned. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "after" names the exclusive lower bound. --) + int64 after_counter = 3; + // How long to wait for a notification when none is retained above + // `after_counter`. + google.protobuf.Duration wait = 4; + // At most this many notifications are returned. Zero means the server's + // default. + int32 max_notifications = 5; + // When set, the call addresses the channel linked to this execution. + // `run_id` is optional and resolves to the current run of a workflow chain, + // as a Signal does. When unset, the call addresses the independent channel + // of that name. + temporal.api.common.v1.Execution execution = 6; +} + +message PollChannelResponse { + repeated temporal.api.notification.v1.Notification notifications = 1; +} + +message DescribeChannelRequest { + string namespace = 1; + string channel = 2; + // When set, the call addresses the channel linked to this execution. + // `run_id` is optional and resolves to the current run of a workflow chain, + // as a Signal does. When unset, the call addresses the independent channel + // of that name. + temporal.api.common.v1.Execution execution = 3; +} + +message DescribeChannelResponse { + repeated temporal.api.notification.v1.ChannelListener listeners = 1; + // The notification with the highest counter the channel retains. + temporal.api.notification.v1.Notification latest = 2; + // How many notifications the channel retains for pollers. + int32 retained_count = 3; + temporal.api.notification.v1.ChannelKind kind = 4; + // The owner of a linked channel and the run that holds it. Empty for an + // independent channel. + // (-- api-linter: core::0140::prepositions=disabled + // aip.dev/not-precedent: "to" names the owner the channel is linked to. --) + temporal.api.common.v1.Execution linked_to = 5; +} + message SignalWithStartWorkflowExecutionRequest { string namespace = 1; string workflow_id = 2; @@ -1212,6 +1327,9 @@ message DescribeWorkflowExecutionResponse { repeated temporal.api.workflow.v1.CallbackInfo callbacks = 6; repeated temporal.api.workflow.v1.PendingNexusOperationInfo pending_nexus_operations = 7; temporal.api.workflow.v1.WorkflowExecutionExtendedInfo workflow_extended_info = 8; + // The notification channels this run stands on: the independent channels it subscribed to and + // the channels linked to it that hold any state. Empty when there are none. + repeated temporal.api.workflow.v1.ChannelSubscriptionInfo channel_subscriptions = 9; } // (-- api-linter: core::0203::optional=disabled diff --git a/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto b/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto index 34f6a73f6..b2ec8a5e9 100644 --- a/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto +++ b/crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/service.proto @@ -483,6 +483,142 @@ service WorkflowService { }; } + // NotifyChannel tells every listener of a channel that a source they consume + // has moved. The writer names no addressee and never learns who listens. The + // server wakes each listener: a Workflow with a Workflow Task, a callback by + // invoking it. Nothing goes to History except the notifications a woken + // Workflow Task carries on its scheduled event. + rpc NotifyChannel (NotifyChannelRequest) returns (NotifyChannelResponse) { + option (google.api.http) = { + post: "/namespaces/{namespace}/channels/{notification.channel}/notify" + body: "*" + additional_bindings { + post: "/api/v1/namespaces/{namespace}/channels/{notification.channel}/notify" + body: "*" + } + additional_bindings { + post: "/namespaces/{namespace}/workflows/{execution.business_id}/channels/{notification.channel}/notify" + body: "*" + } + additional_bindings { + post: "/api/v1/namespaces/{namespace}/workflows/{execution.business_id}/channels/{notification.channel}/notify" + body: "*" + } + additional_bindings { + post: "/namespaces/{namespace}/activities/{execution.business_id}/channels/{notification.channel}/notify" + body: "*" + } + additional_bindings { + post: "/api/v1/namespaces/{namespace}/activities/{execution.business_id}/channels/{notification.channel}/notify" + body: "*" + } + }; + } + + // RegisterChannelListener registers a callback as a listener of a channel. A + // Workflow registers itself with the `SubscribeNotificationChannel` command + // instead. + rpc RegisterChannelListener (RegisterChannelListenerRequest) + returns (RegisterChannelListenerResponse) { + option (google.api.http) = { + post: "/namespaces/{namespace}/channels/{channel}/listeners" + body: "*" + additional_bindings { + post: "/api/v1/namespaces/{namespace}/channels/{channel}/listeners" + body: "*" + } + additional_bindings { + post: "/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}/listeners" + body: "*" + } + additional_bindings { + post: "/api/v1/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}/listeners" + body: "*" + } + additional_bindings { + post: "/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}/listeners" + body: "*" + } + additional_bindings { + post: "/api/v1/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}/listeners" + body: "*" + } + }; + } + + // UnregisterChannelListener removes a listener from a channel. + // + // (-- api-linter: core::0136::http-method=disabled + // aip.dev/not-precedent: Removing a listener is a delete of that listener. --) + rpc UnregisterChannelListener (UnregisterChannelListenerRequest) + returns (UnregisterChannelListenerResponse) { + option (google.api.http) = { + delete: "/namespaces/{namespace}/channels/{channel}/listeners/{listener_id}" + additional_bindings { + delete: "/api/v1/namespaces/{namespace}/channels/{channel}/listeners/{listener_id}" + } + additional_bindings { + delete: "/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}/listeners/{listener_id}" + } + additional_bindings { + delete: "/api/v1/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}/listeners/{listener_id}" + } + additional_bindings { + delete: "/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}/listeners/{listener_id}" + } + additional_bindings { + delete: "/api/v1/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}/listeners/{listener_id}" + } + }; + } + + // PollChannel is a long poll for clients. It returns the retained + // notifications of a channel with a counter above `after_counter`, waiting + // up to `wait` for one when none is retained yet. + rpc PollChannel (PollChannelRequest) returns (PollChannelResponse) { + option (google.api.http) = { + get: "/namespaces/{namespace}/channels/{channel}/notifications" + additional_bindings { + get: "/api/v1/namespaces/{namespace}/channels/{channel}/notifications" + } + additional_bindings { + get: "/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}/notifications" + } + additional_bindings { + get: "/api/v1/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}/notifications" + } + additional_bindings { + get: "/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}/notifications" + } + additional_bindings { + get: "/api/v1/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}/notifications" + } + }; + } + + // DescribeChannel returns the listeners of a channel and its latest + // notification. + rpc DescribeChannel (DescribeChannelRequest) returns (DescribeChannelResponse) { + option (google.api.http) = { + get: "/namespaces/{namespace}/channels/{channel}" + additional_bindings { + get: "/api/v1/namespaces/{namespace}/channels/{channel}" + } + additional_bindings { + get: "/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}" + } + additional_bindings { + get: "/api/v1/namespaces/{namespace}/workflows/{execution.business_id}/channels/{channel}" + } + additional_bindings { + get: "/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}" + } + additional_bindings { + get: "/api/v1/namespaces/{namespace}/activities/{execution.business_id}/channels/{channel}" + } + }; + } + // SignalWithStartWorkflowExecution is used to ensure a signal is sent to a workflow, even if // it isn't yet started. // 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 7287ab1c1..7aa7becc0 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,8 +13,10 @@ import "google/protobuf/empty.proto"; import "temporal/api/failure/v1/message.proto"; import "temporal/api/update/v1/message.proto"; import "temporal/api/common/v1/message.proto"; +import "temporal/api/stream/v1/message.proto"; import "temporal/api/enums/v1/failed_cause.proto"; import "temporal/api/enums/v1/workflow.proto"; +import "temporal/api/notification/v1/message.proto"; import "temporal/sdk/core/activity_result/activity_result.proto"; import "temporal/sdk/core/child_workflow/child_workflow.proto"; import "temporal/sdk/core/common/common.proto"; @@ -140,6 +142,14 @@ 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; + // Notifications from the channels the workflow subscribed to. + NotificationsReceived notifications_received = 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. @@ -147,6 +157,34 @@ 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; +} + +// Hand a workflow the notifications the server folded for its channels. +// +// They come from the scheduled event of the Workflow Task this activation +// belongs to. History is the record, so a replay yields the same job with the +// same notifications at the same point. +message NotificationsReceived { + repeated temporal.api.notification.v1.Notification notifications = 1; +} + // 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 f36c84cef..282c1efea 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 @@ -54,9 +54,32 @@ message WorkflowCommand { UpdateResponse update_response = 20; ScheduleNexusOperation schedule_nexus_operation = 21; RequestCancelNexusOperation request_cancel_nexus_operation = 22; + // 23 to 30 are taken by the stream commands, which share this message. + SubscribeNotificationChannel subscribe_notification_channel = 31; + UnsubscribeNotificationChannel unsubscribe_notification_channel = 32; } } +// Subscribe this workflow to a notification channel, so the scheduled event of +// each later Workflow Task carries the notifications folded for it. +// +// The notifications live in History rather than arriving by a side channel, so +// a replay reads the same ones the live run saw. +message SubscribeNotificationChannel { + // Name of the channel, scoped to the namespace. + string channel = 1; +} + +// End this workflow's subscription to a notification channel. Notifications +// already recorded on a scheduled event still reach that Workflow Task. +// +// The server records an event for every unsubscribe, also one naming a channel +// the run is not subscribed to, so replay can hold each command to its event. +message UnsubscribeNotificationChannel { + // Name of the channel, scoped to the namespace. + string channel = 1; +} + 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 39733b600..a0c050f88 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -1277,6 +1277,16 @@ 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()) + } } } } @@ -1477,6 +1487,18 @@ pub mod coresdk { use crate::protos::temporal::api::{common::v1::Payloads, enums::v1::QueryResultType}; use std::fmt::{Display, Formatter}; + impl Display for SubscribeNotificationChannel { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "SubscribeNotificationChannel({})", self.channel) + } + } + + impl Display for UnsubscribeNotificationChannel { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "UnsubscribeNotificationChannel({})", self.channel) + } + } + impl Display for WorkflowCommand { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { match &self.variant { @@ -1825,6 +1847,18 @@ pub mod temporal { CommandType::ScheduleActivityTask } Attributes::StartTimerCommandAttributes(_) => CommandType::StartTimer, + Attributes::SubscribeStreamCommandAttributes(_) => { + CommandType::SubscribeStream + } + Attributes::AppendStreamRecordsCommandAttributes(_) => { + CommandType::AppendStreamRecords + } + Attributes::SubscribeNotificationChannelCommandAttributes(_) => { + CommandType::SubscribeNotificationChannel + } + Attributes::UnsubscribeNotificationChannelCommandAttributes(_) => { + CommandType::UnsubscribeNotificationChannel + } Attributes::CompleteWorkflowExecutionCommandAttributes(_) => { CommandType::CompleteWorkflowExecution } @@ -2337,6 +2371,10 @@ pub mod temporal { | EventType::TimerStarted | EventType::UpsertWorkflowSearchAttributes | EventType::WorkflowPropertiesModified + | EventType::WorkflowStreamSubscribed + | EventType::WorkflowStreamRecordsAppended + | EventType::WorkflowNotificationChannelSubscribed + | EventType::WorkflowNotificationChannelUnsubscribed | EventType::NexusOperationScheduled | EventType::NexusOperationCancelRequested | EventType::WorkflowExecutionCanceled @@ -2436,6 +2474,16 @@ 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 + } + Attributes::WorkflowNotificationChannelUnsubscribedEventAttributes(_) => { + false + } Attributes::WorkflowExecutionStartedEventAttributes(_) => false, Attributes::WorkflowExecutionCompletedEventAttributes(_) => false, Attributes::WorkflowExecutionFailedEventAttributes(_) => false, @@ -2523,6 +2571,10 @@ 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 } Attributes::WorkflowExecutionCompletedEventAttributes(_) => { EventType::WorkflowExecutionCompleted } Attributes::WorkflowExecutionFailedEventAttributes(_) => { EventType::WorkflowExecutionFailed } @@ -2594,6 +2646,11 @@ pub mod temporal { tonic::include_proto!("temporal.api.namespace.v1"); } } + pub mod notification { + pub mod v1 { + tonic::include_proto!("temporal.api.notification.v1"); + } + } pub mod operatorservice { pub mod v1 { tonic::include_proto!("temporal.api.operatorservice.v1"); @@ -2640,6 +2697,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-c-bridge/src/client.rs b/crates/sdk-core-c-bridge/src/client.rs index 580874fd3..a0ae53b82 100644 --- a/crates/sdk-core-c-bridge/src/client.rs +++ b/crates/sdk-core-c-bridge/src/client.rs @@ -669,6 +669,9 @@ async fn call_workflow_service( "DescribeBatchOperation" => { rpc_call_on_trait!(client, call, WorkflowService, describe_batch_operation) } + "DescribeChannel" => { + rpc_call_on_trait!(client, call, WorkflowService, describe_channel) + } "DescribeDeployment" => { rpc_call_on_trait!(client, call, WorkflowService, describe_deployment) } @@ -802,6 +805,9 @@ async fn call_workflow_service( "ListWorkflowRules" => { rpc_call_on_trait!(client, call, WorkflowService, list_workflow_rules) } + "NotifyChannel" => { + rpc_call_on_trait!(client, call, WorkflowService, notify_channel) + } "PatchSchedule" => rpc_call_on_trait!(client, call, WorkflowService, patch_schedule), "PauseActivity" => rpc_call_on_trait!(client, call, WorkflowService, pause_activity), "PauseActivityExecution" => { @@ -810,6 +816,9 @@ async fn call_workflow_service( "PauseWorkflowExecution" => { rpc_call_on_trait!(client, call, WorkflowService, pause_workflow_execution) } + "PollChannel" => { + rpc_call_on_trait!(client, call, WorkflowService, poll_channel) + } "PollActivityExecution" => { rpc_call_on_trait!(client, call, WorkflowService, poll_activity_execution) } @@ -862,6 +871,9 @@ async fn call_workflow_service( "RecordWorkerHeartbeat" => { rpc_call_on_trait!(client, call, WorkflowService, record_worker_heartbeat) } + "RegisterChannelListener" => { + rpc_call_on_trait!(client, call, WorkflowService, register_channel_listener) + } "RegisterNamespace" => { rpc_call_on_trait!(client, call, WorkflowService, register_namespace) } @@ -1030,6 +1042,9 @@ async fn call_workflow_service( "TriggerWorkflowRule" => { rpc_call_on_trait!(client, call, WorkflowService, trigger_workflow_rule) } + "UnregisterChannelListener" => { + rpc_call_on_trait!(client, call, WorkflowService, unregister_channel_listener) + } "UnpauseActivity" => { rpc_call_on_trait!(client, call, WorkflowService, unpause_activity) } diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index c0a8d74b9..ecee07383 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -1398,7 +1398,7 @@ fn validate_completion( .collect::, EmptyWorkflowCommandErr>>() .map_err(|_| CompleteWfError::MalformedWorkflowCompletion { reason: "At least one workflow command in the completion contained \ - an empty variant" + an empty or unsupported variant" .to_owned(), run_id: completion.run_id.clone(), })?; @@ -1480,7 +1480,7 @@ impl LocalResolution { } #[derive(thiserror::Error, Debug, derive_more::From)] -#[error("Lang provided workflow command with empty variant")] +#[error("Lang provided workflow command with an empty or unsupported variant")] struct EmptyWorkflowCommandErr; /// [DrivenWorkflow]s respond with these when called, to indicate what they want to do next. @@ -1634,6 +1634,16 @@ impl TryFrom for WFCommand { workflow_command::Variant::RequestCancelNexusOperation(s) => { WFCommandVariant::RequestCancelNexusOperation(s) } + // This layer carries the protos only. Dropping the command would let the + // workflow go on as if subscribed while the server never heard of it. + workflow_command::Variant::SubscribeNotificationChannel(_) => { + return Err(EmptyWorkflowCommandErr); + } + // Same reason, reversed: the workflow would go on as if unsubscribed + // while the server keeps delivering. + workflow_command::Variant::UnsubscribeNotificationChannel(_) => { + return Err(EmptyWorkflowCommandErr); + } }; Ok(Self { variant, diff --git a/crates/sdk/src/workflow_future.rs b/crates/sdk/src/workflow_future.rs index f18604ce0..68085a5a5 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -336,6 +336,21 @@ 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. + debug!("Channel notifications received and ignored"); + push_polled_context!(ActivationJobContext::Passive); + } 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 7c6b3a033..4dfb4159e 100644 --- a/crates/workflow/src/runtime/instance.rs +++ b/crates/workflow/src/runtime/instance.rs @@ -1118,6 +1118,25 @@ 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. + ActivationJobResult::None + } Some(ActivationVariant::RemoveFromCache(_)) => ActivationJobResult::None, None => { return Err(Box::new(Failure {