Skip to content

Delivered consumed stream ranges to workflows, live and on replay. - #21

Open
moedash wants to merge 3 commits into
moe/AI-198-st-core-2-machinesfrom
moe/AI-198-st-core-3-delivery
Open

moedash wants to merge 3 commits into
moe/AI-198-st-core-2-machinesfrom
moe/AI-198-st-core-3-delivery

Conversation

@moedash

@moedash moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner

This PR hands a workflow the stream ranges it consumed, both live and on replay.

What changed?

  • Core reads the slices on the poll response. An untagged slice is the range for the task about to run. A tagged one re-supplies what an earlier task consumed, keyed by the completion that recorded it.
  • Live, the ranges arrive as DeliverStreamRecords jobs in a fixed order by stream and offset, after the channel notifications.
  • On replay, each recorded range comes back in the activation of the task that consumed it. Core looks ahead to the completion that closes the task, fetching the next history page when needed, so a read-then-publish task gets its input before it reissues its commands.
  • A task that consumed a range gets its own activation on replay rather than being folded into a heartbeat chain.
  • When a recorded range has content but the server sent no records for it, or records that don't match, the task fails with MissingRecords before the workflow runs. A legacy query sent that way goes unanswered, so the server retries it on the normal task queue, where the records travel with it.
  • HistoryForReplay::with_stream_slices lets a language replayer pass the records it fetched from the stream service.
  • The changelog has entries for streams.

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

Why?

History holds only the offsets a task consumed, so replay depends on the server sending the bytes back. Replaying a task on less input than it ran on produces different commands, and that would surface later as an unrelated nondeterminism error. Failing early with a clear cause is better than failing late. The order ranges arrive in is something the workflow can branch on, so Core fixes it instead of taking whatever order the server produced.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

core_tests::streams covers live delivery, empty ranges, replay order, the lookahead, a page boundary, data-only tasks, a cached run, missing and mismatched records, the legacy query and pushed histories. A channel case checks that notifications come ahead of a stream range in the same activation. Every workspace lib suite passes, the lints and fmt are clean, and each commit builds.

The poll response carries the records for the task about to run and re-supplies what earlier tasks consumed, keyed by the completion that recorded each range. Core hands them over in a fixed order, gives a consuming task its own activation on replay, and fails a task whose recorded range arrived without its records.
The cases cover live delivery, replay across history pages, the lookahead to the closing completion, missing and mismatched records, and notifications ordered ahead of a stream range.
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