Skip to content

Added the Nexus front that hides the store behind one endpoint. - #5

Closed
moedash wants to merge 58 commits into
moe/AI-198-streams-provider-workflow-streamsfrom
moe/AI-198-streams-provider-nexus
Closed

moedash wants to merge 58 commits into
moe/AI-198-streams-provider-workflow-streamsfrom
moe/AI-198-streams-provider-nexus

Conversation

@moedash

@moedash moedash commented Sep 15, 2026 •

Copy link
Copy Markdown
Owner

This PR adds the Nexus front, which puts a stream store behind one Temporal-authenticated endpoint, and the Nexus operation that consumes a stream through its notification channel.

What changed?

  • temporalio/streams/providers/nexus.py has NexusStreams, the outside half of a provider over the endpoint. It also has TemporalStreamsHandler, which serves the endpoint by fronting a storage provider's handles.
  • temporal_streams.nexusrpc.yaml is the contract. An append carries serialized StreamRecord protos and the producer's batch index. A read long-polls, honors last_n, and says when the store ended it. scripts/gen_streams_nexus_api.py generates _nexus_generated/ with stock nexgen and refuses any version but the pin. Regeneration is idempotent, and check-protos guards drift.
  • Both operations address a stream by a StreamRef. Its optional members are nullable on the wire. The handler maps the ref onto the store's accessor and refuses an owner the store can't host with StreamUnsupportedError.
  • On the caller side, client.get_stream_handle(ref) opens a ref through the front, and NexusStreamHandle.ref() hands one back. So an operation can return a stream as data.
  • The caller applies the codec before bytes leave the process. It commits a batch index only after the endpoint answers. A failure reaches it as the store's StreamError or an RPCError.
  • The handler holds one batch and one read at a time per address, in weakly held lock maps. A parked read is released after subscription_idle, a minute by default.
  • stream_consumer_operation(consume, initial=..., listener_url=...) builds StreamConsumerOperation, an asynchronous operation handler whose input is a StreamRef. On start it registers a callback listener on the stream's channel with the operation token in a header and reads the stream from its start through the client's provider. Each delivery the server posts leads to one read from the cursor to the head, a delivery with nothing new reads nothing, and the close (closed in the notification's metadata, a FINISH record, or the store ending the read) completes the operation through the caller's completion callback with the folded value. Cancel unregisters and reports canceled. The cursor lives in the process, keyed by the token.
  • nexus_consumer_service.py is the standalone handler for the demo: an aiohttp process that serves Nexus start and cancel over HTTP for a nexusrpc handler, plus the /deliveries route the channel posts to. One command runs it and can register an external endpoint pointing at itself.
  • The default channel_for is temporalio.client.stream_channel, the server's names for native streams (stream/<stream id> independent, stream/<topic> linked to the owning workflow or standalone activity, stream/<activity id>/<topic> for a workflow activity's). A rule answers with a ChannelAddress or a (channel, Execution | None) pair, and the registration and unregistration carry the owner of a linked channel as its execution. A store with its own rule passes channel_for.

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

Why?

Producers and consumers should talk to one endpoint while the store behind it swaps by operator action. The IDL is how this surface reaches other SDK languages, since a client generated from the same file speaks the same wire. The ref on the contract is the spec's Nexus decision 1, option (b). The consumer operation is option 2 of the Nexus question: the operation registers as the stream's consumer and the server tells it, so the Nexus design and the channel design share one delivery path instead of the handler polling the front.

How did you test it?

I ran uv run poe lint and uv run pytest tests/streams. The live lane sets STREAMS_LIVE=nexus, TEMPORAL_ADDRESS and TEMPORAL_HTTP, then runs tests/streams/test_nexus_provider.py through the HTTP ingress of a server built from the series. The tests register and delete their own endpoint. They cover cursor resume across the hop, ciphertext at rest with a no-codec control read, and an operation that returns a StreamRef the client then reads.

tests/streams/test_nexus_consumer.py drives the consumer's state machine with a client stand-in over the memory store: the registration and the first read, deliveries read from the cursor once and in order, a folded burst costs one read, a delivery with nothing new is a no-op, the close and a FINISH complete through the callback with the right headers, cancel unregisters, an unknown token is ignored, a consume failure fails the operation and a read failure leaves it for the retry. The standalone service is checked against the Nexus HTTP protocol with aiohttp's test client. The live cases run with -E host:port against a server with channels and TEMPORAL_HTTP for its ingress: a caller workflow starts the operation on an external endpoint, the producer appends to the memory store behind the front and notifies the stream's channel by hand (this layer carries no store that notifies on its own), records before and after the registration are consumed once and in order, the close completes the workflow with the values, and the listener is gone from describe_channel after completion and after a cancel. The worker-hosted variant is tried once: the frontend accepts the channel's post as a start of a companion operation, and the completion is left as a follow-up because a worker's start context names temporal://system as the callback. A native case waits behind needs_stream_channel_server for a server whose streams notify their own channel.

  • Unit Tests
  • Staging
  • End to End Tests

@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 32378c0 to 707dbbb Compare September 16, 2026 01:57
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from 810282a to 6773de6 Compare September 16, 2026 01:57
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 707dbbb to 0ecca73 Compare September 16, 2026 22:08
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from 6773de6 to 38b7e69 Compare September 16, 2026 22:15
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 0ecca73 to 7b12920 Compare September 16, 2026 23:22
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from 61fdb67 to 86f0e81 Compare September 16, 2026 23:23
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 7b12920 to 17ed5b1 Compare September 17, 2026 02:51
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from 86f0e81 to 4bfe42e Compare September 17, 2026 02:51
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 17ed5b1 to 72efaad Compare September 18, 2026 23:51
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from 7d5d47c to 18196e5 Compare September 19, 2026 00:03
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 72efaad to 0c700b1 Compare September 19, 2026 00:10
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from 18196e5 to db48755 Compare September 19, 2026 00:10
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 0c700b1 to c2fe72e Compare September 19, 2026 00:13
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from db48755 to ecbe7b7 Compare September 19, 2026 00:13
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from c2fe72e to a5d0fc5 Compare September 21, 2026 09:06
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from ecbe7b7 to 2db64db Compare September 21, 2026 09:15
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from a5d0fc5 to 4a0f7f8 Compare September 21, 2026 10:47
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch 2 times, most recently from 16d6936 to c055e27 Compare September 21, 2026 11:27
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from 61a91a2 to f4e1c40 Compare September 21, 2026 17:27
@moedash
moedash force-pushed the moe/AI-198-streams-provider-nexus branch from c055e27 to 8e46845 Compare September 21, 2026 17:33
@moedash
moedash force-pushed the moe/AI-198-streams-provider-workflow-streams branch from f4e1c40 to 1802e8c Compare September 21, 2026 17:53
The changelog checkpoint requires an entry for any user-facing change.
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.
The handler closes the delegate's subscription as soon as it answers, and
the interface only promises an iterator, so the close is guarded by what
the delegate actually returned.
The changelog checkpoint requires an entry for any user-facing change.
`asyncio.timeout` arrived in 3.11 and this package supports 3.10, so the
collecting loop moves into a coroutine and `wait_for` bounds it.
The stream endpoint's shape lived only in hand-written Python, so no other
language could be given the contract. This defines it once in a nexusrpc
yaml contract and generates the Python bindings from it. The generated
_definitions helper is hidden from the API docs and skipped by pydocstyle,
the same way the other generated trees are.
The store now hosts a workflow's activity streams as reserved topics, so
the case that checks a store's refusal reaching the caller under its own
class needs an owner it still refuses.
…flow-streams' into moe/AI-198-streams-provider-nexus
basedpyright refuses a cast from None to Client, and the in-process endpoint
is reached by id, so the test says outright that it opens without one.
…flow-streams' into moe/AI-198-streams-provider-nexus
The SDK writes null for every member a StreamRef does not use, and a
generated caller refuses an explicit null on a member the contract did not
make nullable, so a ref an operation returns as plain data could not cross
a generated converter.
NexusStreams has no open method, so pydoctor found no link target and the docs build failed.
The handler registers a callback listener on the stream's channel and reads
from its cursor on each delivery, so the Nexus and channel designs share one
path instead of the handler polling the front. The aiohttp service is the
standalone host the PoC demonstrates.
Unit cases drive the state machine with a client stand-in over the memory
store. Live cases put the memory store behind the front and notify the
channel by hand, since this layer has no store that notifies on its own. The
worker-hosted variant stops before the close: a worker's start context names
temporal://system as the callback, which an HTTP post cannot reach.
A standalone stream notifies stream/<stream id>, an owned one stream/<topic>
linked to the owning workflow, so the registration carries the owner. The
rule is the channel_for default and a store with its own passes it.
….com:moedash/sdk-python into moe/AI-198-streams-provider-nexus

# Conflicts:
#	tests/streams/conftest.py
channel_for answers with a ChannelAddress or a (channel, workflow_id) pair
and defaults to the main chain's rule, so the consumer and the server agree
on a native stream's channel. The native case also waits for a layer that
carries a native provider, which this one does not.
….com:moedash/sdk-python into moe/AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus

# Conflicts:
#	tests/streams/conftest.py
The helper hands the owner `stream_channel` names to the registration as
`execution`, so a stream a standalone activity owns is consumed through the
channel linked to the activity. A `channel_for` rule answers with an address
or a `(channel, execution)` pair.
@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