Skip to content

Replayed a workflow that read a server-side stream. - #14

Closed
moedash wants to merge 84 commits into
moe/AI-198-py-06-native-providerfrom
moe/AI-198-py-07-native-replay
Closed

moedash wants to merge 84 commits into
moe/AI-198-py-06-native-providerfrom
moe/AI-198-py-07-native-replay

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR replays a workflow that read a server-side stream, and follows a workflow reset.

What changed?

  • Replayer(stream_client=) fetches every range the completed tasks recorded from the stream service and hands the records to Core with the history. A range the stream no longer holds fails that replay with StreamNotFoundError. Without a client, a consuming workflow's history is refused with both remedies in the message.
  • The replayer splits a history at each reset marker and fetches each era from the run whose stream holds it. HistoryPusher.push_history takes the slices as serialized StreamSlice messages.
  • WorkflowHistory.stream_slices carries the records. to_json and from_json write them as a streamSlices list beside the events, and Replayer.fetch_stream_slices(client, history) fills them while the stream is retained. The exported file is then the whole replay input.
  • The native handle follows a reset through describe, across a chain of resets. A reset run's inherited stream starts at its floor, so a chain-following read never asks below it.
  • temporalio/worker/_stream_ranges.py holds the range readers and fill_short_stream_ranges, the worker's half of paged re-supply. Nothing triggers it yet, because the server and Core don't pass a short slice through.
  • The replayer's channel comes from the client's connection and closes when the replay ends.
  • CHANGELOG.md gets the entry for the native feature.

This is the tip of the native chain. It ends with an empty merge of the replaced #1's head, so #6 and every pin naming that head stay valid.

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

Why?

History records the offsets each task consumed, never the records, so a workflow that read a stream can't be replayed without this layer. A support case arrives as a history file, so the export has to carry its records. The sharp edge is a sticky task to a worker that evicted the run, which gets no records. The pinned Core fails that task before the workflow runs, so a legacy query goes back to the server and is retried non-sticky.

How did you test it?

The lint set from #8 ran on a cold .mypy_cache. test_replayer.py and test_stream_resupply.py run without a server, and test_workflow_stream_e2e.py runs on a server from moedash/temporal#17. It covers the replayer with a client, tampered ranges, a deleted stream and no client, a query on a cold worker, a sticky query to an evicted run, a reset and a chain of two, a pinned read of a reset run from BEGINNING, a codec on both halves, and an exported history with no client. The reset case needs a server that seeds the reset run's stream.

  • Unit Tests
  • Staging
  • End to End Tests

Porting the agent harness onto the interface surfaced three things the small
example never needed. A reader resumes after the record a cursor names, so
it can store the last cursor it handled without advancing an opaque token. A
consumer reports its latest position, so a client can follow a turn it is
about to start without the workflow reporting a position. And a producer can
append onto a topic of the stream the workflow publishes, so an activity's
live output lands next to the workflow's own records for one outside reader.
The buffer holds a body, topic and offset rather than the wire message,
and the cap moved to `workflow_read_stream_messages`.
The pinned Core carries the stream additions to the public API, so the
command, enum, history and workflow service messages were stale. Only the
vendored stream service protos still come from the server checkout.
CI runs `ruff check --select I` and `ruff format --check` over the whole
tree, and `check-protos` reformats what it regenerates.
A Standalone Activity has no `workflow_id`, so opening its Workflow's
stream now says that rather than sending `None` on the wire. `Optional`
and `typing.Sequence` are deprecated spellings the linter rejects.
The changelog checkpoint requires an entry for any user-facing change.
This package has to load on a protobuf 3 runtime, and generated code
refuses to load on a runtime older than the one it was built against. The
protoc of that generation has no `--pyi_out`, so the stubs come from
`mypy-protobuf`, which is what the other generated protos already use.
The stream client needs the replay slice lookahead fix that branch finished.
The new server_streams module gives WorkflowStream and WorkflowStreamClient a second definition, so pydoctor can no longer resolve the short names.
That head matches upstream's nexus model WIT for the workflow-id policies, which the current generator requires.
The generated module was missing the local-activity marker argument flag its own stub already declared, so the two disagreed about the message.
The append cursor names the last record and is None for a batch the server deduplicated, so the stream client reports what the server said about the append. Producer identity is resolved by the package, inbound frames carry no topic, and the inbound id comes from the shared helper.
A handle on a workflow's stream resolves the run once when it opens, so a follower is not redirected to a successor's empty stream after continue-as-new. The channel cache is keyed by the loop object and closed through the provider's close hook, and latest() reads a missing stream as empty but lets every other failure through.
A cancel landing inside the flusher's append unwound with the batch it had taken off the buffer. The client now asks it to stop and waits. The handle is pinned to the run it opened on, from the activity's own run id or one describe, and the channel cache is the shared one.
A range is recorded as consumed once and never resent, so a repeated, skipped or mis-sized delivery would hand the workflow duplicate or shifted bodies with nothing to say so. The worker also refuses a publish over the server's per-message and per-batch byte limits before the command exists, for the reason the count limit already gave.
The native activation job and the two native commands moved to the field numbers the combined Core tree uses, so both trees agree.
The activation job is now deliver_stream_records and the command
append_stream_records, both carrying StreamRecord, so the payload visitor
reaches each record's body and metadata.
…lishes.

The delivered job and the publish command carry StreamRecord protos, and a
task's publishes on one stream are held until the task completes so they
become one command and one History event.
…rors.

Every stub call translates the transport's failure into StreamNotFoundError
or RPCError, and the client's payload codec is applied to bodies so the two
sides of a namespace with a codec agree.
The workflow subscribes to the topic by name and publishes with the batched
command; outside code appends and long-polls the same stream, following the
run chain by cursor unless pinned.
This layer looks a successor up per topic and returns where to resume, so the activity handle's override, which never has a successor, takes the same arguments.
…-native-replay

# Conflicts:
#	temporalio/streams/providers/native.py
#	tests/worker/test_workflow_stream_e2e.py
The replayer fetched recorded ranges over a bare channel to the client's host,
so a client with TLS or an API key named the right server and could not
authenticate. It now opens the shared channel from the client's connection and
closes it by the same key. A range refused below the retention floor arrives
typed now and is reported as the stream no longer holding it, as before.
…ran.

The server re-supplies a replaying workflow's recorded ranges within a budget,
and a range it cannot fit has to come from somewhere or the task cannot run.
The worker now reads the missing offsets of a short DeliverStreamRecords job
through the client's stream channel before the activation runs, with the range
readers the replayer already had, moved where both can use them. The server
still refuses an over-budget task outright and the pinned Core fails a short
slice itself, so this half waits on a server flag and a Core pass-through.
The server now seeds the reset run's inherited stream with the range the
re-run task had consumed, at the offsets it held, so that task decides on the
same input twice over and the inherited stream's floor holds the seeded
record. The reset case follows that contract from the srv chain at 217eed936.
…-native-replay

# Conflicts:
#	temporalio/worker/_workflow.py
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36, #37, #38, #39, #40, #41, #42, #43, #44, #45, #46, #47, #48, #49, #50, #51, #52, #53, #54 and #55.

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

@moedash moedash 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