Conversation
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
1 of 3 tasks
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. |
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.
This PR adds the client and activity stream accessors and ties in the original interface head.
What changed?
Client.get_stream_handle()mirrorsget_workflow_handle: without arun_idit follows the chain. It also takes anactivity_id, aStreamRefor astream_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, andscope="activity"gives a workflow's activity its own streams. Adefactivity is told it needsasync def.StreamProvidergainsget_activity_stream_handle(). The memory provider hosts activity owners and standalone streams, with the policy trimmed on every append andclose()sealing.Client.connectandClientConnectConfigtakestream_provider.hosts_standalone_streams.streams_demo/shows each context asking for its stream.CHANGELOG.mdhas the entry.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 lintanduv run pytest tests/streams tests/test_client.py, plus the plugin, worker and client-exports suites.test_stream_accessors.pyproves a worker inherits the client's provider and carries a ref through a workflow argument and result.test_activity_streams.pycovers the three accessor arms and a retry that surfacesSUPERSEDED.git diffagainst #3's head confirms the tie-in merge is empty.