Skip to content

Added the stream accessors and tied in the interface head. - #13

Closed
moedash wants to merge 75 commits into
moe/AI-198-py-03-workflow-runtimefrom
moe/AI-198-py-04-accessors
Closed

moedash wants to merge 75 commits into
moe/AI-198-py-03-workflow-runtimefrom
moe/AI-198-py-04-accessors

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR adds the client and activity stream accessors and ties in the original interface head.

What changed?

  • Client.get_stream_handle() mirrors get_workflow_handle: without a run_id it follows the chain. It also takes an activity_id, a StreamRef or a stream_id. client.create_stream() creates a standalone stream with its policy.
  • activity.stream_handle() resolves by a static rule, never by probing. A workflow's activity reaches its workflow's stream, a standalone activity reaches its own, and scope="activity" gives a workflow's activity its own streams. A def activity is told it needs async def.
  • StreamProvider gains get_activity_stream_handle(). The memory provider hosts activity owners and standalone streams, with the policy trimmed on every append and close() sealing.
  • Client.connect and ClientConnectConfig take stream_provider.
  • The conformance cases open through the client accessor. They gain ref, standalone and retention cases behind hosts_standalone_streams.
  • streams_demo/ shows each context asking for its stream. CHANGELOG.md has the entry.
  • The last tie-in commit is an empty merge of the replaced Added the stream interface, the memory provider and the conformance suite. #3's head. It gives the chain that head as an ancestor, so the provider branches that grew from Added the stream interface, the memory provider and the conformance suite. #3 merge onto it cleanly.

Part of AI-198 (epic AI-37).

Why?

Building a provider is process setup, so it happens once and every context inherits it. Workflow code that opens a stream before that gets a clear error, not a racy lazy build. The static rule matters because an owned stream is created by its first write. Probing would send attempt 1 and its retry to different streams.

How did you test it?

I ran uv run poe lint and uv run pytest tests/streams tests/test_client.py, plus the plugin, worker and client-exports suites. test_stream_accessors.py proves a worker inherits the client's provider and carries a ref through a workflow argument and result. test_activity_streams.py covers the three accessor arms and a retry that surfaces SUPERSEDED. git diff against #3's head confirms the tie-in merge is empty.

  • Unit Tests
  • Staging
  • End to End Tests

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.
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.
Handles from client.get_stream_handle() and activity.stream_handle() reach the same default topic as the workflow side, so the accessor docs and the changelog say so and a test pins it across all three contexts.
A standalone activity now reaches its own stream, which used to raise, and
scope="activity" gives a workflow's activity its own streams without changing
what a bare call reaches. The rule never probes, because a stream is created
by its first write and probing would split attempts across two streams.
Covers all three arms of the rule, a retry writing to the same stream, and the misaddressed calls, on every provider the tree can stand up.
…4-accessors

# Conflicts:
#	temporalio/streams/providers/memory.py
#	tests/streams/test_streams_conformance.py
…4-accessors

# Conflicts:
#	temporalio/streams/_provider.py
#	temporalio/streams/providers/memory.py
#	tests/streams/test_streams_conformance.py
A StreamRef in place of the workflow id opens the stream it names on whatever provider the client carries, and its topic becomes the handle's default. A standalone stream is created with client.create_stream and reached by its stream_id.
A standalone stream lives in the provider with its policy and a sealed flag. The policy trims on append, a seal wakes parked readers so their reads end, and a stream id that was never created is refused at the call, so the conformance suite can run the standalone cases without a server.
The conformance suite opens a ref, creates and seals a standalone stream and checks its retention policy on every provider that hosts one. The accessor tests carry a ref through a workflow argument and result and open it from the activity and the client.
…4-accessors

# Conflicts:
#	tests/streams/test_streams_conformance.py
…e's client.

A provider that cannot bound a standalone stream by bytes is now expected to refuse the policy, and one that trims by age only after close skips that check. The converter cases derive their client from the provider's own, so a live provider is not asked to reach the fixture's dev server.
A store that refuses the append crossing its byte cap is checked with a cap that holds one stored record, hash metadata included, and not two. The age subcase waits for the older record to leave BEGINNING and accepts the newer one aging out as well, since a store reclaims on its own schedule.
…4-accessors

# Conflicts:
#	temporalio/streams/__init__.py
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36, #37, #38, #39, #40, #41, #42, #43, #44, #45, #46, #47, #48, #49, #50, #51, #52, #53, #54 and #55.

Same content, split into 37 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

@moedash moedash closed this Oct 3, 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.

1 participant