Skip to content

Carried the external stream task and replay wake fixes. - #2

Closed
moedash wants to merge 35 commits into
task/python-sdk-streamingfrom
moe/AI-198-external-core-repairs
Closed

moedash wants to merge 35 commits into
task/python-sdk-streamingfrom
moe/AI-198-external-core-repairs

Conversation

@moedash

@moedash moedash commented Sep 14, 2026 •

Copy link
Copy Markdown
Owner

This PR fixes the external stream task lifetime and replay wakes, and resumes parked external stream waits from channel notifications.

What changed?

  • Cache zero: a retained external stream task stays until its normal boundary, as a local activity does. See retains_task_for_external_streams in managed_run.rs and the eviction check in workflow_stream.rs.
  • Replay wakes: a wake Signal that replay reaches inside a History page resumes the reconstructed waits before the task reports empty.
  • Output replacement: a stale wait from an earlier task doesn't force an empty replacement task. A completion that stages a commit with output still buffered does force one. The one-outstanding-activation rule logs through dbg_panic! in release builds too.
  • Notification receipt: the channel notifications on a task's scheduled event reach lang as one NotificationsReceived job, folded per channel with the highest counter kept, and resume parked waits the way an unparked wake Signal does. The Signal stays as the fallback transport.
  • Deferred subscribe: lang reports the channels a run listens on with WorkflowStreamChannels, and Core issues SubscribeNotificationChannel for a new channel on the completion that ends the Workflow Task, after the marker and never on a run-ending one, with the same deferral on replay. The task that opens the first reader stays retained and parks as before.
  • Unsubscribe: the UnsubscribeNotificationChannel command variant and its machine, matched against WorkflowNotificationChannelUnsubscribed on replay the way the subscribe is. Core issues it for a channel that left the reported set, ahead of any subscribe on the same completion, so a run rotating channels frees the slot first.
  • Protos and client: the stream and notification protos land in the vendored api tree, including the linked channel contract (Notification.linked_to, ChannelKind, the optional execution on the channel requests, and kind and linked_to on DescribeChannelResponse, all naming the owner as a temporal.api.common.v1.Execution by type and business id), the unsubscribe command, event and failed cause, and ChannelSubscriptionInfo with DescribeWorkflowExecutionResponse.channel_subscriptions. The channel calls (NotifyChannel, RegisterChannelListener, UnregisterChannelListener, PollChannel, DescribeChannel) are on the raw client in crates/client/src/grpc.rs and in the C bridge dispatch.
  • Backports: the nexus request-timeout header fix (upstream Tolerate out-of-grammar nexus request-timeout headers temporalio/sdk-rust#1497, authorship kept), two integ test fixes, the per-runner integ timeout, and two WIT declarations matched to upstream.

This layer sits on Max's task/python-sdk-streaming. #3 merges it forward onto current main.

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

Why?

At cache zero the worker evicted the retained task, then spun on empty replacement tasks while unread records sat in the store. A channel notification rides the scheduled event of the task it produces, so a provider wakes a parked workflow without a Signal event. Replay doesn't change, because the read lives in the marker and the notification only decides when the live run looks. The WIT and header fixes keep nexgen codegen and nexus_metrics working against the dev server the Python SDK runs.

How did you test it?

The full cargo suite and cargo test-lint ran on the branch. New tests in core_tests/external_streams.rs cover zero-cache retention, buffered output, the stale wait, a wake Signal reached mid-replay, and the notified scheduled event: it yields the job and resumes a parked wait, a failed task hands its notifications to the retry, and replay yields the same job and resumes the reconstructed waits. For the deferred subscribe: a retained task holds the command until the park or finalization that ends it, a command-carrying completion orders it after the marker, a run-ending one issues none, a later task adds nothing for the same channel, replay reissues it from the same report, and a different channel on replay is nondeterminism. For the unsubscribe: a channel that left the set is unsubscribed on the completion that ends the task, a swap unsubscribes before it subscribes, a run-ending completion issues none, replay reissues it from the same report, and a replay still naming an unsubscribed channel is nondeterminism, plus the machine's own unit and replay cases. moedash/sdk-python#15 pins this head and runs the live Redis provider suite on it. The C bridge entries have no Rust test.

  • Unit Tests
  • Staging
  • End to End Tests

moedash and others added 19 commits September 16, 2026 14:41
`clippy::type_complexity` rejects the inline nested type, and the lint
is denied through `-D warnings` in `cargo test-lint`.
The reporter panics on any machine name it has no visualizer for, so
every state machine has to be listed there.
Backported from upstream `42c7bfe9`. This branch forked before that fix,
and the unscoped query also matches the shared-namespace worker's series,
whose order is not the test's to choose.
Backported from upstream `c87e2060`. The dev server needs
`system.enableCancelActivityWorkerCommand`, and the eager case needs a
server newer than the one the pinned CLI bundles.
The `string` declaration made the generated Python Nexus model require
`SignalWithStartWorkflowResponse.first_execution_run_id`, which the dev server
bundled with this Core never sends. Upstream omits the field instead.
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.
Backported from upstream `84be3c1f`. This branch forked before that fix, and
the dev server bundled with the current Core writes the header with Go's
duration syntax, which the old parser rejects, so `nexus_metrics` hangs.
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 zero-cache tests now poll the eviction at the boundary and expect the run
to be gone, for a quiescent task and for buffered output. The wake test loses
its wall-clock bound and gains a paginated variant whose wake is decoded while
replay is still in progress.
The caller now says whether it is completing the outstanding activation, and a
debug assertion checks that against the run's state, so a new call site cannot
break the one-outstanding-activation rule silently.
`cargo test-lint` refuses a `MutexGuard` held across an await point, and
the assertions only need the guard for two lines.
Regenerating `temporalio/api` from this Core dropped the stream messages,
commands and events the Python branch vendors, so its proto check failed.
The api branch's diff is applied onto the api tree this Core pins, and the
protos crate lists the new package and the two payload fields.
Narrowing the replacement condition to a current wait took buffered
output with it. The flush deadline does nothing once the task is gone, so
the publish latency lang asked for was dropped without a word.
A breach reorders activations rather than crashing, which a debug-only
assertion leaves silent in exactly the builds where it would bite.
Applied from api 75de7aa. The poll response schema in the vendored OpenAPI files lacks some fields the api fork has, so its wakes hunk was placed by hand.
A wake resumes waits without an event in History, so Core takes it from the poll response and treats it as an unparked wake. It is held until replay reaches the live task.
The wake copy carries only the three proto files, matching the main-lineage copy, so the later merges of both lineages agree.
Applied from api 086d85d. Comment-only.
…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 command gets its own machine, matched against the subscribed event on replay. The
notifications on a scheduled event reach lang as one job and resume parked waits, and such a
task is no longer folded into the one before it on replay.
The server clears what it puts on a scheduled event, so the retry's event carries only what
arrived since. The notifications are joined in History order, and the lang protos read as the
protos branch does.
…receipt.

The notification channel replaces the wake in the PoC, so the Core carries
neither the api messages nor the poll-response receipt. The reserved wake
Signal and the channel notifications keep resuming parked external waits.
A linked channel lives in one workflow's state and is addressed by
workflow id, so the channel calls take an optional workflow_execution,
a notification names the channel it is linked to, and describe reports
the kind and the owner.
The api tree now adds fields to Notification, so the test helper takes
what it does not set from Default instead of naming every field.
The server-bound subscribe command ended the task that opened the first reader, so that
task could never be retained or parked. Lang now reports the channels it listens on with
`WorkflowStreamChannels`, and Core issues the subscribe for a new channel after the marker
on the completion that ends the task, never a run-ending one, with the same deferral on
replay so the recorded event matches.
The command, its event and failed cause, and the workflow's channel subscriptions on
DescribeWorkflowExecution, as the api fork defines them.
A channel that leaves the reported set goes out as `UnsubscribeNotificationChannel` on the
completion that ends the task, ahead of the subscribes, and is matched against its event on
replay. The lang-issued command gets the same machine.
An activity can own a linked channel too, so the channel calls and the
notification's owner take a temporal.api.common.v1.Execution with a type
and a business id, as the api fork defines them.
moedash added a commit that referenced this pull request Oct 3, 2026
Carries the vendored api change c239d35..4304fd8: the five channel requests take a common.v1.Execution and both linked_to fields become Execution. Core-3 7a0d774 applies the same patch, so this merge brings #3 to the main lineage's head without merging core-3 itself.
moedash added a commit that referenced this pull request Oct 3, 2026
Carries #2 7006981 through #3 093c55e: the five channel requests take a common.v1.Execution and both linked_to fields become Execution. The native stream protos are untouched.
The workflow-addressed shape was never released, so the execution fields
take its numbers and nothing is reserved, as the api fork defines them.
moedash added a commit that referenced this pull request Oct 3, 2026
Carries the vendored api change 4304fd8..071feb8: the five channel requests' execution fields and both linked_to fields reuse the old numbers and nothing is reserved. Core-3 fbc064b applies the same patch.
moedash added a commit that referenced this pull request Oct 3, 2026
Carries #2 677d468 through #3 2a516e1: the execution and linked_to fields reuse the old numbers and nothing is reserved. The native stream protos are untouched.
@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

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants