Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
5 changes: 5 additions & 0 deletions crates/common/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
1 change: 1 addition & 0 deletions crates/protos/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
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 @@ -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;
}
}

Expand All @@ -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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, temporal.api.common.v1.Payload> metadata = 2;
Expand All @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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.
Expand All @@ -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
Expand Down
Loading
Loading