Conversation
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.
This was referenced Sep 25, 2026
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.
2 of 3 tasks
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.
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. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This PR repairs Max's external streams base, repins Core and wakes its readers through the notification channel.
What changed?
_WorkflowInstanceImpl, and the runtime is installed before the constructor runs. A publish or subscribe from@workflow.initworks.send_wakein_wake.pynotifies the stream's channel withNotifyChannel, which wakes the channel's listeners and writes nothing but the notification on the scheduled event of the task it causes. OnUNIMPLEMENTEDit steps down to the reserved Signal and remembers per client that the server has no channels.wake_transporton the backend isauto,channelorsignal, and both sides read it.channel_forformats 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 atemporal.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_forturns it into the counter, and a wake without a position takes the counter fromStreamBackend.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_FOUNDfrom the channel raises the sameRPCErrora Signal to an ended chain does.wait_idorder with repeats dropped, in aWorkflowStreamChannelscommand 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, withDescribeChannelon 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. Thenotifications_receivedjob only records the latest position per channel, and Core resumes the parked waits from it.workflow.subscribe_channel()andworkflow.linked_channel()return aChannelSubscriptionwithreceive()and async iteration overNotificationvalues, and the client hasnotify_channel,poll_channel,describe_channel,register_channel_listenerandunregister_channel_listener, each takingexecution=as atemporalio.common.Executionto address the channel linked to a workflow or a standalone activity, withworkflow_id=andrun_id=as the short form for a workflow, and withChannelDescription,ChannelKindandChannelListener. A notification names the execution its channel is linked to inlinked_to, and so does the description inkindandlinked_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.needs_unsubscribe_serverand runs only against a server named with-E host:portthat accepts it. The analysis is inspec/subscribe-retention.mdandspec/unsubscribe.md.temporalio/bridge/sdk-corepinsmoe/AI-198-external-core-repairsonmoedash/sdk-rust, which carries thetemporal.api.notification.v1package, the five channel calls, the subscribe and unsubscribe commands with their events, the workflow's channel subscriptions on describe, theNotificationsReceivedjob, the linked channel contract (Notification.linked_to,ChannelKind,executionon the channel requests and the kind and owner on the description, all astemporal.api.common.v1.Execution), and theWorkflowStreamChannelsreport Core turns into those commands. The vendored api and generated clients come from it.pydoctor.This layer has no
temporalio.streamscontent. #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-developanduv run poe lintare clean, anduv run poe gen-docsbuilds.tests/nexuspasses, andtests/contrib/external_workflow_streamsruns on a dev server withTEMPORAL_TEST_REDIS_URLset. 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_FOUNDand the positions the producer and the Worker report.test_channels.pydrives 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 markedneeds_channel_serverorneeds_linked_serverand run against a server named with-E host:portthat serves the kind.test_channel_report.pydrives 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.pysubscribes 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 withchannel.linkedKindEnabledoff, 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.