Conversation
Workflow code reads and writes a stream through the provider's workflow half, made once per instance so its state dies with the instance.
The start hook runs on the workflow's loop before the first task's handlers, so a provider's handlers exist for an Update that arrives with that task. The finish hook skips eviction.
The hooks, the reader's start positions and a workflow reading and writing through the memory provider.
This was referenced Oct 3, 2026
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 lets workflow code read and write a stream through whatever provider the worker has.
What changed?
workflow.stream_reader()andworkflow.stream_writer(), withStreamReaderandStreamWriter. They default to the default topic, and a reader can start atENDor at the last N records.__init__and before the first task's Signals and Updates are handled. A provider that serves outside readers through handlers needs them in place for an Update that arrives with the first task. The finish hook runs on return, failure, cancellation and continue-as-new, but not on eviction.Part of AI-198 (epic AI-37).
Why?
A workflow shouldn't have to know which store holds its stream. With the provider on the worker, the same workflow code runs on memory, Workflow Streams or a server-side stream. The hook placement used to live in the Workflow Streams provider. It's runtime behaviour every provider relies on, so it moves here, where both Python chains share it.
How did you test it?
Link to a test plan if any -
poe lintis clean. The hook, reader and workflow stream tests pass against the memory provider on the dev server the fixtures start. The workflow, replayer and interceptor suites still pass.