Conversation
2 of 3 tasks
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 16, 2026 01:57
32378c0 to
707dbbb
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 16, 2026 01:57
810282a to
6773de6
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 16, 2026 22:08
707dbbb to
0ecca73
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 16, 2026 22:15
6773de6 to
38b7e69
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 16, 2026 23:22
0ecca73 to
7b12920
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 16, 2026 23:23
61fdb67 to
86f0e81
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 17, 2026 02:51
7b12920 to
17ed5b1
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 17, 2026 02:51
86f0e81 to
4bfe42e
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 18, 2026 23:51
17ed5b1 to
72efaad
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 19, 2026 00:03
7d5d47c to
18196e5
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 19, 2026 00:10
72efaad to
0c700b1
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 19, 2026 00:10
18196e5 to
db48755
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 19, 2026 00:13
0c700b1 to
c2fe72e
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 19, 2026 00:13
db48755 to
ecbe7b7
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 21, 2026 09:06
c2fe72e to
a5d0fc5
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 21, 2026 09:15
ecbe7b7 to
2db64db
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 21, 2026 10:47
a5d0fc5 to
4a0f7f8
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
2 times, most recently
from
September 21, 2026 11:27
16d6936 to
c055e27
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 21, 2026 17:27
61a91a2 to
f4e1c40
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-nexus
branch
from
September 21, 2026 17:33
c055e27 to
8e46845
Compare
moedash
force-pushed
the
moe/AI-198-streams-provider-workflow-streams
branch
from
September 21, 2026 17:53
f4e1c40 to
1802e8c
Compare
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.
…AI-198-streams-provider-nexus
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
…AI-198-streams-provider-nexus
…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.
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
NexusStreams has no open method, so pydoctor found no link target and the docs build failed.
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
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.
…AI-198-streams-provider-nexus
…AI-198-streams-provider-nexus
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 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.pyhasNexusStreams, the outside half of a provider over the endpoint. It also hasTemporalStreamsHandler, which serves the endpoint by fronting a storage provider's handles.temporal_streams.nexusrpc.yamlis the contract. An append carries serializedStreamRecordprotos and the producer's batch index. A read long-polls, honorslast_n, and says when the store ended it.scripts/gen_streams_nexus_api.pygenerates_nexus_generated/with stocknexgenand refuses any version but the pin. Regeneration is idempotent, andcheck-protosguards drift.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 withStreamUnsupportedError.client.get_stream_handle(ref)opens a ref through the front, andNexusStreamHandle.ref()hands one back. So an operation can return a stream as data.StreamErroror anRPCError.subscription_idle, a minute by default.stream_consumer_operation(consume, initial=..., listener_url=...)buildsStreamConsumerOperation, an asynchronous operation handler whose input is aStreamRef. 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 (closedin the notification's metadata, aFINISHrecord, 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.pyis the standalone handler for the demo: an aiohttp process that serves Nexus start and cancel over HTTP for anexusrpchandler, plus the/deliveriesroute the channel posts to. One command runs it and can register an external endpoint pointing at itself.channel_foristemporalio.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 aChannelAddressor a(channel, Execution | None)pair, and the registration and unregistration carry the owner of a linked channel as itsexecution. A store with its own rule passeschannel_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 lintanduv run pytest tests/streams. The live lane setsSTREAMS_LIVE=nexus,TEMPORAL_ADDRESSandTEMPORAL_HTTP, then runstests/streams/test_nexus_provider.pythrough 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 aStreamRefthe client then reads.tests/streams/test_nexus_consumer.pydrives 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 aFINISHcomplete 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:portagainst a server with channels andTEMPORAL_HTTPfor 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 fromdescribe_channelafter 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 namestemporal://systemas the callback. A native case waits behindneeds_stream_channel_serverfor a server whose streams notify their own channel.