Skip to content

Added the native provider over one owned stream per topic. - #9

Closed
moedash wants to merge 51 commits into
moe/AI-198-py-05-native-wirefrom
moe/AI-198-py-06-native-provider
Closed

moedash wants to merge 51 commits into
moe/AI-198-py-05-native-wirefrom
moe/AI-198-py-06-native-provider

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR adds NativeStreams, the native provider, with one owned server stream per topic.

What changed?

  • The workflow half lives in _WorkflowInstanceImpl behind private stream helpers. A task's publishes on one stream become one command and one History event, ahead of any command that ends the run. A failed task's publishes never become commands.
  • Delivered ranges go into a per-stream buffer that checks continuity. A repeated, skipped or mis-sized range fails the task, and the message names the way out. The worker also splits publishes to fit the server's per-record and per-batch limits.
  • temporalio/streams/providers/native.py holds NativeStreams. A cursor is native:<run_id>:<offset>. A handle without a run id follows continue-as-new, and one with a run id stays on that run. Read starts map onto the server: BEGINNING to earliest, END to tail, last=N to last N, a cursor to an offset. DEFAULT_TOPIC is the server's default stream.
  • get_activity_stream_handle() pins one activity execution. Standalone streams, ref() and close() go over CreateStream and CloseStream. An owned handle refuses close().
  • temporalio/client_stream.py builds its channel from the client's ConnectConfig: target, TLS, API key, headers, keep-alive and proxy. Calls retry on sdk-core's codes with the default RetryConfig.
  • translate_error maps the server's reason tokens to StreamProducerError, StreamCursorError and StreamClosedError. A policy mismatch on create is a ValueError.
  • Bodies go through the client's data converter, codec and external storage. Each record carries the SHA-256 of its plaintext body under temporal.io/content-hash, and CommandAwarePayloadVisitor leaves that key alone.
  • temporalio/contrib/server_streams/ is the Workflow Streams API over the same log.
  • tests/streams/test_stream_channel_e2e.py follows a native stream through the channel stream_channel derives: a standalone stream through its independent channel and a workflow's stream through the channel linked to its owner, one notification per append and one for the close, with a callback registered on the derived name. It carries needs_stream_channel_server and passed against the server layer on which streams drive channels. A standalone stream hands its channel the latest change through one outstanding task, so the case waits for each append's notification before the next.

Replay of a consuming workflow is left to #14.

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

Why?

The memory provider can't show the two rules that matter for a workflow. A publish commits with its Workflow Task, and a read is an observation the server re-supplies on replay. Both live in the server's commit, so only a server can test them. The content hash keeps a codec with a fresh nonce per call from turning a retried append into a divergent write.

How did you test it?

The lint set from #8 ran on a cold .mypy_cache. Without a server, test_workflow_stream.py, the test_client_stream_* modules and the native handle tests ran, with an in-process gRPC server recording the headers it gets. Live, on a server from moedash/temporal#17, the conformance, standalone, activity, e2e, payload and test_server_streams.py suites ran. They cover a publish failing with its task, cold-cache replay and starts, batching, continue-as-new, a storage driver on both halves and a nonce codec on a retry, which needs a server that dedupes on the hash.

  • Unit Tests
  • Staging
  • End to End Tests

