Skip to content

Replayed workflows that read native streams. - #54

Open
moedash wants to merge 4 commits into
moe/AI-198-st-py-3-native-providerfrom
moe/AI-198-st-py-4-native-replay
Open

moedash wants to merge 4 commits into
moe/AI-198-st-py-3-native-providerfrom
moe/AI-198-st-py-4-native-replay

Conversation

@moedash

@moedash moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner

This PR lets the replayer replay a workflow that read a native stream, live or from an exported history.

What changed?

  • History records only the offsets each task consumed. The worker fills short ranges before handing them to the workflow, and the replayer fetches the records from the stream service with Replayer(stream_client=).
  • A replay across a reset fetches the ranges recorded before the reset point from the run it was reset from. A range the stream no longer holds fails the replay with StreamNotFoundError.
  • Replayer.fetch_stream_slices(client, history) attaches the records to a WorkflowHistory while the stream still holds them. to_json() and from_json() carry them as streamSlices, so that history replays with no server. The bridge's push_history passes the slices through.
  • The changelog entry for server-side streams.

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

Why?

A workflow that reads a native stream can't be debugged or tested with the replayer unless the replayer can get the same records back. History deliberately holds offsets, not records, so the replayer has to fetch them, or carry them in an exported history for offline use.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

The bridge builds, and cargo clippy and poe lint are clean. The re-supply, replayer, offline-history and native stream workflow cases pass live, against a local server built from the stream server PRs with STREAMS_LIVE=native. Those include the cases that replay against a live native stream and across a reset. The whole streams suite passes there too.

History records only the offsets each task consumed, so the replayer fetches those records from the stream service, or from slices carried with an exported history, and hands them to the replay.
Re-supply of short ranges, the replayer against a live stream, reset runs and offline histories with stream slices.
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