Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
85b71d7
Update api_upstream
tconley1428 Aug 28, 2026
5209468
Carried the stream protos in the Core's api tree.
moedash Sep 22, 2026
98d9621
Matched the nexus model WIT to upstream for the workflow-id policies.
moedash Sep 25, 2026
bc0c1c9
Honored the per-runner timeout for the integ test matrix.
moedash Sep 25, 2026
9163b11
Vendored the api's stream protos into Core.
moedash Sep 25, 2026
7603658
Added the stream delivery job to the bridge protos.
moedash Sep 25, 2026
2e05ed4
Carried the review round's stream proto changes.
moedash Sep 25, 2026
56750a5
Renamed the stream command and event fields.
moedash Sep 25, 2026
b512279
Merged the current upstream Core into the stream protos branch.
moedash Sep 25, 2026
b3c75b5
Vendored the api's subscribe start position.
moedash Sep 28, 2026
0b69924
Carried the subscribe start position in the Core's api tree.
moedash Sep 28, 2026
1410ad5
Carried the wake protos in the Core's api tree.
moedash Oct 1, 2026
1e35b60
Merged the wake protos into the Core's protos branch.
moedash Oct 1, 2026
0ab3810
Exposed the wake call on the raw workflow client.
moedash Oct 1, 2026
209ca4b
Revert "Exposed the wake call on the raw workflow client."
moedash Oct 1, 2026
27690d3
Exposed the wake call on the raw workflow client.
moedash Oct 1, 2026
88a116c
Merged the wake client call from the protos branch.
moedash Oct 1, 2026
3fafca0
Carried the wake's fold-rule wording in the Core's api tree.
moedash Oct 1, 2026
7e22ab7
Merged the wake's fold-rule wording into the Core's protos branch.
moedash Oct 1, 2026
ae087c6
Dispatched the wake call in the C bridge.
moedash Oct 1, 2026
8085667
Carried the notification channel protos and client calls in the Core'…
moedash Oct 2, 2026
5aef7fb
Carried the wrapped notification channel lines in the Core's api tree.
moedash Oct 2, 2026
c8ba0a3
Merge commit '8085667b' into moe/AI-198-core-1-protos
moedash Oct 2, 2026
43f2bd1
Carried the channel command and job protos for the Python bridge.
moedash Oct 2, 2026
b481c81
Carried the subscribe-notification-channel failed cause in the Core's…
moedash Oct 2, 2026
432b2f0
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
moedash Oct 2, 2026
a9a025b
Compiled the channel command and job without handling them on the pro…
moedash Oct 2, 2026
678864b
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
moedash Oct 2, 2026
23307d8
Dropped the point-to-point wake from the Core's api tree and client.
moedash Oct 2, 2026
41d0dd6
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
moedash Oct 2, 2026
0c21f1c
Carried the linked channel contract in the Core's api tree.
moedash Oct 2, 2026
59f1db1
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
moedash Oct 2, 2026
8103755
Carried the channel subscriptions on describe in the Core's api tree.
moedash Oct 2, 2026
acd3091
Carried the unsubscribe channel command in the Core's api tree.
moedash Oct 2, 2026
83f965c
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
moedash Oct 2, 2026
7a0d774
Addressed a linked channel by execution in the Core's api tree.
moedash Oct 3, 2026
5ae2095
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
moedash Oct 3, 2026
fbc064b
Reused the old numbers for the execution fields in the Core's api tree.
moedash Oct 3, 2026
56113a0
Merge remote-tracking branch 'moedash/moe/AI-198-api-protos-on-85b71d…
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
2 changes: 1 addition & 1 deletion .github/workflows/per-pr.yml
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,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
45 changes: 45 additions & 0 deletions crates/client/src/grpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -955,6 +955,51 @@ proxier! {
r.extensions_mut().insert(labels);
}
);
(
notify_channel,
NotifyChannelRequest,
NotifyChannelResponse,
|r| {
let labels = namespaced_request!(r);
r.extensions_mut().insert(labels);
}
);
(
register_channel_listener,
RegisterChannelListenerRequest,
RegisterChannelListenerResponse,
|r| {
let labels = namespaced_request!(r);
r.extensions_mut().insert(labels);
}
);
(
unregister_channel_listener,
UnregisterChannelListenerRequest,
UnregisterChannelListenerResponse,
|r| {
let labels = namespaced_request!(r);
r.extensions_mut().insert(labels);
}
);
(
poll_channel,
PollChannelRequest,
PollChannelResponse,
|r| {
let labels = namespaced_request!(r);
r.extensions_mut().insert(labels);
}
);
(
describe_channel,
DescribeChannelRequest,
DescribeChannelResponse,
|r| {
let labels = namespaced_request!(r);
r.extensions_mut().insert(labels);
}
);
(
signal_with_start_workflow_execution,
SignalWithStartWorkflowExecutionRequest,
Expand Down
8 changes: 8 additions & 0 deletions crates/common/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -829,6 +829,14 @@ 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",
// 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
2 changes: 2 additions & 0 deletions crates/protos/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,15 @@ const SERDE_DERIVE_PREFIXES: &[&str] = &[
".temporal.api.nexus",
".temporal.api.nexusoperation",
".temporal.api.nexusservices",
".temporal.api.notification",
".temporal.api.operatorservice",
".temporal.api.protocol",
".temporal.api.query",
".temporal.api.replication",
".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 @@ -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;
}
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;
}
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;
}
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
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
Loading
Loading