Conversation
moedash
force-pushed
the
moe/AI-198-stream-client
branch
from
September 15, 2026 22:57
e1583a7 to
16c7035
Compare
2 of 3 tasks
moedash
force-pushed
the
moe/AI-198-stream-client
branch
3 times, most recently
from
September 16, 2026 23:23
df741be to
3713d74
Compare
moedash
force-pushed
the
moe/AI-198-streams-interface
branch
from
September 17, 2026 02:51
a92fe0a to
6b8b35e
Compare
moedash
force-pushed
the
moe/AI-198-stream-client
branch
5 times, most recently
from
September 21, 2026 17:30
9f5a1cf to
8109d26
Compare
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
force-pushed
the
moe/AI-198-stream-client
branch
from
September 21, 2026 17:54
8109d26 to
d5e8f4b
Compare
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.
This was referenced Sep 25, 2026
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 |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changed?
This PR adds the native provider on the v2 interface:
NativeStreamsis a worker plugin whose topic is one owned server stream named after the topic, the workflow subscribes to it by name and publishes with oneAppendStreamRecordscommand per task, and outside code appends and long-polls the same stream through the stream client on the record vocabulary. A cursor isnative:<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 asStreamNotFoundErrororRPCError.Replayertakes astream_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 withStreamNotFoundError. 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.
WorkflowHistorygainsstream_slices, carried byto_json()andfrom_json()as astreamSliceslist 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
Replayersplits 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_providerconfig key. The Core pin stays on this branch's record-vocabulary head, and the vendoredtemporalio/apiregenerates 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
Replayerat 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 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; theReplayergiven 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 theReplayeron 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.