Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
42 commits
Select commit Hold shift + click to select a range
76e2931
Delivered server-side stream ranges to workflows.
moedash Aug 27, 2026
05c910c
Let a workflow subscribe to a stream from lang.
moedash Aug 27, 2026
287ba43
Let a workflow publish to a stream from lang.
moedash Sep 14, 2026
338beea
Replayed recorded stream ranges before matching commands.
moedash Sep 14, 2026
a1c15f4
Registered the stream machines with the coverage reporter.
moedash Sep 16, 2026
d2651a2
Kept the closing completion visible to a history update's last task.
moedash Sep 18, 2026
7e55e39
Checked replayed stream commands against the recorded event.
moedash Sep 18, 2026
a02ccce
Gave a data-only stream task its own replay activation.
moedash Sep 18, 2026
0a45a5e
Ordered stream ranges and checked re-supplied slices against History.
moedash Sep 18, 2026
e32ea19
Covered the stream delivery guards and each range's activation.
moedash Sep 19, 2026
1519aee
Renumbered the stream protos to shared numbers and moved a doc comment.
moedash Sep 19, 2026
419b548
Added the changelog entry for the native stream commands and job.
moedash Sep 19, 2026
6d1e2bc
Renamed the stream command, job and record to the api's record vocabu…
moedash Sep 21, 2026
268e004
Carried stream slices on a history pushed for replay.
moedash Sep 21, 2026
3be51c8
Failed a task owed stream records the response did not carry.
moedash Sep 21, 2026
f32dee0
Withheld a legacy query owed records the response did not carry.
moedash Sep 21, 2026
a2f49db
Merged the unified branch with the stream protos into the all branch.
moedash Sep 22, 2026
21ad8d0
Held the stream commands to more of what their events record.
moedash Sep 25, 2026
7622031
Named the unreachable stream machine responses fatal.
moedash Sep 25, 2026
ab70e2e
Failed a re-supply that disagrees with History as the worker's.
moedash Sep 25, 2026
473a1a4
Scoped the data-only task flag to the task it describes.
moedash Sep 25, 2026
a7c3b1c
Said why the extra boundary page cannot be fetched more narrowly.
moedash Sep 25, 2026
93d7678
Restored the comment break the append proto lost.
moedash Sep 25, 2026
20064ee
Kept the command match arms in their original order.
moedash Sep 25, 2026
0c45d00
Merged the repairs and unified review round into the all branch.
moedash Sep 25, 2026
618e8a0
Dropped the subscribe start offset check as unsound.
moedash Sep 25, 2026
04ce5c2
Renamed the stream command and event fields.
moedash Sep 25, 2026
748f1c3
Renamed the lang stream command fields to match the wire.
moedash Sep 25, 2026
32ec55d
Carried the stream field rename into the delivery tests.
moedash Sep 25, 2026
d2b9dac
Merged the current upstream Core into the all branch.
moedash Sep 25, 2026
a55d2d9
Moved the stream changelog entries back under Unreleased.
moedash Sep 26, 2026
6da86ea
Vendored the api's subscribe start position.
moedash Sep 28, 2026
e03957e
Passed the subscribe start position from lang to the server.
moedash Sep 28, 2026
6ae3829
Merge commit 'f62993b1' into to-all
moedash Oct 1, 2026
899511e
Merge commit '736ecc8f' into to-all
moedash Oct 2, 2026
614cb65
Merge commit '741ff2bb' into to-all
moedash Oct 2, 2026
5b4b857
Merge commit '1de8b920' into to-all
moedash Oct 2, 2026
d7684be
Merge commit 'ba05a851' into to-all
moedash Oct 2, 2026
c3b4b35
Merge commit 'c3d0a4b3' into to-all
moedash Oct 2, 2026
67426be
Merge commit '12b429e9' into to-all
moedash Oct 2, 2026
c017e35
Merged the execution-addressed unified core into to-all.
moedash Oct 3, 2026
0871616
Merged the renumbered unified core into to-all.
moedash 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
Original file line number Diff line number Diff line change
Expand Up @@ -339,8 +339,11 @@ message Command {
// 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;
// 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;
Expand All @@ -353,16 +356,22 @@ message AppendStreamRecordsCommandAttributes {
// 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.
// 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -390,8 +390,11 @@ message WorkflowTaskCompletedEventAttributes {
// 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;

// 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 @@ -972,7 +975,9 @@ 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.
// 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
Expand Down Expand Up @@ -1005,14 +1010,20 @@ 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.
// Name of the 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;
// 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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +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 payload codec
// applies here as it does to any other payload.
// 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 @@ -28,23 +33,31 @@ message StreamRecord {
// 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.
// 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
Expand All @@ -62,13 +75,39 @@ message StreamSlice {
// 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 @@ -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/failed_cause.proto";
import "temporal/api/enums/v1/workflow.proto";
import "temporal/api/notification/v1/message.proto";
Expand Down Expand Up @@ -156,6 +157,11 @@ message WorkflowActivationJob {
// Runtime-internal: encode the terminal boundary for a marker Core is about to write.
// Runs no user code.
FinalizeExternalStreams finalize_external_streams = 20;
// The number below is shared with the native stream tree, which leaves 17
// to 20 to the jobs above, 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
Expand Down Expand Up @@ -231,6 +237,25 @@ message FinalizeExternalStreams {
coresdk.external_data.ParkReason reason = 3;
}

// 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
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import "temporal/api/common/v1/message.proto";
import "temporal/api/enums/v1/workflow.proto";
import "temporal/api/failure/v1/message.proto";
import "temporal/api/sdk/v1/user_metadata.proto";
import "temporal/api/stream/v1/message.proto";
import "temporal/api/sdk/v1/event_group_marker.proto";
import "temporal/sdk/core/child_workflow/child_workflow.proto";
import "temporal/sdk/core/nexus/nexus.proto";
Expand Down Expand Up @@ -61,6 +62,10 @@ message WorkflowCommand {
ExternalStreamFinalized external_stream_finalized = 26;
WorkflowOutputStreamCommit workflow_output_stream_commit = 27;
WorkflowOutputStreamBuffered workflow_output_stream_buffered = 28;
// The two numbers below are shared with the native stream tree, which
// leaves 23 to 28 to the commands above, and must not be reused.
SubscribeStream subscribe_stream = 29;
AppendStreamRecords append_stream_records = 30;
SubscribeNotificationChannel subscribe_notification_channel = 31;
UnsubscribeNotificationChannel unsubscribe_notification_channel = 32;
WorkflowStreamChannels workflow_stream_channels = 33;
Expand Down Expand Up @@ -186,6 +191,45 @@ message WorkflowOutputStreamBuffered {
google.protobuf.Duration max_publish_latency = 1;
}

// Append a batch of records to a stream this workflow owns.
//
// The bodies go to the stream's own log rather than into History, which gets
// one fixed-size event naming the offset range. That is what makes the batch
// size free: a thousand records cost the same in History as one. The server
// stores each record with an empty producer id, because the workflow is the
// producer here.
message AppendStreamRecords {
// 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;
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. 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 SubscribeStream {
// 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. The server refuses a negative value.
int64 start_offset = 2;
// Where to start: an offset, the earliest record held, the tail, or the
// last N records. The server resolves it and records the resulting offset,
// so replay does not resolve it again and Core does not keep it.
temporal.api.stream.v1.StreamStartPosition start_position = 3;
}

message StartTimer {
// Lang's incremental sequence number, used as the operation identifier
uint32 seq = 1;
Expand Down
47 changes: 47 additions & 0 deletions crates/protos/src/protos/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1374,6 +1374,13 @@ pub mod coresdk {
fin.reason()
)
}
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())
}
Expand Down Expand Up @@ -1586,6 +1593,23 @@ pub mod coresdk {
}
}

impl Display for SubscribeStream {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "SubscribeStream({})", self.stream_name_or_id)
}
}

impl Display for AppendStreamRecords {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(
f,
"AppendStreamRecords({}, {} records)",
self.stream_name,
self.records.len()
)
}
}

impl Display for StartTimer {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "StartTimer({})", self.seq)
Expand Down Expand Up @@ -2094,6 +2118,29 @@ pub mod temporal {
}
}

impl From<workflow_commands::AppendStreamRecords> for command::Attributes {
fn from(s: workflow_commands::AppendStreamRecords) -> Self {
Self::AppendStreamRecordsCommandAttributes(
AppendStreamRecordsCommandAttributes {
stream_name: s.stream_name,
records: s.records,
},
)
}
}

impl From<workflow_commands::SubscribeStream> for command::Attributes {
fn from(s: workflow_commands::SubscribeStream) -> Self {
Self::SubscribeStreamCommandAttributes(
SubscribeStreamCommandAttributes {
stream_name_or_id: s.stream_name_or_id,
start_offset: s.start_offset,
start_position: s.start_position,
},
)
}
}

impl From<workflow_commands::SubscribeNotificationChannel> for Attributes {
fn from(s: workflow_commands::SubscribeNotificationChannel) -> Self {
Self::SubscribeNotificationChannelCommandAttributes(
Expand Down
Loading
Loading