Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
4e7910b
Added the stream commands to the bridge protos.
moedash Sep 25, 2026
e3e429c
Let a workflow subscribe to a stream and append records to it.
moedash Sep 25, 2026
e0a6d15
Registered the stream machines with the coverage reporter.
moedash Sep 25, 2026
3750076
Covered the stream commands and their replay checks.
moedash Sep 25, 2026
36f218b
Held the stream commands to more of what their events record.
moedash Sep 25, 2026
1e3532a
Named the unreachable stream machine responses fatal.
moedash Sep 25, 2026
99b7609
Dropped the subscribe start offset check as unsound.
moedash Sep 25, 2026
6d63fc3
Merged the stream proto rename into the machines branch.
moedash Sep 25, 2026
5b4088e
Renamed the lang stream command fields to match the wire.
moedash Sep 25, 2026
2ad52df
Merged the current upstream Core into the stream machines branch.
moedash Sep 25, 2026
dac8b78
Merged the subscribe start position protos into the machines branch.
moedash Sep 28, 2026
07fcbfa
Passed the subscribe start position from lang to the server.
moedash Sep 28, 2026
8f9e6a4
Merged the wake protos and client call into the machines branch.
moedash Oct 1, 2026
a611440
Merged the relocated wake client call into the machines branch.
moedash Oct 1, 2026
6bf384d
Merged the wake's fold-rule wording into the machines branch.
moedash Oct 1, 2026
eaffd07
Merged the C bridge wake dispatch into the machines branch.
moedash Oct 1, 2026
4c6314c
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 2, 2026
245d729
Subscribed workflows to notification channels by command.
moedash Oct 2, 2026
52cd208
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 2, 2026
dd12013
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 2, 2026
833d4cf
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 2, 2026
9668471
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 2, 2026
4249775
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 2, 2026
2a20a57
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
moedash Oct 3, 2026
abf6ce9
Merge branch 'moe/AI-198-core-1-protos' into moe/AI-198-core-2-machines
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 @@ -14,9 +14,9 @@ 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";
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";
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 @@ -54,12 +55,55 @@ 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.
// 23 to 28 are taken by the external stream commands, which are developed
// alongside these and share this message. The two numbers below are fixed
// with that family and must not be reused.
SubscribeStream subscribe_stream = 29;
AppendStreamRecords append_stream_records = 30;
SubscribeNotificationChannel subscribe_notification_channel = 31;
UnsubscribeNotificationChannel unsubscribe_notification_channel = 32;
}
}

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

// Subscribe this workflow to a notification channel, so the scheduled event of
// each later Workflow Task carries the notifications folded for it.
//
Expand Down
60 changes: 55 additions & 5 deletions crates/protos/src/protos/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1487,6 +1487,21 @@ pub mod coresdk {
use crate::protos::temporal::api::{common::v1::Payloads, enums::v1::QueryResultType};
use std::fmt::{Display, Formatter};

impl Display for WorkflowCommand {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match &self.variant {
None => write!(f, "Empty"),
Some(v) => write!(f, "{v}"),
}
}
}

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

impl Display for SubscribeNotificationChannel {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "SubscribeNotificationChannel({})", self.channel)
Expand All @@ -1499,12 +1514,14 @@ pub mod coresdk {
}
}

impl Display for WorkflowCommand {
impl Display for AppendStreamRecords {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match &self.variant {
None => write!(f, "Empty"),
Some(v) => write!(f, "{v}"),
}
write!(
f,
"AppendStreamRecords({}, {} records)",
self.stream_name,
self.records.len()
)
}
}

Expand Down Expand Up @@ -1912,6 +1929,39 @@ 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(
SubscribeNotificationChannelCommandAttributes {
channel: s.channel,
},
)
}
}

