Skip to content

Added the workflow-side stream reader and writer. - #12

Closed
moedash wants to merge 35 commits into
moe/AI-198-py-02-streams-packagefrom
moe/AI-198-py-03-workflow-runtime
Closed

moedash wants to merge 35 commits into
moe/AI-198-py-02-streams-packagefrom
moe/AI-198-py-03-workflow-runtime

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR adds workflow.stream_reader, workflow.stream_writer, workflow.subscribe_channel and workflow.linked_channel, carries the provider into the workflow runtime, and adds the notification channel calls, the channel subscriptions on describe and stream_channel to the client.

What changed?

  • temporalio/workflow/_streams.py has StreamReader and StreamWriter. A reader is an async iterator over its topic, and it takes after=END or last=N. A writer's publish returns at once, because the Workflow Task is the visibility boundary.
  • With no topic, both calls address DEFAULT_TOPIC. A topic definition is resolved to its name before the provider sees it.
  • temporalio/workflow/_channels.py has subscribe_channel, linked_channel, ChannelSubscription and Notification. The first subscribe_channel call for a channel in a run issues the SubscribeNotificationChannel command; a second call shares the subscription. linked_channel opens 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_to carries the owner of a linked channel as a temporalio.common.Execution, a workflow or a standalone activity with the run that received it, and is None for an independent one. Execution and ExecutionType are new in temporalio.common, with Execution.workflow(...) and Execution.activity(...) as the two spellings. ChannelSubscription.linked says which kind a handle is on.
  • ChannelSubscription.unsubscribe() ends an independent subscription. It records the UnsubscribeNotificationChannel command once, a second call is a no-op, and the handle is closed: what it had queued can still be read, then iteration ends and receive() raises RuntimeError. A later subscribe_channel on the same name opens a new subscription with a new command. On a linked handle it raises ValueError, since a linked channel is part of the run and has no subscription to end.
  • WorkflowExecutionDescription.channel_subscriptions lists the channels a run stands on as ChannelSubscriptionInfo: 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 a StreamRef, as a ChannelAddress of name and owning Execution: 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, and stream/<stream id> independent for a standalone stream. ChannelAddress.workflow_id reads 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.
  • _WorkflowInstanceImpl handles the NotificationsReceived activation job by routing each notification to the handle of its kind on its channel: linked_to set 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_listener and unregister_channel_listener go through the outbound interceptor with one input each. Each takes execution to address a linked channel instead of the independent one of that name, with workflow_id and run_id as 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.
  • ChannelDescription gains kind, a ChannelKind enum mirroring the proto's, and linked_to, an Execution. A linked channel of a running execution exists by construction, so a describe with execution answers the linked kind with nothing retained for a name nobody has notified yet, where the independent kind answers not found.
  • _Runtime gains workflow_streams(), workflow_subscribe_channel(), workflow_linked_channel() and workflow_is_evicting(), and _WorkflowInstanceImpl implements them. The workflow half of the provider is made per instance, so its state dies with the instance.
  • _StreamHooksInterceptor in temporalio/worker/_workflow.py goes innermost when the worker has a provider. It brackets the workflow function with the provider's start and finish hooks.
  • Worker takes the provider from its own option or from the client's. Replayer passes 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_to picks the handle rather than the name alone.

How did you test it?

I ran uv run poe lint and uv run pytest tests/streams. test_streams_workflow.py runs the handles in real workflows on the memory provider. test_stream_reader.py covers the reader's buffer, cancellation and reopen rules. test_stream_hooks.py drives 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.py drives 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, a NotificationsReceived job resolves receive(), 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 carry linked_to while a same-named subscription gets the rest. A client case against the dev server replaces the service calls and checks that each request carries execution only when an owner was given, by execution or by the workflow_id shorthand, that the two together are refused, and that describe maps kind and linked_to. A linked live case under needs_execution_server describes, notifies and polls by execution and checks the typed owner that comes back, and waits for the server layer that addresses a linked channel by execution. The live cases carry the needs_channel_server marker and skip unless -E host:port names a server that serves channels; they passed against the server from the moedash/temporal chain. The linked live cases carry needs_linked_server and 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 with linked_to set 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 carry needs_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 onto ChannelSubscriptionInfo, and another checks the names stream_channel derives for each owner kind. The live describe cases carry needs_describe_server and 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 carries needs_unsubscribe_server and needs_unsubscribe_core and 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.

  • Unit Tests
  • Staging
  • End to End Tests

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.
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.
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.
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.
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.
moedash added 12 commits October 1, 2026 19:33
… 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.
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.
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.
@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

skip-changelog Changelog entry rides another PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant