Conversation
The record, cursors, typed topic definitions and the error family are the vocabulary; the two provider protocols say what a store implements and which half runs on the workflow thread. Nothing here does I/O, so workflow code can import all of it.
`ProviderPlugin` is the plugin half every provider in this tree shares, so one registration reaches the client, the worker and the replayer through their `stream_provider` option. The memory provider exists so the conformance cases can exercise the whole outside surface without a store, and it documents in one file what a provider owes.
This was referenced Sep 25, 2026
The conformance file is what a new provider has to answer, so the cases that unit-test the shared wire, policy and id helpers belong beside them rather than in it.
There is one slot and it decides where a client, worker or replayer reads and publishes, so a plugin that overwrites another provider's registration silently is a configuration bug with nothing to see.
A repeat at a sequence the store already holds is only a retry when it carries the same content; answering a different one with the original position drops a record the writer meant to send.
The interface says a caller that stops early can aclose() the generator and get back what the provider parked; a cancelled wait left its waiter on the topic for the life of the process.
One provider serves a client, the workers built from it and every handle opened outside them, so no single one of those can close it without cutting off the others.
Two of the three providers here have no way to tell their store a reader has gone, so a flat promise of an immediate release was one the interface could not keep.
The wire no longer carries a sentinel for an unnumbered record, so the memory producer starts its own numbering at one and leaves zero to the producers that never set it.
The server already resolves an unnamed workflow stream to "output", so the SDK uses the same name and a no-topic call on any provider lands on the stream the server would pick.
BEGINNING is documented as the oldest record a stream still retains, which on a truncated stream is not offset zero. END follows the tail from when a read starts and last=N starts at the newest records; after= stays the only way to resume.
A provider owes a record body the codec and external storage the SDK gives every payload. The shared helper takes the retry fingerprint before either runs and stamps the plaintext hash on the record, so a nondeterministic codec cannot turn a retry into a divergent write and the store can compare retries without the plaintext.
A stream with an id of its own and no owner is created on purpose with a retention policy and sealed on purpose. The memory provider refuses both calls for now, and a workflow's handle refuses close, since its stream ends with the workflow.
A handle is bound to its client and provider, so a stream crosses a process boundary as its owner and topic in plain data, with no cursor and no provider name. The default converter carries it as JSON, so it can be a workflow argument, an activity result or a Nexus operation input or result.
A store may bound a standalone stream by record count and age but not by bytes, or apply retention only once the stream is closed. The two flags let such a provider say so and have the suite hold it to what it declares.
A store may keep a standalone stream under its byte bound by refusing the append that would cross it rather than by dropping its oldest records. The flag lets a provider say which it does, and memory declares that it trims.
The workflow and activity accessors and Client.get_stream_handle arrive on a later layer, so pydoctor finds no target for them here and the docs build fails. They stay named as literal text.
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 the
temporalio.streamspackage and the memory provider, the reference provider.What changed?
temporalio/streams/_record.py,_topic.py,_errors.pyand_ref.pyhold the vocabulary: the record and its kinds,CursorwithBEGINNINGandEND, typed topics,DEFAULT_TOPIC, theStreamErrorfamily andStreamRef. AStreamRefnames an owner and a topic, with no cursor and no provider name, so it travels as plain JSON._provider.pystates a provider in two halves.WorkflowStreamProviderruns on the workflow thread and does no I/O.StreamProvideris the process half that hands outStreamHandles. It also declares standalone streams andStreamHandle.close()._wire.py,_policy.py,_ids.pyand_body.pyare the helpers every provider shares. They cover the proto mapping, synthesized supersession, collision-free store keys, and the body codec path with the plaintext hash undertemporal.io/content-hash.providers/__init__.pyhasProviderPlugin. It sets itself as thestream_providerof the client, the worker and the replayer, which is why those three configs gain that key. A second provider on the same slot is refused.providers/memory.pyis the reference provider. Its docstring says what it doesn't keep: it's not replay-safe and a publish is visible before the task is accepted.tests/streams/test_streams_conformance.pyis the conformance core, parametrised over providers. A provider declares what it lacks, likedetects_divergent_retriesortruncates, and those cases skip with a reason.The memory provider refuses standalone streams on this layer. #13 hosts them.
Part of AI-198 (epic AI-37).
Why?
The two-halves split keeps I/O out of everything workflow code imports. It also lets a language that bundles workflow code separately ship the halves in two packages. The memory provider lets the conformance cases exercise the whole outside surface with no store and no server. Its outside behavior is faithful: producer identity, retry dedupe, start positions, supersession and cursors.
How did you test it?
I ran
uv run poe lintanduv run pytest tests/streams. That covers the conformance cases on memory, including each read start, a truncated topic and a divergent retry. It also covers the shared helpers intest_streams_internals.py. Workflow code can't reach a stream until #12, so nothing here runs inside a workflow.