Skip to content

Added the Workflow Streams provider over the shipped wire format. - #4

Closed
moedash wants to merge 53 commits into
moe/AI-198-py-04-accessorsfrom
moe/AI-198-streams-provider-workflow-streams
Closed

moedash wants to merge 53 commits into
moe/AI-198-py-04-accessorsfrom
moe/AI-198-streams-provider-workflow-streams

Conversation

@moedash

@moedash moedash commented Sep 15, 2026 •

Copy link
Copy Markdown
Owner

This PR adds WorkflowStreamsProvider, the stream interface over the shipped Workflow Streams wire format.

What changed?

  • temporalio/streams/providers/workflow_streams.py maps a topic onto the shipped topic of the same name. A record rides in the item payload as the StreamRecord proto. A Query serves a closed run's tail, paged and filtered by topic in the workflow.
  • An outside append is the __temporal_streams_publish Update. The workflow answers with the position, so append() returns a cursor. A divergent or stale repeat is refused as StreamProducerError. publish_transport="signal" picks the shipped Signal instead.
  • A cursor names its run, because the log doesn't cross continue-as-new. A handle without a run id reads the chain run by run and follows a reset into the run describe names.
  • A workflow's activity keeps its own streams as reserved activity/<id>/<name> topics in the workflow's log. A standalone activity and a standalone stream raise StreamUnsupportedError.
  • Bodies meet the codec and external storage at the transport envelope, so the worker's converter has to match the clients'.
  • temporalio/contrib/workflow_streams/ gains public accessors for code that rides its wire.
  • worker/_workflow.py and _workflow_instance.py run on_workflow_start on the workflow's event loop, so handlers register before the first task's Updates run. That change belongs on Added the workflow-side stream reader and writer. #12 and moves there in a later cleanup.

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

Why?

It's the on-ramp. The same wire format means current streams interoperate and old histories replay, so adopting temporalio.streams needs nothing deployed. The transport's limits stay, and one contract rule weakens: other readers see a producer's records only after the workflow observes them.

How did you test it?

I ran uv run poe lint and uv run pytest tests/streams. The conformance suite runs over this provider, with the standalone and byte-cap cases declared absent. test_workflow_streams_provider.py runs live against the test environment's server. It covers retry dedupe, supersession, the closed-run tail, continue-as-new and reset, activity streams, the Signal fallback, and external storage at the envelope.

  • Unit Tests
  • Staging
  • End to End Tests

@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch 3 times, most recently from 0ecca73 to 7b12920 Compare September 16, 2026 23:22
@moedash
moedash force-pushed the moe/AI-198-streams-interface branch from a92fe0a to 6b8b35e Compare September 17, 2026 02:51
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch 7 times, most recently from 61a91a2 to f4e1c40 Compare September 21, 2026 17:27
A provider that answers outside readers through handlers on the workflow has
to register them before the first task completes, and one that parks a
long-poll update against the run has to let go before the workflow returns.
Both are no-ops on the storage providers, so workflow code calls them
unconditionally and stays portable.
Porting the agent harness onto the interface surfaced three things the small
example never needed. A reader resumes after the record a cursor names, so
it can store the last cursor it handled without advancing an opaque token. A
consumer reports its latest position, so a client can follow a turn it is
about to start without the workflow reporting a position. And a producer can
append onto a topic of the stream the workflow publishes, so an activity's
live output lands next to the workflow's own records for one outside reader.
…ort.

The poll Update stops answering once the workflow is closing, so a reader
between polls at that moment lost whatever the final task published, and
nothing could read the stream after completion at all. The log is workflow
state, so one Query serves it for as long as the History is retained; the
consumer asks for the tail when its subscription ends.
The transport already resolves its payload converter and falls back to
the default when it holds no client, so the provider asks it rather than
reaching through a client that may not be there.
The changelog checkpoint requires an entry for any user-facing change.
A provider that reuses the shipped transport needs the publish signal's name, the log past an offset, and the client's handle and converter. Reading them off private attributes breaks silently on the next contrib change.
The SDK rebuilds an evicted workflow from history as a new object. A process map keyed by run id handed that object the previous instance's log, with no handlers registered on the new one and the records of a failed task still inside. The stream is now found through the handler the shipped class registers on the instance, and the provider reads the contrib module only through its public surface.
A retry after an ambiguous failure used to carry fresh sequences, so the shipped dedupe let the same batch land twice. The batch now stays pending under its signal sequence until the server accepts it, and goes out first if the caller moves on. append returns None because this transport learns positions at read time.
Inbound frames carry no topic and a topic filter on an inbound stream is rejected through the shared check; producer identity comes resolved from the package; an unreadable frame is skipped with a warning; read is declared as the generator it is.
…topic.

Only a missing handler or a vanished history means there is no tail to serve; anything else, such as a query timeout with no worker polling, is now an error rather than a silent empty stream. An owner-stream read with a topic subscribes to that shipped topic alone instead of dropping inbound items client-side.
…ver.

The tests take the shared client fixture instead of an opt-in live gate, the loop workflow calls the portable lifecycle hooks, and new cases cover a cold workflow cache, the tail of a closed run, and the resend of a batch whose signal failed.
Topics are the shipped topics of the same name, the record is the
StreamRecord proto inside the item payload, and cursors name the run
because a log is not carried across continue-as-new.
The provider registers in SETUPS with a host workflow that owns the
stream, and its own tests cover the tail Query, a read that ends, and
a handle that follows continue-as-new run by run.
…one.

The workflow half now keeps the WorkflowStream it found at start, so the
finish hook lets go of this run's pollers rather than of whatever the
signal handler lookup answers on the thread at the time.
The shipped subscribe loop swallows a cancellation of the caller's task,
so a bounded read resubscribed forever instead of timing out. The poll
name and its failure types are public on the contrib package for this.
An Update in a run's first task runs ahead of the start hook, so a live read hit it on every successor after continue-as-new. A second rejection on a running run now names the missing provider instead.
… moe/AI-198-streams-provider-workflow-streams
The provider serves a read that starts at END or at the last N records,
defaults the topic when a call names none, and refuses a stream an
activity owns, since its log lives inside a running workflow. The
conformance setup gained the truncate hook and the worker the shared
cases expect.
An activity a workflow scheduled keeps its streams under
activity/<escaped id>/<name> in the workflow's log, the way native reserves
that prefix, so activity.stream_handle(scope="activity") and the client's
activity handle work here. A standalone activity stays refused because no
workflow hosts its log, and the read ends with the activity as the workflow
describes it.
The base run closes with nothing in its own History about the reset, so a
chain-following read asks describe for the run reset into it. That run
rebuilt the base run's log up to the reset point by replay, so the items
before it sit at the same offsets and the read carries on where it was.
A Signal has no response, so the workflow could only drop a divergent
retry. The publish Update answers with the batch's position, so append()
returns a cursor, and refuses a conflicting repeat with a typed error
before accepting it. A producer falls back to the Signal on a workflow
whose worker predates the Update; the Action cost is one per batch either way.
… moe/AI-198-streams-provider-workflow-streams
…treams.

Handles name their stream as a StreamRef and refuse close(), standalone
streams are refused with a housing workflow noted as the option, and each
record carries the plaintext hash of its body, which the workflow now
matches a repeated batch by. Bodies meet the codec and external storage at
the transport's envelope, since the workflow thread reads them from state.
… moe/AI-198-streams-provider-workflow-streams

# Conflicts:
#	tests/streams/conftest.py
#	tests/streams/test_streams_conformance.py
… moe/AI-198-streams-provider-workflow-streams
…reams.

The provider hosts no standalone streams, so neither retention policy
exists here and the retention case skips with the rest of them.
The SDK handles a task's Signals and Updates ahead of the workflow
function, so a publish or poll Update that arrived with the first task was
rejected and, on a worker with no cache, on every rebuilt instance too. The
start hook now runs when the instance's loop first turns, and Workflow
Streams binds the shipped stream object lazily so a workflow's own stays.
… moe/AI-198-streams-provider-workflow-streams
The provider hosts no standalone streams, so there is no byte cap to
refuse an append past; the flag says so next to the other three.
@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