Skip to content

Put the shared stream interface on the repaired external SDK. - #2

Closed
moedash wants to merge 70 commits into
task/python-sdk-streamingfrom
moe/AI-198-stream-handle-on-fixed-core
Closed

moedash wants to merge 70 commits into
task/python-sdk-streamingfrom
moe/AI-198-stream-handle-on-fixed-core

Conversation

@moedash

@moedash moedash commented Sep 14, 2026 •

Copy link
Copy Markdown
Owner

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 his external_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:

  • Retention by trimming. 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's StreamIntegrityError (the external storage cause, the integrity counter) and a message that names the window; an outside read(after=) below the trim raises StreamCursorError on its first pull instead of resuming from the first retained record; a fully trimmed topic reads as empty and its latest() is BEGINNING. The replay driver lets a backend's own integrity error through where it used to wrap one as a transient storage failure.
  • A start cursor for the workflow reader. workflow.stream_reader(topic, after=) seeds the transport's subscription from a redis:in:<ms>-<seq> cursor, the form a workflow reader returns, so a run can start after a record its predecessor saw. His subscribe() 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 with StreamCursorError, 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-protos regenerates the vendored temporalio/api byte-for-byte from the pin, and the branch carries the stream_provider key in ClientConnectConfig.

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 -

  • Unit Tests
  • Staging
  • End to End Tests

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_tool under latest deps, caused by the old contrib base; upstream passes the byte-identical test.

The live module also covers retention: a max_len window that the loop's records fit inside and six later appends push them out of, with the offline Replayer passing before and failing loudly after, the outside read below the trim refused and a read inside the window landing where append said; a retention window 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 the Replayer for both runs.

mdashti and others added 23 commits September 5, 2026 20:49
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
moedash force-pushed the moe/AI-198-stream-handle-on-fixed-core branch from ce66293 to 14b5042 Compare September 17, 2026 04:23
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
moedash force-pushed the moe/AI-198-stream-handle-on-fixed-core branch from f82a7af to c4bed86 Compare September 17, 2026 04:36
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.
`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.
@moedash

moedash commented Sep 25, 2026

Copy link
Copy Markdown
Owner Author

This PR is replaced by #15, #16 and #17, which carry the same tree split into the base repairs, the stream interface and the Redis provider. The tip of #17 is an empty merge of this branch's head d3aeb5691427, so #6's union merge and every sha pin that names this branch stay 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.

3 participants