Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
9109f60
Added the stream messages, commands, events, and task carriage.
moetemp Sep 14, 2026
7caf5c3
Exempted the stream proto from the prepositions rule.
moetemp Sep 16, 2026
2b05496
Regenerated the OpenAPI specs.
moetemp Sep 16, 2026
e65008a
Added workflow task failure causes for the stream commands.
moetemp Sep 18, 2026
228b478
Documented the stream proto fields and relocated stream_cursors.
moetemp Sep 18, 2026
1aace12
Regenerated the OpenAPI specs.
moetemp Sep 18, 2026
43546df
Named the stream record on the wire and carried its producer identity.
moetemp Sep 21, 2026
5efc277
Renamed the stream command, event and range to the record vocabulary.
moetemp Sep 21, 2026
2cffd47
Regenerated the OpenAPI specs.
moetemp Sep 21, 2026
f2156b1
Named the stream command fields for what they address.
moetemp Sep 25, 2026
b3947f9
Reserved the history numbers held for the main line.
moetemp Sep 25, 2026
92b0399
Gave the appended event the range fields every range has.
moetemp Sep 25, 2026
9df53d9
Dropped the unnumbered sequence sentinel.
moetemp Sep 25, 2026
35f43bd
Scoped the payload codec claim to the paths this API owns.
moetemp Sep 25, 2026
6e49c7b
Regenerated the OpenAPI specs.
moetemp Sep 25, 2026
656645b
Added a start position to the subscribe command.
moetemp Sep 28, 2026
75de7aa
Added the wake, a coalesced event-less Workflow Task request.
moetemp Oct 1, 2026
086d85d
Described the wake's fold rule as pending-and-undelivered only.
moetemp Oct 1, 2026
3f48439
Regenerated the OpenAPI specs for the wake's fold rule.
moetemp Oct 1, 2026
2313c60
Added the notification channel: notify, listen, poll and describe.
moetemp Oct 2, 2026
d0c100d
Wrapped the notification channel lines that ran past 100 characters.
moetemp Oct 2, 2026
5fa4026
Added the failed cause for a refused SubscribeNotificationChannel com…
moetemp Oct 2, 2026
9cd8b40
Removed the point-to-point wake call and message.
moetemp Oct 2, 2026
a9e6517
Added the channel linked to a workflow to the notification contract.
moetemp Oct 2, 2026
30918a0
Added a workflow's channel subscriptions to DescribeWorkflowExecution.
moetemp Oct 2, 2026
a446be7
Added the UnsubscribeNotificationChannel command and its event.
moetemp Oct 2, 2026
07a46d0
Put the channel first on the unsubscribe event.
moetemp Oct 2, 2026
c239d35
Restored the file's field order on the unsubscribe event.
moetemp Oct 2, 2026
4304fd8
Addressed a linked channel by execution rather than by workflow.
moetemp Oct 3, 2026
071feb8
Reused the old numbers for the execution fields.
moetemp Oct 3, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3,132 changes: 2,787 additions & 345 deletions openapi/openapiv2.json

Large diffs are not rendered by default.

2,414 changes: 2,298 additions & 116 deletions openapi/openapiv3.yaml

Large diffs are not rendered by default.

64 changes: 64 additions & 0 deletions temporal/api/command/v1/message.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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;
}
4 changes: 4 additions & 0 deletions temporal/api/enums/v1/command_type.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
15 changes: 15 additions & 0 deletions temporal/api/enums/v1/event_type.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
13 changes: 13 additions & 0 deletions temporal/api/enums/v1/failed_cause.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
77 changes: 77 additions & 0 deletions temporal/api/history/v1/message.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
}

Expand Down
75 changes: 75 additions & 0 deletions temporal/api/notification/v1/message.proto
Original file line number Diff line number Diff line change
@@ -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<string, temporal.api.common.v1.Payload> 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;
}
Loading
Loading