A task's publishes on one stream are held until it completes, so they become one command and one History event however many there were, and a failed task publishes nothing. Delivered ranges are buffered because the server records a range as consumed and never sends it again.
The workflow subscribes to the topic by name and publishes with the batched command; outside code appends and long-polls the same stream, following the run chain by cursor unless pinned.
The same API as temporalio.contrib.workflow_streams on a Temporal-owned log, so an application swaps the import and keeps its code. The LangGraph docstring qualifies its stream references, because this module gives both names a second definition pydoctor cannot resolve.
The module says an application swaps the import and keeps its code, but the append carried no identity and took the batch off the buffer before the await, so a retried Activity wrote twice and a failed append lost what it held.
One channel is shared per loop, target and namespace, so closing a provider took out channels another one in the same process was still reading through.
The check fails the Workflow Task and keeps failing it, because the range is recorded as consumed and is never sent again, so the message has to say what a reader can do about it. The copied batch limits now say what breaks if the server's differ.
There is no unsubscribe command, so the server keeps delivering for the life of the run and the buffer grew for a reader nobody would read again. The raw stream API that takes stream ids and protos is private now: the typed surface beside it is the one an application should reach for.
Zero on the wire says a producer does not number its records, and this one
does, so it cannot start there.
Core names these fields after what they carry, so the writes and the helpers that feed them read the same way.
The server resolves an unnamed stream to the same name as DEFAULT_TOPIC, so the native handle takes no topic and sends the name explicitly, which keeps a record's topic equal to its stream's name.
A standalone activity is its own owner on the server and a workflow's activity is reached through its workflow; the handle pins one activity execution and never follows a chain, because a retry writes to the same stream.
BEGINNING now asks the server for the earliest record rather than offset zero, which a stream with a raised floor refuses. END and last=N become the tail and last-N positions, resolved by the server on the first poll or at subscription and recorded, so replay never resolves them again.
The stream client has a channel of its own, so Core's retry never covered
it and a reader under frontend rate limiting raised mid-read. An append
without a producer id goes again only on a refusal the server sent before
doing anything, because the server has nothing to deduplicate it against.
The stream service was reached on a bare insecure channel, so a client with
TLS, an API key, default headers or its own retry policy lost all of them on
that path. Connection reads the client's ConnectConfig and opens a grpcio
channel with the same target, TLS material, bearer header, headers, keep-alive
and proxy, and the shared client retries under the client's retry_config.
A divergent or stale producer repeat is a StreamProducerError and a poll below
the retention floor is a StreamCursorError, the fix the spec lists as pending.
The server names each refusal with a reason token at the front of the message,
so one function matches the status code and the token, keeps the phrases an
older server sends, and leaves every other failure an RPCError.
The server deduplicated a producer's repeat on the encoded bytes, so a codec
with a fresh nonce per call made every retry look divergent. The provider now
puts the hex SHA-256 of the converted body under temporal.io/content-hash in
the record's metadata, before the codec on the outside path and before the
worker's payload pass on the workflow path; the server dedupes on it when
present. The live nonce-codec case needs a server with that dedupe.
The outside half applied the payload codec alone, so an ExternalStorage
driver on the client never saw a stream body and a body the worker offloaded
could not be read back. Encode and decode now run the codec and the external
store in the worker's order, an offloaded body is stored under the owning
execution, and the fingerprint is taken before either.
The payload visitor walks a stream record's metadata payloads like any other,
so a codec encrypted the hash a workflow publish declared and an external
store could offload it. The server reads that value as sent and refuses one
that is not hex, which would have failed the workflow's own append. The key
and the digest now live in the streams wire module, where both halves and the
worker read them.
… moe/AI-198-py-06-native-provider

# Conflicts:
#	tests/streams/test_streams_conformance.py
A stream with an id of its own is created and sealed over the existing
CreateStream and CloseStream calls, read through a handle whose cursors name
the stream where an owned stream's name a run, and a read on an id nobody has
created yet relies on the server's parked poll. Every native handle answers
ref(), the provider opens a StreamRef, and a create of an id that exists is
answered from the stream's own policy since the server does not say
ALREADY_EXISTS. The content hash and the body encoding now come from the
shared streams helper the interface chain added, and an append on a sealed
stream is a StreamClosedError.
The server is giving an append on a sealed stream a reason token of its own,
so the mapping takes the token beside the three it already knew and keeps the
whole-message match for a server built before it.
The server is answering ALREADY_EXISTS for an id that exists with the same
policy and a STREAM_POLICY_MISMATCH refusal for one that differs, so the
provider returns the handle on the first and lets the second arrive as the
ValueError the contract names. The describe comparison stays for a server that
answers both with a generic failure.
The server's lifecycle now carries max_bytes, refuses an append that would
cross it, reclaims an open standalone stream's records by age, and answers a
create of an existing id with ALREADY_EXISTS or a STREAM_POLICY_MISMATCH
refusal. The provider passes the cap through, drops the describe comparison,
and the native conformance setup declares the byte cap as a refusal rather
than a trim and waits for the store's age reclaim.
This layer maps the STREAM_CURSOR_BELOW_FLOOR token to StreamCursorError, which is not an RPCError, so the live case written on the wire layer no longer matched the exception it got.
The two cases poll the channel stream_channel derives, one standalone and one owned by a workflow, and register a callback on the derived name. They need a server on which streams drive channels and skip otherwise.
A standalone stream hands its channel the latest change through one outstanding task, so a burst of appends arrives as its newest change. The case now observes every change rather than a fold.
…ive-provider

# Conflicts:
#	tests/streams/test_activity_streams.py
#	tests/streams/test_streams_conformance.py
…ve cases.

The native e2e cases read the owner as an execution, and a new live case
polls the channel linked to a standalone activity, which waits for the
server layer that addresses a linked channel by execution.
@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

skip-changelog Changelog entry rides another PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant