Skip to content

Woke parked external readers through the stream's channel. - #42

Open
moedash wants to merge 2 commits into
moe/AI-198-if-pyext-5-core-wake-pinfrom
moe/AI-198-if-pyext-6-channel-wake
Open

moedash wants to merge 2 commits into
moe/AI-198-if-pyext-5-core-wake-pinfrom
moe/AI-198-if-pyext-6-channel-wake

Conversation

@moedash

@moedash moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner

This PR makes notification channels the external stream runtime's wake transport, with the reserved Signal as the fallback.

What changed?

  • _wake.py gains WakeTransport (auto, channel, signal), channel_for, send_wake, the three-way server probe, WakeRequest.channel and channel_execution, and a per-client memory that steps down to the Signal after the server answers UNIMPLEMENTED. ChannelAddress comes from temporalio.client, so there's one copy.
  • Producers and the Worker's sweep send wakes through send_wake with a position and a counter from the backend's order. Backends get wake_transport, wake_counter_for and wake_counter_now.
  • A subscription asks the run to listen on its stream's channel through subscribe_stream_channel. From then on every completion carries a WorkflowStreamChannels report of the open readers' channels, and Core subscribes and unsubscribes from it.
  • The Worker probes the server once, on a running workflow, and records the answer as lang flags 100 and 101. A replay takes the path the live run took.
  • The runtime notes each channel's latest notification. Core resumes the parked waits itself.
  • Tests: test_wake.py, test_channel_report.py, the flag cases in test_channels.py, two live wake cases, and Signal pins in the cases that count Signals.

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

Why?

A reader parked on an external stream needs a wake when someone appends. The reserved Signal costs a History event per wake, and it only reaches a run that's already parked. A channel notification rides the scheduled event of the task it causes, folds repeats by counter, and reaches whatever run is current. Lang flags keep the choice deterministic, so a History written on a channel server replays on any Worker. The subscribe goes out in a report rather than as a command because a completion carrying a server-bound command can't be retained, and the task that opens a reader is exactly the one that should stay open and park.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

poe lint is clean, and the external stream suite passes on the dev server at each commit. Against a channel server, the channel, report, wake and integration cases pass. The live wake through the linked channel passed with the linked kind on. The wake through the independent channel passed with it off.

A wake now notifies the stream's channel and steps down to the reserved Signal
on a server without channels. The run reports the channels its open readers
listen on with every completion, so Core keeps the server's subscriptions in
step, and two lang flags record what the server offered for replay.
The unit cases cover the transports, the step-down memory, the three-way probe,
the report and the flags. The live cases wake a reader through the independent
and the linked channel against a server named with -E.
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.

1 participant