Skip to content

Added the Nexus front for stream providers. - #97

Open
moetemp wants to merge 4 commits into
moe/AI-198-s1-py-7-workflow-streams-providerfrom
moe/AI-198-s1-py-8-nexus-front
Open

moetemp wants to merge 4 commits into
moe/AI-198-s1-py-7-workflow-streams-providerfrom
moe/AI-198-s1-py-8-nexus-front

Conversation

@moetemp

@moetemp moetemp commented Oct 6, 2026

Copy link
Copy Markdown
Owner

This PR puts one Nexus endpoint in front of any stream provider, so callers reach a stream without knowing which store holds it.

What changed?

  • temporal_streams.nexusrpc.yaml defines two sync operations, append and read. The nexgen bindings in _nexus_generated come from it, and poe gen-protos regenerates them through gen-streams-nexus-api.
  • TemporalStreamsHandler runs in a worker next to any storage provider. It serves those operations through the provider's own handles. It keeps one parked read per stream ref, so an idle caller's long poll is reused rather than abandoned, and it dedupes appends by producer attempt.
  • NexusStreams is an outside-only provider that appends and reads through the endpoint. It batches, handles cursors, and runs the caller's codec before records leave the process. A read is a long poll on the handler, so a caller following a stream needs nothing else to learn of new records.
  • A record crosses as the serialized StreamRecord. A stream is addressed by a wire StreamRef.
  • pyproject.toml keeps the generated module out of mypy, pydocstyle and the docs. The changelog gets an entry.

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

Why?

An operator should be able to change the store behind a stream without touching the callers. With a Nexus endpoint in front, callers only know the endpoint, and Temporal's auth guards it. The contract is a nexgen yaml, so other SDKs can generate the same client.

The earlier series had a second Nexus PR above this one: an async operation that consumed a stream and woke on the stream's notification channel. The channel is paused per the 2026-10-05 design review, so that consumer is out of this phase and comes back with the channel. The front never depended on it: it has its own long-poll read, unchanged.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

poe lint is clean. A second run of the generators reproduces the bindings exactly. tests/streams passes with the front's cases in it, and the front passes the conformance suite as its own lane on the dev server the fixtures start. The live handler and caller cases against a Nexus endpoint ran higher in the series, where they pass.

The contract is a nexgen yaml, so any language nexgen targets gets the same append and read operations. gen-protos regenerates the bindings.
NexusStreams appends and reads through one endpoint, and TemporalStreamsHandler serves it in front of any storage provider, so callers never name the store.
The handler and caller cases on the dev server's Nexus endpoints, and the front as a conformance case.
@moetemp moetemp mentioned this pull request Oct 6, 2026
1 of 3 tasks
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