Skip to content

Added the stream interface, the memory provider and the conformance suite. - #3

Closed
moedash wants to merge 36 commits into
mainfrom
moe/AI-198-streams-interface
Closed

moedash wants to merge 36 commits into
mainfrom
moe/AI-198-streams-interface

Conversation

@moedash

@moedash moedash commented Sep 15, 2026 •

Copy link
Copy Markdown
Owner

What changed?

This PR adds temporalio.streams, the stream interface a workflow reads, decides on and writes: workflow.stream_reader() and workflow.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 the StreamProvider and WorkflowStreamProvider protocols 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 vendored temporal.api.stream.v1.StreamRecord proto, a topic is the one noun on both sides, cursors name their provider, a SUPERSEDED record sits before the record that triggered it, and read() ends when the owning chain is closed and its tail delivered; MemoryStreams is the in-memory reference provider and the conformance suite states each of those rules as a case.

Two later fixes: ClientConnectConfig lists stream_provider, because the TypedDict mirrors Client.connect's keyword arguments and the client test compares the two; and the Core submodule pins a protos-only branch, the commit main pins plus the stream protos, so check-protos regenerates the vendored temporalio/api byte-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 -

  • Unit Tests
  • Staging
  • End to End Tests

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.

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
moedash force-pushed the moe/AI-198-streams-interface branch from a92fe0a to 6b8b35e Compare September 17, 2026 02:51
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.
@moedash moedash changed the title Added the stream interface, provider registry, and conformance tests. Added the stream interface, the memory provider and the conformance suite. Sep 21, 2026
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.
@moedash

moedash commented Sep 25, 2026

Copy link
Copy Markdown
Owner Author

Replaced by four stacked PRs over the same tree: #10 the vendored stream protos and the Core pin, #11 the temporalio.streams package with the memory provider and the conformance core, #12 the workflow-side reader and writer, and #13 the accessors. #13 ends with an empty merge of this branch's head 1a8d19fe, so the tip tree is identical to this one and #4 and #5 keep their diffs and merge cleanly after the chain lands. The branch stays; only the PR closes.

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