Skip to content

Added workflow.stream_reader and workflow.stream_writer. - #29

Open
moedash wants to merge 3 commits into
moe/AI-198-if-py-3-memory-providerfrom
moe/AI-198-if-py-4-workflow-runtime
Open

moedash wants to merge 3 commits into
moe/AI-198-if-py-3-memory-providerfrom
moe/AI-198-if-py-4-workflow-runtime

Conversation

@moedash

@moedash moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner

This PR lets workflow code read and write a stream through whatever provider the worker has.

What changed?

  • workflow.stream_reader() and workflow.stream_writer(), with StreamReader and StreamWriter. They default to the default topic, and a reader can start at END or at the last N records.
  • The worker takes the provider from its own config or from the client's, and hands it to each workflow instance and to the replayer.
  • A lifecycle interceptor runs the provider's start hook on the workflow's own event loop, after the workflow's __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 -

  • Unit Tests
  • Staging
  • End to End Tests

poe lint is 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.

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.
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