Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
f88293a
Preserve external stream task and replay wake boundaries
mdashti Sep 6, 2026
975bf0e
Retain incomplete stream tasks with workflow caching disabled
mdashti Sep 6, 2026
e25905c
Factored the flush window record into a type alias.
moedash Sep 16, 2026
c81363e
Registered ExternalStreamMachine with the coverage reporter.
moedash Sep 16, 2026
d1b3032
Scoped the Prometheus metrics query to the test's task queue.
mjameswh Aug 14, 2026
7af3a47
Fixed activity cancellation against the current dev server.
mjameswh Aug 14, 2026
f57b074
Matched the workflow-service WIT to upstream for first-execution-run-id.
moedash Sep 16, 2026
9e5676e
Honored the per-runner timeout for the integ test matrix.
moedash Sep 16, 2026
6a7e6b2
Tolerated out-of-grammar nexus request-timeout headers.
mjameswh Aug 14, 2026
31af053
Matched the nexus model WIT to upstream for the workflow-id policies.
moedash Sep 17, 2026
08599ac
Asserted the release of retained zero-cache tasks and a mid-replay wake.
moedash Sep 19, 2026
6d52436
Stated when a stream resolve may be queued behind an activation.
moedash Sep 19, 2026
7df8a5b
Folded the two Fixed sections of the changelog into one.
moedash Sep 19, 2026
8000526
Released the recorded-output lock before the shutdown await.
moedash Sep 22, 2026
c713b5e
Carried the stream protos in the Core's api tree.
moedash Sep 22, 2026
34c1018
Forced a replacement task for output left buffered behind a commit.
moedash Sep 25, 2026
2e8455e
Logged the one-outstanding-activation breach in release builds too.
moedash Sep 25, 2026
ff4e0a6
Carried the wake protos and RPC in the Core's api tree.
moedash Oct 1, 2026
0efd317
Received wakes from the poll response for parked external streams.
moedash Oct 1, 2026
b8e948a
Kept the vendored OpenAPI files at their base.
moedash Oct 1, 2026
8449014
Carried the wake fold rule comments from the api fork.
moedash Oct 1, 2026
30969f9
Carried the notification channel protos and client calls in the Core'…
moedash Oct 2, 2026
6fc86f1
Carried the wrapped notification channel lines in the Core's api tree.
moedash Oct 2, 2026
2fadc9b
Subscribed workflows to notification channels by command.
moedash Oct 2, 2026
1497b92
Carried the subscribe-notification-channel failed cause in the Core's…
moedash Oct 2, 2026
757b985
Handed a failed task's notifications to its retry in one job.
moedash Oct 2, 2026
0a3e3a9
Folded a task's notifications per channel again, highest counter first.
moedash Oct 2, 2026
cdc4260
Removed the point-to-point wake from the Core's api tree, client and …
moedash Oct 2, 2026
bdca403
Carried the linked channel contract in the Core's api tree.
moedash Oct 2, 2026
c6af517
Defaulted the new notification fields in the Core's tests.
moedash Oct 2, 2026
207a295
Subscribed a run's channels on the completion that ends the task.
moedash Oct 2, 2026
0df91dd
Carried the unsubscribe and describe protos in the Core's api tree.
moedash Oct 2, 2026
4afe3e5
Unsubscribed a run from the channels its readers left.
moedash Oct 2, 2026
7006981
Addressed a linked channel by execution in the Core's api tree.
moedash Oct 3, 2026
677d468
Reused the old field numbers for the execution fields in the api tree.
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 @@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,13 @@ per-topic record-count vector to each segment, including empty segments. Replay
capacity and latency policy, validates the recorded logical manifests, and performs the live number
of event-loop drains once across both directions.

An input schedule also retains activations with no input observation. If the first subscription
is created late in a task, earlier empty activations in that task precede its first observed
segment. This keeps the two schedules aligned when an Activity, timer, or another stream causes
an intervening activation. A prerelease combined marker with different input and output segment
counts is rejected: the positions of omitted empty input segments cannot generally be inferred
from those counts, and dropping output segments would weaken replay verification.

When another publish would exceed the record or logical-byte limit, `publish()` waits, the current
batch is staged with an output-capacity terminal, and Core forces a replacement Workflow Task. A
single oversized record or manifest is rejected before unsafe external I/O or marker growth.
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 @@ -820,6 +820,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 @@ -29,13 +29,15 @@ const SERDE_DERIVE_PREFIXES: &[&str] = &[
".temporal.api.namespace",
".temporal.api.nexus",
".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 @@ -136,24 +136,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
3 changes: 2 additions & 1 deletion crates/protos/protos/api_upstream/nexus/workflow-service.wit
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,8 @@ interface workflow-service {
started: option<bool>,
/// @nexus.omit
signal-link: placeholder,
first-execution-run-id: string,
/// @nexus.omit
first-execution-run-id: placeholder,
}

/// @nexus.doc
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,59 @@ 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 {
// 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;
}

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

enum StartChildWorkflowExecutionFailedCause {
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,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 +966,53 @@ 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 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;
// 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 +1336,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