Conversation
Workflow code reaches its stream the way it reaches the payload converter, through the runtime on the thread. The workflow half of the provider is made per instance, so whatever it keeps dies with the instance the way handlers do.
The worker and the replayer hand the provider to each workflow instance, and the worker brackets the workflow function with the provider's lifecycle hooks so no workflow code has to call them. The finish hook stays off an evicted run and off a collected coroutine, because neither is the workflow ending.
This was referenced Sep 25, 2026
…-workflow-runtime
Reading the eviction state off a private attribute name failed open: rename it upstream and the finish hook starts running during eviction, with every test still green.
stream_writer() hands out a new object per call, so a guard kept on one of them let a publish land behind the finish marker it was meant to refuse.
A second loop on one reader, a cancelled read, a reopened topic and a workflow that touches no stream were all promised in a docstring and pinned by no case.
…-workflow-runtime
…-workflow-runtime
…-workflow-runtime
…-workflow-runtime
A workflow with one stream should not have to invent a topic name for it; with no topic both resolve to DEFAULT_TOPIC and decode as a string-named topic does.
…-workflow-runtime
A workflow could only start from BEGINNING or a cursor, so it had no way to follow from now. The provider resolves the start outside the workflow and records it, so replay reproduces it; last= is passed only when given, which keeps older providers working.
…-workflow-runtime
…-workflow-runtime
…-workflow-runtime
…-workflow-runtime
…-workflow-runtime
…-workflow-runtime
A workflow subscribes to a notification channel once per run and gets the notifications the server folded for it with each Workflow Task. They travel in History, so a replay sees the same ones at the same points.
Notify, poll, describe, register and unregister go through the outbound interceptor with one input each. Metadata values are encoded with the data converter on the way out and come back as payloads, the shape workflow code sees.
A describe that races the subscribe is answered not found, and a worker whose Core refuses the subscribe command leaves the task in a timeout loop that the worker's shutdown waits on. The case now tolerates the first, names the second and bounds the shutdown.
…mand. The protos-only Core refuses the subscribe command, so the case is gated on a conftest constant the native layers flip with their pin. A client-only case covers what the server keeps for pollers: a notify with no listener is retained and answers zero, a counter at or below the latest is not kept, and a bounded poll above the latest comes back empty.
2 of 3 tasks
… case. A channel nobody has touched is not found, and the first notify on a fresh channel shows up in describe with its position and counter before any poll.
…-workflow-runtime
…-workflow-runtime
A subscription used to end only with the run. The handle now records the unsubscribe command once and closes, so a workflow can rotate channels under the per-run cap. The description lists what a run stands on, and stream_channel derives the channel a native stream notifies so a client can follow a stream without asking the server.
The owner's state takes a notification in the write that accepts it, so the counter is the run's at once, and with the first task already scheduled the notification waits behind it as the pending entry. A closed run keeps its listing.
…-workflow-runtime
A standalone activity can own a channel too, so `temporalio.common.Execution` names the owner of a linked channel, the five channel calls take it, and `workflow_id` with `run_id` stays as the shorthand for a workflow owner. `stream_channel` names a standalone activity's channel as `stream/<topic>` linked to the activity execution.
The interceptor inputs name the owner once, as `execution`, and the client resolves the `workflow_id` and `run_id` shorthand before building them, so an interceptor sees one spelling.
…-workflow-runtime
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 adds
workflow.stream_reader,workflow.stream_writer,workflow.subscribe_channelandworkflow.linked_channel, carries the provider into the workflow runtime, and adds the notification channel calls, the channel subscriptions on describe andstream_channelto the client.What changed?
temporalio/workflow/_streams.pyhasStreamReaderandStreamWriter. A reader is an async iterator over its topic, and it takesafter=ENDorlast=N. A writer'spublishreturns at once, because the Workflow Task is the visibility boundary.DEFAULT_TOPIC. A topic definition is resolved to its name before the provider sees it.temporalio/workflow/_channels.pyhassubscribe_channel,linked_channel,ChannelSubscriptionandNotification. The firstsubscribe_channelcall for a channel in a run issues theSubscribeNotificationChannelcommand; a second call shares the subscription.linked_channelopens the channel linked to this workflow and issues nothing, since the owner is the listener by construction.receive()and async iteration wait on a future the delivery resolves, so no command or timer is involved.Notification.linked_tocarries the owner of a linked channel as atemporalio.common.Execution, a workflow or a standalone activity with the run that received it, and isNonefor an independent one.ExecutionandExecutionTypeare new intemporalio.common, withExecution.workflow(...)andExecution.activity(...)as the two spellings.ChannelSubscription.linkedsays which kind a handle is on.ChannelSubscription.unsubscribe()ends an independent subscription. It records theUnsubscribeNotificationChannelcommand once, a second call is a no-op, and the handle isclosed: what it had queued can still be read, then iteration ends andreceive()raisesRuntimeError. A latersubscribe_channelon the same name opens a new subscription with a new command. On a linked handle it raisesValueError, since a linked channel is part of the run and has no subscription to end.WorkflowExecutionDescription.channel_subscriptionslists the channels a run stands on asChannelSubscriptionInfo: the kind, the subscribe event, the counter last accepted, the pending and scheduled notifications and, for the linked kind, the listener, retained and accepted counts.temporalio.client.stream_channel(ref)derives the channel a native stream notifies from aStreamRef, as aChannelAddressof name and owningExecution:stream/<topic>linked to the owning workflow,stream/<activity id>/<topic>linked to the workflow for a workflow activity's stream,stream/<topic>linked to the activity execution for a standalone activity's, andstream/<stream id>independent for a standalone stream.ChannelAddress.workflow_idreads the owner's id when it is a workflow. The server derives the same name, so a client polls or registers on it without asking._WorkflowInstanceImplhandles theNotificationsReceivedactivation job by routing each notification to the handle of its kind on its channel:linked_toset goes to the linked handle, unset to the subscription. A notification on a channel the run never asked for is dropped with a debug log.Client.notify_channel,poll_channel,describe_channel,register_channel_listenerandunregister_channel_listenergo through the outbound interceptor with one input each. Each takesexecutionto address a linked channel instead of the independent one of that name, withworkflow_idandrun_idas the shorthand for a workflow owner, and refuses both at once. The run id is optional and resolves to the chain's current run, as a Signal does. Notify and register set a request id. Metadata values are encoded with the data converter on the way out and come back as payloads a converter reads, the same shape workflow code sees.ChannelDescriptiongainskind, aChannelKindenum mirroring the proto's, andlinked_to, anExecution. A linked channel of a running execution exists by construction, so a describe withexecutionanswers the linked kind with nothing retained for a name nobody has notified yet, where the independent kind answers not found._Runtimegainsworkflow_streams(),workflow_subscribe_channel(),workflow_linked_channel()andworkflow_is_evicting(), and_WorkflowInstanceImplimplements them. The workflow half of the provider is made per instance, so its state dies with the instance._StreamHooksInterceptorintemporalio/worker/_workflow.pygoes innermost when the worker has a provider. It brackets the workflow function with the provider's start and finish hooks.Workertakes the provider from its own option or from the client's.Replayerpasses the one it was given.The client and activity accessors for streams come in #13.
Part of AI-198 (epic AI-37).
Why?
A provider that serves outside readers has to register its handlers before the first task completes. One that parks a poll has to let go before the workflow returns. Making that the worker's job keeps workflow code portable across providers. The finish hook skips eviction and collected coroutines, because at those points the runtime on the thread may belong to another run.
A channel separates the notification from the data. A writer outside Temporal says that a source moved; the listening workflow runs a task and reads the source from its own cursor. The notifications ride on the scheduled event of the Workflow Task, so History is the record and a replay sees the same ones at the same points, with no side channel to reconcile.
An unsubscribe lets a run rotate channels under the per-run cap instead of holding every name it ever listened on until it closes. The instance drops the handle with the command, so a notification the server folded onto a task before the command landed finds nothing and is dropped, and a replay matches the command to its event.
The linked kind is for the common case of one listener. It costs the owner one write per notification and no subscribe command or event, and a writer reaches it by execution, a workflow id or a standalone activity id, so a continue-as-new successor is reached by the same call. The independent kind stays for fan-out to many workflows and callbacks. A name may be open as both, which is why the notification's
linked_topicks the handle rather than the name alone.How did you test it?
I ran
uv run poe lintanduv run pytest tests/streams.test_streams_workflow.pyruns the handles in real workflows on the memory provider.test_stream_reader.pycovers the reader's buffer, cancellation and reopen rules.test_stream_hooks.pydrives the interceptor for eviction and a collected coroutine. Two cases are expected failures on memory: a publish committed with its task, and a read that replay re-supplies. A storage provider turns them into passes.test_channels.pydrives the workflow instance with activations the way Core does, since the dev server does not accept the subscribe command: the first subscription is a command and the second shares it, aNotificationsReceivedjob resolvesreceive(), notifications arrive in order, a foreign channel is dropped, an empty name is refused, a linked handle costs no command and gets only the notifications that carrylinked_towhile a same-named subscription gets the rest. A client case against the dev server replaces the service calls and checks that each request carriesexecutiononly when an owner was given, byexecutionor by theworkflow_idshorthand, that the two together are refused, and that describe mapskindandlinked_to. A linked live case underneeds_execution_serverdescribes, notifies and polls byexecutionand checks the typed owner that comes back, and waits for the server layer that addresses a linked channel by execution. The live cases carry theneeds_channel_servermarker and skip unless-E host:portnames a server that serves channels; they passed against the server from the moedash/temporal chain. The linked live cases carryneeds_linked_serverand passed against the server layer that keeps a channel in the owner's state: a client notifies a running workflow's linked channel by workflow id with the owner counted as the one listener, the workflow receives it withlinked_toset and no subscribe event in History, an untouched linked name of a running workflow describes as the linked kind while the independent name is not found, a wrong run id or an unknown or closed workflow is not found, and a linked channel is polled by workflow id or run id until the run closes. The two of them in which the workflow receives also carryneeds_linked_core, since the notifications reach workflow code as the job Core builds from the scheduled event and the protos-only Core this layer pins ignores that job; they run from the layer that pins the delivery Core. The unsubscribe is driven with activations too: the command goes out once with a second call a no-op, a notification queued at the unsubscribe is still read and then the iteration ends, a late job for the closed handle is dropped, a linked handle refuses, and a new subscription after an unsubscribe is a new command. A unit case maps every field of the describe response ontoChannelSubscriptionInfo, and another checks the namesstream_channelderives for each owner kind. The live describe cases carryneeds_describe_serverand passed against the server layer that lists the subscriptions: a subscribed run lists the independent kind with the subscribe event's id and a zero counter, the counter follows the task that carried a notification, a closed run keeps its listing, and a linked channel is listed once a notify by workflow id gives it state, with the notification pending behind the first task, the callback count following a register and an unregister, and the listing kept after the run ends. The unsubscribe live case carriesneeds_unsubscribe_serverandneeds_unsubscribe_coreand passed against the server layer that accepts the command: after the first notification the run records the unsubscribed event naming the subscribe event, the channel shows no workflow listener, the description lists nothing, a later notify wakes nobody, and the run finishes with its handle closed and nothing drained.