Conversation
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.
This was referenced Sep 25, 2026
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.
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
NativeStreams, the native provider, with one owned server stream per topic.What changed?
_WorkflowInstanceImplbehind 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.temporalio/streams/providers/native.pyholdsNativeStreams. A cursor isnative:<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:BEGINNINGto earliest,ENDto tail,last=Nto last N, a cursor to an offset.DEFAULT_TOPICis the server's default stream.get_activity_stream_handle()pins one activity execution. Standalone streams,ref()andclose()go overCreateStreamandCloseStream. An owned handle refusesclose().temporalio/client_stream.pybuilds its channel from the client'sConnectConfig: target, TLS, API key, headers, keep-alive and proxy. Calls retry on sdk-core's codes with the defaultRetryConfig.translate_errormaps the server's reason tokens toStreamProducerError,StreamCursorErrorandStreamClosedError. A policy mismatch on create is aValueError.temporal.io/content-hash, andCommandAwarePayloadVisitorleaves that key alone.temporalio/contrib/server_streams/is the Workflow Streams API over the same log.tests/streams/test_stream_channel_e2e.pyfollows a native stream through the channelstream_channelderives: 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 carriesneeds_stream_channel_serverand 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, thetest_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 andtest_server_streams.pysuites 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.