Skip to content

Added server-side streams to the Python SDK. - #1

Closed
moedash wants to merge 43 commits into
moe/AI-198-streams-interfacefrom
moe/AI-198-stream-client
Closed

moedash wants to merge 43 commits into
moe/AI-198-streams-interfacefrom
moe/AI-198-stream-client

Conversation

@moedash

@moedash moedash commented Sep 14, 2026 •

Copy link
Copy Markdown
Owner

What changed?

This PR adds the native provider on the v2 interface: NativeStreams is a worker plugin whose topic is one owned server stream named after the topic, the workflow subscribes to it by name and publishes with one AppendStreamRecords command per task, and outside code appends and long-polls the same stream through the stream client on the record vocabulary. A cursor is native:<run_id>:<offset>, so a handle without a run id follows the chain run by run and ends once the last run is closed and drained, while a handle with a run id stays pinned. The bridge is repinned to the record-vocabulary Core, the payload visitor covers each record's body so a codec applies on both sides and the codec refusal is gone, and every gRPC failure surfaces as StreamNotFoundError or RPCError.

Replayer takes a stream_client. Given one, it fetches every range a history's completed tasks recorded from the stream service, resolving a subscribed name as the server does (a stream the workflow owns by that name first, else a standalone stream by that id), 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 the remedy in the message, because History records only the offsets.

Offline replay from an export. WorkflowHistory gains stream_slices, carried by to_json() and from_json() as a streamSlices list beside the events; without slices the JSON is unchanged. Replayer.fetch_stream_slices(client, history) returns a copy with the records attached while the stream is retained, so the exported file is the whole replay input and a replayer given it needs no server. The message for a history that carries neither slices nor a client names both remedies.

A handle without a run id follows a reset as well as a continue-as-new. A run that closed without continuing as new is described, and the run it was reset into, which only describe names, is read next from the floor its stream reports: a stream the reset run inherited a subscription to continues the base run's offset space, one the base run only published to starts at zero. A handle pinned to the base run ends with it. The Replayer splits a reset run's history into eras at each reset marker and fetches each era from the run whose stream holds it, the run reset from for the ranges before the marker.

Merged from the interface branch: the stream_provider config key. The Core pin stays on this branch's record-vocabulary head, and the vendored temporalio/api regenerates from it with no diff. In addition, the eras docstring marks the reset event name as a literal, because pydoctor read the single-backtick span as a link and failed the docs build on it.

Fixes AI-198

Why?

Binds the shared interface to the server prototype so the same workflow file runs on server-side storage. Wakes arrive through task dispatch, so there is no signal-race retry on this provider. A consuming workflow could not be replayed by the Replayer at all before, since its records live in the server's stream and not in History; a reader following a workflow that was reset stopped at the base run; and an exported history could not be replayed once the stream was gone.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

Unit plus live stream tests against a server built from the stream branch (STREAMS_LIVE=native), including the shared conformance suite on the native provider and, inside a workflow, the publish that fails with its task and the cold-cache replay the server re-supplies. Also live: an encrypting payload codec on both halves of the native provider, with the stored bodies read raw through the stream service; the Replayer given a client on a consuming workflow's history, on the same history with its recorded ranges tampered, on a history whose stream is gone, without a client, and on a workflow that consumed a standalone stream; a query against a cold worker and a sticky query against a worker that evicted the run, both answered from replayed state; a reset of a consuming workflow, with the reset run replaying the base run's ranges, a chain-following read crossing into it on both an inherited and a fresh stream, a pinned base reader ending with the base run, and the Replayer on the reset run's history; and a history exported with its records, round-tripped through JSON and replayed with no client, for a plain run and for a reset run, with the stripped export refused. A unit test pins the JSON shape with and without slices.

@moedash
moedash force-pushed the moe/AI-198-stream-client branch from e1583a7 to 16c7035 Compare September 15, 2026 22:57
@moedash
moedash changed the base branch from main to moe/AI-198-streams-interface September 15, 2026 22:57
@moedash
moedash force-pushed the moe/AI-198-stream-client branch 3 times, most recently from df741be to 3713d74 Compare September 16, 2026 23:23
@moedash
moedash force-pushed the moe/AI-198-streams-interface branch from a92fe0a to 6b8b35e Compare September 17, 2026 02:51
@moedash
moedash force-pushed the moe/AI-198-stream-client branch 5 times, most recently from 9f5a1cf to 8109d26 Compare September 21, 2026 17:30
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 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.
@moedash
moedash force-pushed the moe/AI-198-stream-client branch from 8109d26 to d5e8f4b Compare September 21, 2026 17:54
The Core pin moves to the head that puts a history's stream slices on the
replay worker's poll response. `push_history` takes them as serialized
`StreamSlice` messages, so a replayer can hand Core the records History
does not hold.
An encrypting codec on the client has to cover the records a workflow
publishes and receives through the worker's payload visitor, and the records
the outside half sends and reads. The stored bodies are read raw through the
stream service to show they are not plaintext.
History records the offsets each task consumed and never the records, so a
workflow that read a server-side stream could not be replayed at all. Given
a `stream_client`, the replayer fetches every recorded range from the stream
service, resolving a subscribed name as the server does, and hands the
records to Core with the history. Without one, such a history is refused
with the remedy in the message.
A sticky task or a sticky legacy query handed to a worker that evicted the
run arrives with partial history and no records, and the history the worker
fetches itself carries none either. Core now fails such a task before the
workflow runs on less input than it had.
A run's close event names a continue-as-new successor but not the run it
was reset into; describe does. The reset run's inherited streams continue
the base run's offset space and refuse anything below their floor, so a
chain-following read starts them at the floor the stream reports. A reset
run's start event is the base run's, copied, so its original run id leads a
handle opened afterwards back to the base run.
The ranges recorded before a reset point were consumed from the base run's
streams, and only the reset marker after them names that run. The replayer
splits the history at each marker the way the server does when it
re-supplies a cache miss.
A recorded range with content that the response carried no records for is
now the worker's failure to reconstruct the run rather than nondeterminism,
so a legacy query on the sticky queue of an evicted run is left for the
server to retry where the records travel, and a task fails without ever
failing the workflow over it.
A cold worker answers a query from replayed state because the non-sticky
query task carries the recorded ranges. A sticky query to a worker that
evicted the run goes unanswered and is retried non-sticky. A reset run
replays the base run's ranges, carries its subscription on at the inherited
offset, and is followed by chain-following readers and the replayer.
History records only the offsets a task consumed, so an export alone could
not be replayed once the stream was gone. `WorkflowHistory.stream_slices`
holds the records, `to_json` and `from_json` carry them beside the events
with the plain shape unchanged, and `Replayer.fetch_stream_slices` fills
them while the stream is retained, so the exported file is the whole replay
input.
mypy-protobuf writes the public API types fully qualified into the stubs,
`temporal.api.common.v1.message_pb2.Payload`, which the generator's import
rewrites did not reach, so a cold mypy run found `temporal` undefined in
three stubs. The rule matches the module form only, so the proto package
names inside the serialized descriptors stay as they are.
The gitlink is resolved to sdk-rust #1's head, `32c66693`, whose
`api_upstream` carries the same stream protos as the protos-only branch
the interface branch pins, plus the native stream machinery this branch
uses.
pydoctor reads a single-backtick span as a link target and fails the
docs build on a proto event name no Python object carries.
@moedash

moedash commented Sep 25, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #8, #9 and #14, which are this same tree split into three reviewable layers: the Core pin with the vendored stream service protos and the stream client, the native provider on both halves, and replay with reset following. The tip of #14 is an empty merge of this branch's head 5590dbc6, so the chain has it as an ancestor, git diff 5590dbc67130 <tip> is empty, and every sha pin that names this head stays valid.

@moedash moedash closed this Sep 25, 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