Skip to content

Repaired the external stream base and repinned Core. - #15

Closed
moedash wants to merge 34 commits into
task/python-sdk-streamingfrom
moe/AI-198-py-08-external-repairs
Closed

moedash wants to merge 34 commits into
task/python-sdk-streamingfrom
moe/AI-198-py-08-external-repairs

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR repairs Max's external streams base, repins Core and wakes its readers through the notification channel.

What changed?

  • The input and output replay schedules keep their empty activations. A replay marker applies before the task's first drain and doesn't earn a drain of its own.
  • The marker format changes. A history written by the unrepaired base fails its Workflow Task for good, rather than replaying a guessed schedule. There are no released histories behind this base.
  • A legacy query on a parked run is answered without the stream snapshot beside it.
  • Per-run state hangs off the SDK's _WorkflowInstanceImpl, and the runtime is installed before the constructor runs. A publish or subscribe from @workflow.init works.
  • send_wake in _wake.py notifies the stream's channel with NotifyChannel, which wakes the channel's listeners and writes nothing but the notification on the scheduled event of the task it causes. On UNIMPLEMENTED it steps down to the reserved Signal and remembers per client that the server has no channels. wake_transport on the backend is auto, channel or signal, and both sides read it.
  • channel_for formats the stream identity into the channel name, once for writers and readers, and returns the owner with it: a stream key names a workflow chain, so the channel is linked to that workflow. The notification carries the owner as a temporal.api.common.v1.Execution, a workflow or a standalone activity named by its business id, so a server with linked channels delivers to the owner's own channel and one without the kind ignores the owner and serves the independent channel of that name. A wake carries the store position, StreamBackend.wake_counter_for turns it into the counter, and a wake without a position takes the counter from StreamBackend.wake_counter_now, the same rule applied to the clock, so the shutdown sweep's wake is not folded under the channel's latest. The Signal's request id doubles as the notification's, so a retry deduplicates. NOT_FOUND from the channel raises the same RPCError a Signal to an ended chain does.
  • A run reports the complete set of channels its open readers listen on, one per reader in wait_id order with repeats dropped, in a WorkflowStreamChannels command on every completion once a reader has opened, the empty set included. The park result and the finalize answers carry none, and Core reads an absent report as unchanged. Core subscribes a channel new to the report on the completion that ends the task, after the progress marker and never on a run-ending one, and unsubscribes a channel that left it, so a reader closed before the task ends costs the server nothing. The worker asks the server once, with DescribeChannel on a probe channel linked to the first workflow it sees, and the answer becomes two lang flags a run reads on its first stream subscription and keeps until continue-as-new: with linked channels a stream the run owns is left out of the report, since the run is its listener by construction, with only independent channels it is reported, and without channels nothing is reported and the Signal wakes the run. A replay reports what the live run did, so Core matches the recorded events. The notifications_received job only records the latest position per channel, and Core resumes the parked waits from it.
  • The public channel surface has the main chain's names: workflow.subscribe_channel() and workflow.linked_channel() return a ChannelSubscription with receive() and async iteration over Notification values, and the client has notify_channel, poll_channel, describe_channel, register_channel_listener and unregister_channel_listener, each taking execution= as a temporalio.common.Execution to address the channel linked to a workflow or a standalone activity, with workflow_id= and run_id= as the short form for a workflow, and with ChannelDescription, ChannelKind and ChannelListener. A notification names the execution its channel is linked to in linked_to, and so does the description in kind and linked_to, so a run holding both kinds under one name routes by it. The public subscription and the stream runtime share one once-per-channel emission, so a workflow that reads a stream and also subscribes by hand pays one command.
  • The report asks nothing of the server, so the task that opens the first reader stays retained and parks as it does with the Signal, on every channel server. The eight cases built on a retained or parked first task run everywhere again. A case that needs the unsubscribe command carries needs_unsubscribe_server and runs only against a server named with -E host:port that accepts it. The analysis is in spec/subscribe-retention.md and spec/unsubscribe.md.
  • temporalio/bridge/sdk-core pins moe/AI-198-external-core-repairs on moedash/sdk-rust, which carries the temporal.api.notification.v1 package, the five channel calls, the subscribe and unsubscribe commands with their events, the workflow's channel subscriptions on describe, the NotificationsReceived job, the linked channel contract (Notification.linked_to, ChannelKind, execution on the channel requests and the kind and owner on the description, all as temporal.api.common.v1.Execution), and the WorkflowStreamChannels report Core turns into those commands. The vendored api and generated clients come from it.
  • CI repairs: main's Nexus and payload visitor generators, two upstream backports for the latest dependency set, and docstrings that build under pydoctor.

This layer has no temporalio.streams content. #16 adds the interface and #17 the Redis provider.

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

Why?

The base had drifted from main, so its generators failed and several suites broke on the latest dependencies, and the repairs land first so #16 and #17 read as interface changes. The constructor fix belongs here because the defect is in the base: a harness agent publishes from its constructor, and on Redis every Workflow Task of that run failed. The channel follows Max's review: the writer names the stream and never learns who listens, and the notifications ride History on the scheduled event, so a workflow may read them and replay sees the same. The linked kind follows his second comment: most streams have one listener, the workflow that owns them, and a channel kept in that workflow's state costs one write per notification and no subscription, so the task that opens the first reader stays retained as it does with the Signal. A burst of appends still costs one Workflow Task and no event per write.

How did you test it?

uv run poe build-develop and uv run poe lint are clean, and uv run poe gen-docs builds. tests/nexus passes, and tests/contrib/external_workflow_streams runs on a dev server with TEMPORAL_TEST_REDIS_URL set. Transport tests on a fake service cover the channel request and the owner it carries, the step down from the channel to the Signal and its per-client memory, the explicit transports, the two-way and the three-way probe, the clock rule, NOT_FOUND and the positions the producer and the Worker report. test_channels.py drives the workflow instance with activations for both channel kinds, for the stream runtime's gate under each probe answer and on replay, and for the client's request mapping. Its live cases are marked needs_channel_server or needs_linked_server and run against a server named with -E host:port that serves the kind. test_channel_report.py drives the instance with a real stream runtime: the report goes out on a completion with a snapshot and on one without, carries the empty set after the last reader closes, is absent before a reader opens and on the park and finalize answers, orders channels by reader without repeats, is left out on a server without channels and empty for the run's own streams on a linked one, goes out on replay, and the stream path emits no subscribe command of its own. test_worker_integration.py subscribes from the constructor and has a channel case per kind that checks the notified scheduled event and no Signal in History, with the subscribed event on the independent kind and none on the linked kind, then replays it. The modules that count the Signal pin their backends to it. The suite passes against the stock dev server, where the probe finds no channels. Against the server of moedash/temporal PR #17 it runs in both modes: linked by default, where every channel case runs and only the independent-only case skips, and independent with channel.linkedKindEnabled off, where the eight retained-first-task cases run on the independent kind and the linked-only cases skip. #16 adds a case that publishes from the constructor, and the channel runs live on #17.

  • Unit Tests
  • Staging
  • End to End Tests

An empty activation still ran one drain, so the marker has to record it or
replay fires wait_condition predicates a different number of times. The
marker also goes in with the first set that drains, rather than after the
signal and update jobs have already published.
A record only reaches a reader after a provider placed it, and the
subscription manager captures the loop it is built on, so the fakes have to
do both. The suite also runs against a dev server rather than the
time-skipping one, where which cases fail drifts between runs.
pydoctor resolves single-backtick prose as symbol references, so shell words,
tuples and file names were read as missing targets. Two tables also lost
their column alignment because pydoctor rewrites role references before
parsing.
The branch predates google-genai 2.21 file downloads and the test fixes that
came with the latest dependency set, and CI runs both. In addition, openai
now carries its own httpx, and the Agents SDK rejects two invocations sharing
one completed call id.
The branch's own Nexus generator shells out to a binary name the nexgen crate
does not install, and the payload visitor script still looked for the service
module under its old filename, so generation died before it could diff
anything.
The pin moves off the mfateev branch onto moedash/sdk-rust, whose head carries
the stream protos, so the vendored temporalio/api regenerates byte-for-byte
from it and the check-protos job passes. The script that stages the api
branch over the pin stays for the next time the two diverge.
A parked Run has no open Workflow Task, so the server dispatches a query on a
task of its own and Core allows nothing beside the answer there. The instance
still reported the registered wait set, Core refused the completion, and the
query timed out.
The id is Core's, and nothing in this file said where it is defined, so a
rename over there would be found by a failing query rather than by grep.
Behind a patch-only job set the install added a second drain, so every
wait_condition predicate fired once more than the recorded task fired it.
The dispatch moved to its own method so the schedule can be asserted.
The suite-wide skip took roughly five hundred offline cases out of the
time-skipping job. Only the cases that ask for a server observe its clock.
Two code-only fixes: a replacement task for output left buffered behind a commit, and the outstanding-activation breach logged in release builds.
The Core head carries the WakeWorkflowExecution call, so the vendored api and the bridge client are regenerated from it.
A wake records no History event and folds per stream, so a burst of appends costs one task. Servers without the call still get the reserved Signal.
Per-Run state lives on the SDK's Run object and the runtime is installed before the constructor runs, since the user's object does not exist yet while @workflow.init executes.
moedash added 11 commits October 1, 2026 18:17
The channel is the stream identity formatted once for writers and readers. Under auto the transport steps down on UNIMPLEMENTED and remembers per client where it stopped.
The subscribe command is gated behind a lang flag the Worker sets after asking the server once, so a stock server never sees it and a replay emits it where the live run did. The notifications job only records positions; Core resumes the parked waits from it.
The Signal-counting modules pin their backends to the Signal, since auto leaves none on a server with the channel or the wake call.
The same names as the main chain: workflow.subscribe_channel() with ChannelSubscription and Notification, and the five channel calls on the client. The external-stream runtime keeps its gated entry beside the public one, and both emit the command once per channel per run.
A plain clock sits far below a counter derived from a store position, so the server would fold the shutdown sweep's wake away. Each backend states its rule once, beside wake_counter_for.
The subscribe command ends the task that opens the first reader, so these cases measure nothing there until the command moves to the leaving completion. The skip fixture lives in the package conftest beside the channel-server marker.
A stream key names a workflow chain, so its channel is addressed to that workflow. A server with the linked kind keeps the channel in the owner's state and the owner listens by construction, so the task that opens the first reader needs no subscribe command and stays retained. The Worker's probe becomes three-way and is recorded as two lang flags.
The five calls take workflow_id and run_id, the description reports the kind and the owner, and a notification names the workflow it is linked to.
The channel module, the client dataclasses and the five calls take the main chain's text, so the union merges without a second copy. A linked channel's description lists the owner beside a registered callback once the channel holds state, so the live case looks for the callback among the listeners.
One test-only commit on top of the linked contract pin; the generated clients and the bridge are unchanged.
The api tree carries the unsubscribe command, its event and failed cause, and
the workflow's channel subscriptions on describe. The Core commands gain
`UnsubscribeNotificationChannel` and `WorkflowStreamChannels`, regenerated here.
The subscribe command is server-bound, so the task that opened the first
reader could never be retained or parked. The run reports the complete set of
channels its open readers listen on with `WorkflowStreamChannels` on every
completion, and Core subscribes and unsubscribes on the completion that ends
the task. The eight retained-first-task cases run on every channel server.
The five channel requests carry `execution` as a
`temporal.api.common.v1.Execution`, and both `linked_to` fields take the
same type. The generated api and clients follow the pin.
A standalone activity can own a channel too, so the client calls, the
stream wake path and the probe name the owner as a
`temporalio.common.Execution`, with `workflow_id` kept as the short form
for a workflow.
The execution fields take the numbers of the workflow-addressed shape,
which was never released, and nothing is reserved. Only the two
generated descriptors change.
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36, #37, #38, #39, #40, #41, #42, #43, #44, #45, #46, #47, #48, #49, #50, #51, #52, #53, #54 and #55.

Same content, split into 37 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.

1 participant