Conversation
2 of 3 tasks
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.
CI runs `ruff check --select I` and `ruff format --check` over the whole tree, and the new files were written by hand.
The changelog checkpoint requires an entry for any user-facing change.
`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.
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.
moedash
force-pushed
the
moe/AI-198-streams-interface
branch
from
September 17, 2026 02:51
a92fe0a to
6b8b35e
Compare
A workflow's publish runs on the workflow thread, so waking a consumer parked on another loop needs call_soon_threadsafe. The memory provider also stops reading idle_timeout as its poll period, which gave the parameter a second meaning.
A caller resuming with read(after=appended) must see only what came later, so the cursor names the last record of the batch. None says nothing was written, or that the transport learns positions at read time.
The activity defaulting rule is an interface rule, and every provider carried its own copy.
Building a provider is process setup. Workers resolve the default through worker_options(); workflow code that opens a stream before that gets a clear error instead of a racy lazy build.
reader() and Consumer.read bind T through overloads, StreamRecord.value admits the control bodies it actually carries, writer() drops the type it ignored, and next() leaves the public surface.
…arning. An inbound stream's name is its whole address, so its records carry no topic and a topic filter there is a ValueError on every provider. Both readers now skip a frame the interface did not write and log it, instead of one failing the task while the other drops it silently.
The reader and writer handles run in a real workflow on the memory provider. Rule 1 and rule 2 are strict expected failures there, so a storage provider running this module turns them into passes.
Joining the workflow id and the stream name with a bare colon let two addresses share a store. Every provider that keys by the pair now goes through stream_key, which percent-encodes both components first.
A provider that holds a connection pool or an HTTP session needs a moment where the process says it is done. It is a no-op when nothing is configured, so a shutdown path can call it unconditionally.
The memory provider always runs; a store branch registers its own setup behind its STREAMS_LIVE gate and declares the capabilities it lacks, which the fixture turns into skips with a reason. Cases read with a topic wherever they can, and two new cases prove colon-bearing addresses never share a store.
…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.
A store that lives inside a running workflow has no stream until that workflow runs, so a provider's setup can hand the cases a host that starts it before they open a handle.
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.
`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.
`moe/AI-198-api-protos-on-85b71d7` is upstream's `85b71d7` plus the api branch's stream diff on its `api_upstream`, so the vendored `temporalio/api` regenerates byte-for-byte from the pin, which is what the `check-protos` job asserts. The bridge protos are unchanged.
This was referenced Sep 25, 2026
Owner
Author
|
Replaced by four stacked PRs over the same tree: #10 the vendored stream protos and the Core pin, #11 the |
2 of 3 tasks
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
temporalio.streams, the stream interface a workflow reads, decides on and writes:workflow.stream_reader()andworkflow.stream_writer()for workflow code,activity.stream_handle()for an activity (its own workflow, pinned to its run),client.get_stream_handle()for any process holding a client, and theStreamProviderandWorkflowStreamProviderprotocols a store implements; a topic is a typed definition,streams.topic("inputs", Token), shared by workflow, activity and client code so the record and value types follow from it, and a plain string names a topic decided at runtime. A provider is registered once, as a plugin on the client (Client.connect(plugins=[provider])); workers built from that client inherit it, the worker carries it into the workflow runtime and the activity context, and brackets the workflow function with the provider's lifecycle hooks, so nothing is global and no workflow code calls a hook. The record on every wire is the vendoredtemporal.api.stream.v1.StreamRecordproto, a topic is the one noun on both sides, cursors name their provider, aSUPERSEDEDrecord sits before the record that triggered it, andread()ends when the owning chain is closed and its tail delivered;MemoryStreamsis the in-memory reference provider and the conformance suite states each of those rules as a case.Two later fixes:
ClientConnectConfiglistsstream_provider, because the TypedDict mirrorsClient.connect's keyword arguments and the client test compares the two; and the Core submodule pins a protos-only branch, the commitmainpins plus the stream protos, socheck-protosregenerates the vendoredtemporalio/apibyte-for-byte from the pin.Fixes AI-198
Why?
The contract has to be provider-independent so storage can swap under an unchanged workflow file. Five rules carry it: transactional publish, recorded reads, producer identity, opaque cursors, relative stream naming. The in-memory provider is the reference semantics; when two providers disagree, it is the tiebreak.
How did you test it?
Link to a test plan if any -
The conformance suite and the demo loop on the memory provider, with the provider registered on the client alone. Every later provider PR runs the same conformance suite against its own store.