Skip to content

Delivered stream records on Workflow Tasks and on replay. - #7

Closed
moetemp wants to merge 30 commits into
moe/AI-198-srv-4-workflow-commandsfrom
moe/AI-198-srv-5-delivery-replay
Closed

moetemp wants to merge 30 commits into
moe/AI-198-srv-4-workflow-commandsfrom
moe/AI-198-srv-5-delivery-replay

Conversation

@moetemp

@moetemp moetemp commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR delivers stream records on Workflow Tasks and re-supplies them on replay.

What changed?

  • recordworkflowtaskstarted/stream_slices.go stages a range per subscription and reads it from an owned stream or, through the routed client, a standalone one. Slices ride RecordWorkflowTaskStarted and the poll response. A duplicate request gets the same slice. Speculative tasks get no range, since a discarded one commits no cursor.
  • The task records only the range it consumed, on WorkflowTaskCompleted. The cursor advance is folded in as that event is built, so both commit together. A subscription with offsets left schedules a task on transaction close, which also covers a fresh subscription to a full stream.
  • Replay walks every page of the history it sends and re-reads each range. A range that's truncated, deleted or past stream.replayMaxRecords, stream.replayMaxBytes or stream.replayMaxPages fails the task with STREAM_RANGE_UNAVAILABLE. The second failure in a row terminates the workflow, since one retry covers a moving shard and a second isn't a race. There's no paged re-supply and no short flag on the wire yet, as chasm/lib/stream/doc.go notes.
  • All routed reads of one task start share stream.RoutedSetBudget, the deadline the completion path uses, because the execution's lock is held across them.
  • GetStreamReplaySlices lets matching re-supply a non-sticky query task, which skips RecordWorkflowTaskStarted. Matching calls it only when GetMutableStateResponse says consumes_streams, so other queries take no history lease. attachReplaySlices returns early on a sticky task, so an SDK that refetches history after eviction gets no slices.
  • tests/testcore/onebox.go starts more than one history host for the cross-host test. tests/sdk_server_host_test.go and tests/stream_publish_cost_test.go are opt-in. TestARefusedTruncateProbesThePinsItWasRefusedFor tests a Added the stream service and its frontend wiring. #5 change. It lands here with the consumer harness it needs.

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

Why?

This is the half of the design that has to hold under replay. A cold worker sees what the original execution saw, and History never holds a payload.

How did you test it?

go build ./..., make proto and make fmt are clean. Unit tests ran for ./chasm/lib/..., ./service/history/..., ./service/matching/ and ./common/.... The functional tests ran with -tags test_dep: delivery, continue-as-new, replay past the first page, a deleted stream, the query path, cross-host routing and speculative tasks.

  • Unit Tests
  • Staging
  • End to End Tests

Records ride the task response and the completion records only the offset range consumed, so replay re-reads the payloads from the stream instead of History carrying them. A range that cannot be served fails the task with its own cause, and a non-sticky query asks History for the same ranges after it has fetched the history.
A discarded speculative task commits no cursor, so the range it was given came
back on the next task while the worker had already consumed it. The replay
budget is dynamic config now, a consumer that stays over it is terminated with
the cause instead of failing the same task forever, and a query only asks
History for ranges when the workflow actually holds a cursor.
…-5-delivery-replay

# Conflicts:
#	service/history/api/recordworkflowtaskstarted/stream_routing.go
A workflow may consume as many streams as the subscription limit allows, and
the execution's lock is held across every read.
The subscribe command names a stream by name or by id, and an unnumbered record no longer rides a sentinel.
…-5-delivery-replay

# Conflicts:
#	api/historyservice/v1/request_response.pb.go
#	api/matchingservice/v1/request_response.pb.go
The ConsumesStreams entry had been hand-placed out of alphabetical order, so
a later mockgen run would have rewritten it.
Main now builds under Go 1.27, where go fix rewrites errors.As and reverse
index loops, so make fmt was left dirty on these two files.
…-5-delivery-replay

# Conflicts:
#	service/history/workflow/mutable_state_impl.go
…-5-delivery-replay

# Conflicts:
#	chasm/lib/stream/config.go
…n the wire.

Paging the re-supply needs a short flag on the poll response or the last
slice, which is a wire change; the package doc says so next to the other
wire-round notes so the gap is not lost.
A wake tells a running execution that a source it consumes moved. It folds into a wake no started task has carried, rides the next task, and is dropped when that task completes, so nothing about it enters History.
@moetemp moetemp closed this Oct 1, 2026
@moetemp moetemp reopened this Oct 1, 2026
@moetemp moetemp mentioned this pull request Oct 1, 2026
1 of 3 tasks
@moetemp

moetemp commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #18, #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36 and #37.

Same content, split into 20 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.

@moetemp moetemp closed this Oct 3, 2026
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