Skip to content

Added a Nexus operation that consumes a stream through its channel. - #34

Open
moedash wants to merge 4 commits into
moe/AI-198-if-py-8-nexus-frontfrom
moe/AI-198-if-py-9-nexus-consumer
Open

moedash wants to merge 4 commits into
moe/AI-198-if-py-8-nexus-frontfrom
moe/AI-198-if-py-9-nexus-consumer

Conversation

@moedash

@moedash moedash commented Oct 3, 2026 •

Copy link
Copy Markdown
Owner

This PR adds stream_consumer_operation, an async Nexus operation that consumes a stream and completes when the stream closes.

What changed?

  • StreamConsumerOperation takes a StreamRef, 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_for defaults to temporalio.client.stream_channel, the same rule the server uses. Another rule can be passed for a store that notifies elsewhere.
  • nexus_consumer_service is a standalone aiohttp host that serves the operation and receives the channel's callbacks. aiohttp comes in through the new streams-nexus extra.
  • The consumer's changelog line. needs_native_provider gates 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 -

  • Unit Tests
  • Staging
  • End to End Tests

poe lint is clean, and uv lock --check accepts 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.

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.
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