Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
f5164b3
Delivered server-side stream ranges to workflows.
moedash Aug 27, 2026
05a34e8
Let a workflow subscribe to a stream from lang.
moedash Aug 27, 2026
3b4e95e
Let a workflow publish to a stream from lang.
moedash Sep 14, 2026
d0884a3
Replayed recorded stream ranges before matching commands.
moedash Sep 14, 2026
280fe2d
Registered the stream machines with the coverage reporter.
moedash Sep 16, 2026
865805b
Matched the stream test mocks to the two-argument completion.
moedash Sep 16, 2026
f88f78c
Honored the per-runner timeout for the integ test matrix.
moedash Sep 16, 2026
4076ef2
Matched the nexus model WIT to upstream for the workflow-id policies.
moedash Sep 17, 2026
d6d66fb
Kept the closing completion visible to a history update's last task.
moedash Sep 18, 2026
76163d4
Checked replayed stream commands against the recorded event.
moedash Sep 18, 2026
83ab23f
Gave a data-only stream task its own replay activation.
moedash Sep 18, 2026
2342913
Ordered stream ranges and checked re-supplied slices against History.
moedash Sep 18, 2026
0cc74f7
Covered the stream delivery guards and each range's activation.
moedash Sep 19, 2026
6cf953e
Renumbered the stream protos to shared numbers and moved a doc comment.
moedash Sep 19, 2026
a677812
Added the changelog entry for the native stream commands and job.
moedash Sep 19, 2026
97cd172
Renamed the stream command, job and record to the api's record vocabu…
moedash Sep 21, 2026
f8c5fa8
Carried stream slices on a history pushed for replay.
moedash Sep 21, 2026
61b0cae
Failed a task owed stream records the response did not carry.
moedash Sep 21, 2026
6446174
Withheld a legacy query owed records the response did not carry.
moedash Sep 21, 2026
32c6669
Matched the stream protos to the api branch head.
moedash Sep 22, 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
2 changes: 1 addition & 1 deletion .github/workflows/per-pr.yml
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,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:
Expand Down
5 changes: 5 additions & 0 deletions crates/common/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -820,6 +820,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",
// 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",
Expand Down
1 change: 1 addition & 0 deletions crates/protos/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,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 @@ -147,24 +147,15 @@ interface model {
/// typescript="common.WorkflowIdReusePolicy"
/// dotnet="Temporalio.Api.Enums.V1.WorkflowIdReusePolicy"
/// typescript-import="@temporalio/common"
enum workflow-id-reuse-policy {
allow-duplicate,
allow-duplicate-failed-only,
reject-duplicate,
terminate-if-running,
}
type workflow-id-reuse-policy = placeholder;

/// @nexus.proto "temporal.api.enums.v1.WorkflowIdConflictPolicy" typescript-import="@temporalio/proto"
/// @nexus.type
/// python="temporalio.common.WorkflowIDConflictPolicy"
/// typescript="common.WorkflowIdConflictPolicy"
/// dotnet="Temporalio.Api.Enums.V1.WorkflowIdConflictPolicy"
/// typescript-import="@temporalio/common"
enum workflow-id-conflict-policy {
fail,
use-existing,
terminate-existing,
}
type workflow-id-conflict-policy = placeholder;

/// @nexus.proto "temporal.api.sdk.v1.UserMetadata" typescript-import="@temporalio/proto"
/// @nexus.flatten-in-api
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 @@ -324,5 +325,38 @@ 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;
}
}

// 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 {
// Empty means the Workflow's default output stream.
string stream_id = 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. A stream in another execution is addressed by its id;
// one this Workflow owns is addressed by the name it was published under.
// The server resolves an owned name first and falls back to a standalone
// id, so a Workflow that owns a stream under this name cannot reach a
// standalone stream with the same id.
string stream_id = 1;
// Where to start. Negative means from wherever the stream is when the
// subscription is registered, which the server resolves and records so
// replay does not resolve it again.
int64 start_offset = 2;
}
Original file line number Diff line number Diff line change
Expand Up @@ -29,4 +29,6 @@ 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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -175,4 +175,12 @@ 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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,14 @@ 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;
}

enum StartChildWorkflowExecutionFailedCause {
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/deployment/v1/message.proto";
import "temporal/api/failure/v1/message.proto";
import "temporal/api/taskqueue/v1/message.proto";
Expand Down Expand Up @@ -379,6 +380,13 @@ 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.
// Numbered 20 to leave 14 through 19 free for fields added on the main line.
repeated temporal.api.stream.v1.StreamRange consumed_stream_ranges = 20;
}

message WorkflowTaskTimedOutEventAttributes {
Expand Down Expand Up @@ -953,6 +961,33 @@ 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;
// Stream the Workflow subscribed to.
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;
// Stream the Workflow appended to.
string stream_id = 2;
// Offset the first record of the batch landed at.
int64 first_offset = 3;
// How many records the batch held. With first_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.
int64 record_count = 4;
}

message WorkflowExecutionUpdateAcceptedEventAttributes {
// The instance ID of the update protocol that generated this event.
string protocol_instance_id = 1;
Expand Down Expand Up @@ -1276,6 +1311,8 @@ 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;
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
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 here as it does to any other payload.
temporal.api.common.v1.Payload body = 1;
// Producer-supplied provenance, stored as sent.
map<string, temporal.api.common.v1.Payload> 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, or -1 when unnumbered.
// Stored as sent; the server does not assign, validate or order by it.
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 {
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.
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.
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 {
string stream_id = 1;
// Inclusive.
int64 from_offset = 2;
// Exclusive.
int64 to_offset = 3;
}

// 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;
}
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/history/v1/message.proto";
import "temporal/api/workflow/v1/message.proto";
import "temporal/api/command/v1/message.proto";
Expand Down Expand Up @@ -383,6 +384,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/enums/v1/workflow.proto";
import "temporal/sdk/core/activity_result/activity_result.proto";
import "temporal/sdk/core/child_workflow/child_workflow.proto";
Expand Down Expand Up @@ -139,13 +140,38 @@ 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;
// 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.
RemoveFromCache remove_from_cache = 50;
}
}

// 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