Skip to content

Vendored the api's stream protos and the stream delivery job. - #5

Closed
moedash wants to merge 39 commits into
mainfrom
moe/AI-198-core-1-protos
Closed

moedash wants to merge 39 commits into
mainfrom
moe/AI-198-core-1-protos

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR vendors the api's stream and notification channel protos into Core and adds the stream delivery job to the bridge protos.

What changed?

  • temporal/api/stream/v1/message.proto and the api files that reference it: the two workflow commands and their events, the enum values, the failed cause, StreamStartPosition, and the stream slices on a poll response. The api branch's diff is applied onto the tree this Core pins.
  • The protos crate wires in temporal.api.stream.v1. crates/common/build.rs exempts StreamRecord.body and StreamRecord.metadata from the blob-size check. A record body never reaches an event, and the server bounds a batch by message count.
  • DeliverStreamRecords on the activation job at field 21. The external family takes 17 to 20, so the number is fixed across both trees.
  • The Rust SDKs have no stream API, so workflow_future.rs and runtime/instance.rs fail on that job. The server already recorded the range as consumed, so dropping it loses data.
  • The integ matrix reads its per-runner timeout in per-pr.yml.
  • The notification channel protos from the api branch: temporal.api.notification.v1, SubscribeNotificationChannelCommandAttributes and its event and failed cause, WorkflowTaskScheduledEventAttributes.notifications, and the channel calls on the raw client with their C bridge entries, plus the linked channel fields: Notification.linked_to, ChannelKind, execution on the channel requests, and kind and linked_to on describe. A linked channel's owner is a temporal.api.common.v1.Execution with a type, a business id and a run id, so an activity can own one as a workflow does. The lang protos carry SubscribeNotificationChannel at 31 and NotificationsReceived at 22. Until Let a workflow subscribe to a stream and append records to it. #6 maps the command, Core refuses it. The Rust SDKs ignore the job, since they cannot subscribe and a notification carries no data.
  • ChannelSubscriptionInfo and DescribeWorkflowExecutionResponse.channel_subscriptions, so a client built on this tree decodes the channels a run stands on. No new RPC.
  • UnsubscribeNotificationChannelCommandAttributes with its event and failed cause, and UnsubscribeNotificationChannel at 32 on the lang command message. Core refuses it here the way it refuses the subscribe, until Delivered server-side stream ranges to workflows. #7 adds the machine.

This is the proto layer, with no stream behavior yet. #6 adds the command machines.

Part of AI-198 (epic AI-37).

Why?

Core can't carry a workflow stream before it has the messages. Pinning the protos here first means every later layer generates the same temporalio/api. The shared field numbers let an artifact generated from the external tree decode correctly against this one. The channel calls sit on the raw client with C bridge entries so the Python bridge can reach them.

How did you test it?

cargo check --workspace, cargo fmt --all --check and cargo test-lint are clean. This layer has no unit tests, since the stream tests arrive with #6.

  • Unit Tests
  • Staging
  • End to End Tests

tconley1428 and others added 6 commits August 28, 2026 11:42
This branch is the Core commit upstream sdk-python pins, plus the api
branch's stream diff applied onto that Core's own api tree, so a Python
branch that vendors the stream protos can regenerate them from its pin.
The protos crate lists the new package and the two payload fields;
nothing else changes.
The nexgen the Python SDK drives rejects @nexus.type on a native declaration, so generation failed for every consumer. Upstream declares both policies as placeholders, which keeps the annotations valid.
The matrix asks for 40 minutes on `macos-intel`, but the job never read that value, so the runner stayed on the 25-minute default it exceeds.
The api adds a stream family: the record, the two workflow commands with their events, and the slices a poll response carries. Core needs those messages before anything can use them. The record bodies are exempt from the blob-size validation because they never reach an event; the server bounds a batch by message count instead.
History records the offsets a task consumed and never the payloads, so the server sends the bytes on the poll response and Core hands them to lang as their own job. The Rust SDKs have no stream API, so they fail loudly on this job rather than ignoring it: the server has already recorded the range as consumed and will not send it again, so dropping it would lose data silently.
The Python branch that vendors these regenerates from this pin, so the pin has to hold the same field names the api branch now publishes.
The append event carries the exclusive end offset instead of a count, so
the three range-carrying messages read the same way. The command fields
say name rather than id, which is what a Workflow actually addresses.
# Conflicts:
#	crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto
@moedash moedash added the skip-changelog Changelog entry rides another PR label Sep 26, 2026
The api adds a StreamStartPosition message and a start_position field on the subscribe command, at a new field number, so a Workflow can ask for the earliest record, the tail or the last N.
The vendored tree sits on an older upstream base, so this applies the api commit's wake diff rather than copying whole files.
# Conflicts:
#	crates/protos/protos/api_upstream/nexus/workflow-service.wit
#	crates/protos/protos/api_upstream/temporal/api/enums/v1/failed_cause.proto
#	crates/protos/protos/api_upstream/temporal/api/workflowservice/v1/request_response.proto
#	crates/protos/src/protos/mod.rs
…s api tree.

Applies the api commit's diff to the vendored tree, exposes the channel calls on the raw client and the C bridge, and classifies Notification.metadata for payload limits. The C bridge also gains the wake call it was missing.
# Conflicts:
#	crates/protos/src/protos/mod.rs
moedash added 14 commits October 1, 2026 18:07
…7' into moe/AI-198-core-1-protos

# Conflicts:
#	crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto
#	crates/protos/src/protos/mod.rs
…7' into moe/AI-198-core-1-protos

# Conflicts:
#	crates/sdk/src/workflow_future.rs
#	crates/workflow/src/runtime/instance.rs
Applies api 9cd8b40 to the vendored tree and removes the raw-client proxier and the C-bridge dispatch entry that the Python bridge generator built on.
Applies api a9e6517 to the vendored tree. The new fields carry no payloads, so the payload-limits table in `crates/common/build.rs` needed no entry.
Applies api 30918a0 to the vendored tree. The new message nests a Notification, whose metadata the payload-limits table already exempts, so the table needed no entry.
Applies api c239d35 to the vendored tree, with the lang command at 32 next to the subscribe. The protos pin refuses the command the way it refuses the subscribe, so a workflow cannot go on as if unsubscribed while the server keeps delivering.
Carries api c239d35..4304fd8: the five channel requests, the notification
and the describe response name the linked owner as an Execution, so an
activity can own a channel too.
Carries api 4304fd8..071feb8. The earlier shape was never released, so
the Execution fields take the numbers the WorkflowExecution fields had
and nothing is reserved.
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #8, #9, #10, #11, #12, #13, #14, #15, #16, #17, #18, #19, #20, #21, #22 and #23.

Same content, split into 16 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

@moedash moedash closed this Oct 3, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

skip-changelog Changelog entry rides another PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants