Conversation
It registers a callback listener on the stream's notification channel and reads on each delivery, so a caller waits on a stream without polling.
A process that is not a worker can serve the operation and receive the channel's callbacks. aiohttp comes in through the streams-nexus extra.
Unit cases with a client double, and live cases that need a channel server named with -E. The native case stays gated until the union.
This was referenced Oct 3, 2026
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
stream_consumer_operation, an async Nexus operation that consumes a stream and completes when the stream closes.What changed?
StreamConsumerOperationtakes aStreamRef, folds the stream's records with a consume function, and returns the result. On start it registers a callback listener on the stream's notification channel, then reads from the start. Each delivery triggers a read up to the head. The close completes the operation through the caller's completion callback.channel_fordefaults totemporalio.client.stream_channel, the same rule the server uses. Another rule can be passed for a store that notifies elsewhere.nexus_consumer_serviceis a standalone aiohttp host that serves the operation and receives the channel's callbacks.aiohttpcomes in through the newstreams-nexusextra.needs_native_providergates the one case that needs native streams until the union.Part of AI-198 (epic AI-37).
Why?
A caller often wants "run this when the stream is done" without holding a reader open. A Nexus operation fits: it's async, it survives restarts, and the caller's workflow just awaits it. Because the channel drives it, nobody polls. This PR is where the interface meets the channel, and it needs both below it.
How did you test it?
Link to a test plan if any -
poe lintis clean, anduv lock --checkaccepts the lock. The unit cases pass with a client double, covering start, delivery, cancel, close, failure and the default channel rule. The live cases ran against a local server built from the channel server PRs, with a Nexus callback template and local callback addresses in its dynamic config. They cover a standalone service consuming an external stream through its channel, cancel unregistering the listener, and a worker-hosted consumer reached through the frontend, and they pass. The native case stays gated until the union.