impl From<workflow_commands::StartTimer> for command::Attributes {
fn from(s: workflow_commands::StartTimer) -> Self {
Self::StartTimerCommandAttributes(StartTimerCommandAttributes {
Expand Down
134 changes: 134 additions & 0 deletions crates/sdk-core/src/core_tests/channels.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
use crate::{
replay::TestHistoryBuilder,
test_help::{MockPollCfg, ResponseType, WorkerTestHelpers, build_mock_pollers, mock_worker},
worker::client::mocks::mock_worker_client,
};
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use temporalio_common::protos::{
coresdk::{
workflow_commands::SubscribeNotificationChannel,
workflow_completion::WorkflowActivationCompletion,
},
temporal::api::{
command::v1::command,
enums::v1::{CommandType, EventType, WorkflowTaskFailedCause},
workflowservice::v1::RespondWorkflowTaskCompletedResponse,
},
};

fn subscribe(channel: &str) -> SubscribeNotificationChannel {
SubscribeNotificationChannel {
channel: channel.to_string(),
}
}

/// The command has to reach the server naming the channel. Only a task that is
/// not being replayed sends commands, so this drives a single open task.
#[tokio::test]
async fn subscribe_channel_command_reaches_the_server() {
let mut t = TestHistoryBuilder::default();
t.add_by_type(EventType::WorkflowExecutionStarted);
t.add_workflow_task_scheduled_and_started();

let mut mock_client = mock_worker_client();
mock_client
.expect_complete_workflow_task()
.times(1)
.returning(|resp, _| {
let cmd = resp.commands.first().expect("a command was sent");
assert_eq!(
cmd.command_type(),
CommandType::SubscribeNotificationChannel
);
match cmd.attributes.as_ref().unwrap() {
command::Attributes::SubscribeNotificationChannelCommandAttributes(a) => {
assert_eq!(a.channel, "orders");
}
other => panic!("wrong attributes: {other:?}"),
}
Ok(RespondWorkflowTaskCompletedResponse::default())
});

let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client);
let core = mock_worker(build_mock_pollers(mock));

let task = core.poll_workflow_activation().await.unwrap();
core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds(
task.run_id,
vec![subscribe("orders").into()],
))
.await
.unwrap();
core.shutdown().await;
}

/// A subscribe reissued on replay has to match the event the original run
/// wrote. Core pops one queued command per command-generated event, so this is
/// what keeps every later command lined up with its own event.
#[tokio::test]
async fn subscribe_channel_command_round_trips_through_replay() {
let mut t = TestHistoryBuilder::default();
t.add_by_type(EventType::WorkflowExecutionStarted);
t.add_full_wf_task();
t.add_notification_channel_subscribed("orders");
t.add_full_wf_task();

let mut mock_client = mock_worker_client();
mock_client
.expect_complete_workflow_task()
.returning(|_, _| Ok(RespondWorkflowTaskCompletedResponse::default()));
mock_client
.expect_fail_workflow_task()
.returning(|_, _, f| panic!("core rejected the reissued subscribe: {f:?}"));

let mock = MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_client);
let core = mock_worker(build_mock_pollers(mock));

let task = core.poll_workflow_activation().await.unwrap();
core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds(
task.run_id,
vec![subscribe("orders").into()],
))
.await
.unwrap();
core.shutdown().await;
}

/// The channel is part of what the recorded event holds the reissued command
/// to. Core sends no commands while replaying, so this check is the only place
/// a subscribe to the wrong channel can be noticed.
#[tokio::test]
async fn a_subscribe_reissued_to_a_different_channel_fails_the_task() {
let mut t = TestHistoryBuilder::default();
t.add_by_type(EventType::WorkflowExecutionStarted);
t.add_full_wf_task();
t.add_notification_channel_subscribed("orders");
t.add_full_wf_task();

let failures = Arc::new(AtomicUsize::new(0));
let counted = failures.clone();
let mut mock =
MockPollCfg::from_resp_batches("wfid", t, [ResponseType::AllHistory], mock_worker_client());
mock.num_expected_fails = 1;
mock.expect_fail_wft_matcher = Box::new(move |_, cause, _| {
counted.fetch_add(1, Ordering::Relaxed);
*cause == WorkflowTaskFailedCause::NonDeterministicError
});
let mut mock = build_mock_pollers(mock);
mock.make_wft_stream_interminable();
let core = mock_worker(mock);

let task = core.poll_workflow_activation().await.unwrap();
core.complete_workflow_activation(WorkflowActivationCompletion::from_cmds(
task.run_id,
vec![subscribe("invoices").into()],
))
.await
.unwrap();
core.handle_eviction().await;
assert_eq!(failures.load(Ordering::Relaxed), 1);
core.shutdown().await;
}
2 changes: 2 additions & 0 deletions crates/sdk-core/src/core_tests/mod.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
mod activity_tasks;
mod channels;
mod event_groups;
mod queries;
mod replay_flag;
mod streams;
mod updates;
mod workers;
mod workflow_cancels;
Expand Down
Loading
Loading