Skip to content

Added the native stream provider. - #53

Open
moedash wants to merge 5 commits into
moe/AI-198-st-py-2-native-wirefrom
moe/AI-198-st-py-3-native-provider
Open

moedash wants to merge 5 commits into
moe/AI-198-st-py-2-native-wirefrom
moe/AI-198-st-py-3-native-provider

Conversation

@moedash

@moedash moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner

This PR adds NativeStreams, which puts a server-side stream behind the stream interface, along with what it needs in the workflow runtime and the stream client.

What changed?

  • The stream client shares one gRPC channel per namespace, mirroring the client's connection settings. It retries transient errors the way the client does, and runs record bodies through the data converter.
  • The workflow runtime gets the native stream commands. A workflow appends with a command its task commits, and reads the ranges the server delivers on its tasks. That keeps native streams replay-safe. The command-aware visitor keys stream records by content hash.
  • contrib.server_streams gives the shipped Workflow Streams surface over a server-side stream.
  • temporalio.streams.providers.native.NativeStreams serves one owned stream per topic for workflows, activities and standalone streams. It joins the conformance suite when STREAMS_LIVE=native.
  • The needs_stream_channel_server marker is for cases where a native stream notifies its channel. The native Nexus consumer case gets it too.

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

Why?

Memory and Workflow Streams cover small streams. A server-side stream is the store meant for large ones, with retention and trimming handled by the server. Putting it behind the same interface means workflow, activity and client code don't change when an app moves onto it.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

poe lint is clean. The unit cases pass on the dev server the fixtures start. Against a local server built from the stream server PRs, these pass with STREAMS_LIVE=native:

  • the native workflow, payload and end-to-end cases
  • contrib.server_streams
  • the stream client suites, covering connection, errors, retry and sharing
  • the whole streams suite, with the native provider in the conformance run
  • the three cases where a native stream notifies its channel

One case failed once in the first full run and passed in three reruns, and I didn't catch which. The skips are capability skips and the native Nexus consumer case, which stays gated until the union.

One channel per namespace mirrors the client's connection, transient errors retry like the client's, and record bodies go through the data converter.
A workflow appends with a command its task commits and reads the ranges the server delivers on its tasks, so a native stream is replay-safe.
The Workflow Streams surface, backed by a server-side stream instead of the workflow's History.
NativeStreams puts a server-side stream behind the stream interface, one owned stream per topic, and joins the conformance suite under STREAMS_LIVE=native.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

skip-changelog Changelog entry rides another PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant