Conversation
The server refuses the wake while the execution is closing, and at that moment its status is still the running one, so a single look read "running" and turned the ordinary ending into an error.
A provider that answers outside readers through handlers on the workflow has to register them before the first task completes, and one that parks a long-poll update against the run has to let go before the workflow returns. Both are no-ops on the storage providers, so workflow code calls them unconditionally and stays portable.
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.
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.
A replay marker carries what its Workflow Task published, and installing it resets the publish records. The activation applied it with the non-query jobs, after the signal and update set had already run and published, so a task whose publishes a signal woke replayed as zero records against a manifest of several, and any query against the completed run failed. The marker now goes in with the first set that drains, after that set's other jobs and before the drain.
CI runs `ruff check --select I` and `ruff format --check` over the whole tree, and the new files were written by hand.
`prepare` and `drain` are hooks only a transport that parks something against the running workflow needs, so they move to their own protocol rather than making every provider carry two empty methods. The demo's teardown called a `close` that no provider setup defined.
The provider imports nothing from `workflow`, and its producer, consumer and `latest` need the docstrings the doc linter asks every public method for.
A record only reaches a reader after a provider placed it, so the fake runtime hands out placed records and the reader can name where each one came from. The subscription manager captures the loop it is built on, which is the Worker's, so its fixture has to build it on one.
Every Workflow Task that closes commits its own marker, and empty input activations now close one too, so a baseline read three steps earlier counts markers that belong to tasks this case says nothing about.
The changelog checkpoint requires an entry for any user-facing change.
A transport that serves outside readers through handlers on the workflow registers them in `prepare`, so without the call an outside reader finds no handler and the demo's follow fails against that provider.
* Support Google ADK 2.7 streaming import * Fix tests with latest dependencies * Apply suggestion from @brianstrauch
The branch's own copy shells out to a binary name the nexgen crate does not install, so generation died before it could diff anything.
Under the latest dependency set openai carries its own httpx as httpx2, so pyright will not accept the httpx.Response the test builds. Going through the class object was tried first and did not silence it.
moedash
force-pushed
the
moe/AI-198-stream-handle-on-fixed-core
branch
from
September 17, 2026 04:23
ce66293 to
14b5042
Compare
The same suite on the same Core is 673 passed and 0 failed against a dev server, and fails under time skipping. Which cases fail drifts between fourteen and seventeen across runs, so the suite is held as a whole rather than by a list of names that was never stable.
moedash
force-pushed
the
moe/AI-198-stream-handle-on-fixed-core
branch
from
September 17, 2026 04:36
f82a7af to
c4bed86
Compare
A script that drives two tool calls handed both invocations the same completed id, which the Agents SDK now rejects outright rather than tolerating. This diverges from upstream, which still hands out a constant there; unique ids are strictly more correct for a builder handing out completed call ids, so the right long-term home is upstream rather than here.
That branch now matches upstream's nexus model WIT for the workflow-id policies, which is what the current generator needs.
The generator needs its matching support file and the upstream lint exclusions, so both came across with it. One test asserted the old model's field names; the only consumers of that model outside the generated package are a generated re-export in temporalio.workflow and three package-level helpers, none of which touch those fields.
The generator renames the service module, and the visitor script still looked for the old filename, so the generation sequence died after the model was written.
…anches. The native and Redis providers import it under that name, so the interface exports the name they were told to use.
The Core submodule does not pin the api branch that carries them yet, so the generator stages the pinned API tree, applies the branch's diff and regenerates only the files it touches.
The registry, its string names and the private frame envelope go. A provider is an object that serves workers as a plugin and hands out handles outside them, and every store keeps the StreamRecord proto.
workflow.stream_reader() and stream_writer() reach it the way payload_converter() does, and the worker brackets the workflow function with the provider's lifecycle hooks so no workflow code has to.
Adds the cases for supersession cursors, append positions, a read that ends when the chain closes, chain following, foreign cursors, the lifecycle hooks and a subscription shared by two readers.
The transport keys every stream by direction, so an outside producer appends each record to the input stream that wakes a subscribed workflow and to the output stream where the workflow's promoted batches land. A workflow's publish is synchronous and stages with the task; an outside read ends with the workflow.
Neither is the workflow ending. An evicted run's state is not to be touched, and at collection time the runtime on the thread belongs to whichever workflow is running, so the hook acted on that one.
A provider is registered once, as a plugin on the client, and every context asks for its stream the same way: the worker inherits the client's provider, an activity's handle is its own workflow pinned to its run, and a client's handle mirrors get_workflow_handle.
A topic is defined once with streams.topic(name, type) and shared by the workflow, its activities and the backend, so the record and value types follow from the definition; a plain string still names a topic decided at runtime. The wire and every provider see only the name.
A read that passes a definition may not also pass result_type, and a record's topic is the definition's name.
The runtime and the manager already take a start cursor and the marker header records it; only `subscribe()` had no way to pass one. The provider now seeds the subscription from a `redis:in:` cursor instead of refusing it, and refuses an outside cursor because the two streams number entries apart.
A backend that can tell a recorded range is gone reports the loss itself. The wrapper filed it under a transient storage failure, which is retried and clears on its own; a trimmed range does neither.
`RedisStreams(retention=, max_len=)` trims a topic's input and output keys behind every append the provider makes, exactly, because the approximate trim never fires on a stream shorter than a macro node. A replay past the window fails its task with the transport's integrity error, an outside cursor below it is refused, and a fully trimmed topic reads as empty.
This was referenced Sep 21, 2026
`Client.connect` takes `stream_provider`, and the TypedDict that mirrors its keyword arguments has to name the same keys. The client test that compares the two failed on every branch since the keyword was added. (cherry picked from commit b068712)
A parked Run has no open Workflow Task, so the server dispatches a query on a task of its own and Core allows nothing beside the answer there. The instance still reported the registered wait set, Core refused the completion, and the query timed out. The test parks a Run and queries it.
The vendored `temporalio/api` now regenerates byte-for-byte from the pin, which is what the `check-protos` job asserts.
This was referenced Sep 25, 2026
Owner
Author
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 puts the v2 stream interface on Max's external-store base and binds Redis as
RedisStreams, a worker plugin that also sets hisexternal_stream_backend. His transport keys every stream by direction, so an outside producer appends each record to the topic's input stream, which wakes a subscribed workflow, and to its output stream, where the workflow's promoted batches land and outside readers follow; a workflow's publish is synchronous and stages with its task, and an outside read ends once the workflow is closed and every promoted record is delivered. The interface commits were cherry-picked one by one and meet his worker keywords in four files, each of which keeps both.Two later additions close the provider's remaining gaps:
RedisStreams(retention=, max_len=)trims a topic's input and output keys behind every append the provider makes, the outside producer's appends and the workflow's staged batches alike. The trim is exact, because Redis's approximate trim only drops whole macro nodes and never fires on a stream shorter than one. This is retention without a consumer floor: a replay that reaches a recorded range the trim removed fails its Workflow Task with the transport'sStreamIntegrityError(the external storage cause, the integrity counter) and a message that names the window; an outsideread(after=)below the trim raisesStreamCursorErroron its first pull instead of resuming from the first retained record; a fully trimmed topic reads as empty and itslatest()isBEGINNING. The replay driver lets a backend's own integrity error through where it used to wrap one as a transient storage failure.workflow.stream_reader(topic, after=)seeds the transport's subscription from aredis:in:<ms>-<seq>cursor, the form a workflow reader returns, so a run can start after a record its predecessor saw. Hissubscribe()takes the boundary and his runtime records it in the marker header, as it does the restored one. An outside cursor (redis:<ms>-<seq>) is refused withStreamCursorError, because the input and output streams number their entries independently.Two later fixes. A legacy query against a parked run timed out: the server dispatches such a query on a task of its own and Core allows nothing beside the answer there, yet the instance still reported the registered wait set with it. The instance now answers a legacy query alone. In addition, the Core submodule moved to the repairs head that carries the stream protos, so
check-protosregenerates the vendoredtemporalio/apibyte-for-byte from the pin, and the branch carries thestream_providerkey inClientConnectConfig.Fixes AI-198
Why?
This is the client-side branch of record; the measured comparison numbers came from it. The base predates main's contrib rewrite in a few places, so small local workarounds live here and are deliberately not propagated to the branches that sit on current main.
How did you test it?
Link to a test plan if any -
Full lint and the memory suites, plus live Redis runs (
STREAMS_LIVE=redis) of the shared conformance suite and the interface loop inside a workflow. One accepted red:test_hosted_mcp_toolunder latest deps, caused by the old contrib base; upstream passes the byte-identical test.The live module also covers retention: a
max_lenwindow that the loop's records fit inside and six later appends push them out of, with the offlineReplayerpassing before and failing loudly after, the outside read below the trim refused and a read inside the window landing whereappendsaid; aretentionwindow crossed by a later append; and a fully trimmed topic. The start cursor is proven by a workflow that continues as new with the cursor of its first record and reads the second again in the successor, live, on a cold query replay, and offline through theReplayerfor both runs.