Skip to content

Added the Nexus front for stream providers. - #33

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

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

Conversation

@moedash

@moedash moedash commented Oct 3, 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, 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 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.

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.

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. The handler and caller cases pass against Nexus endpoints on the dev server the fixtures start. The front also runs as a conformance case.

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