Conversation
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
3 times, most recently
from
September 16, 2026 23:22
0ecca73 to
7b12920
Compare
moedash
force-pushed
the
moe/AI-198-streams-interface
branch
from
September 17, 2026 02:51
a92fe0a to
6b8b35e
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
7 times, most recently
from
September 21, 2026 17:27
61a91a2 to
f4e1c40
Compare
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.
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
…vider-workflow-streams
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
WorkflowStreamsProvider, the stream interface over the shipped Workflow Streams wire format.What changed?
temporalio/streams/providers/workflow_streams.pymaps a topic onto the shipped topic of the same name. A record rides in the item payload as theStreamRecordproto. A Query serves a closed run's tail, paged and filtered by topic in the workflow.__temporal_streams_publishUpdate. The workflow answers with the position, soappend()returns a cursor. A divergent or stale repeat is refused asStreamProducerError.publish_transport="signal"picks the shipped Signal instead.activity/<id>/<name>topics in the workflow's log. A standalone activity and a standalone stream raiseStreamUnsupportedError.temporalio/contrib/workflow_streams/gains public accessors for code that rides its wire.worker/_workflow.pyand_workflow_instance.pyrunon_workflow_starton 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.streamsneeds 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 lintanduv run pytest tests/streams. The conformance suite runs over this provider, with the standalone and byte-cap cases declared absent.test_workflow_streams_provider.pyruns 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.