Skip to content

Added the temporalio.streams package and the memory provider. - #11

Closed
moedash wants to merge 27 commits into
moe/AI-198-py-01-protosfrom
moe/AI-198-py-02-streams-package
Closed

moedash wants to merge 27 commits into
moe/AI-198-py-01-protosfrom
moe/AI-198-py-02-streams-package

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR adds the temporalio.streams package and the memory provider, the reference provider.

What changed?

  • temporalio/streams/_record.py, _topic.py, _errors.py and _ref.py hold the vocabulary: the record and its kinds, Cursor with BEGINNING and END, typed topics, DEFAULT_TOPIC, the StreamError family and StreamRef. A StreamRef names an owner and a topic, with no cursor and no provider name, so it travels as plain JSON.
  • _provider.py states a provider in two halves. WorkflowStreamProvider runs on the workflow thread and does no I/O. StreamProvider is the process half that hands out StreamHandles. It also declares standalone streams and StreamHandle.close().
  • _wire.py, _policy.py, _ids.py and _body.py are 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 under temporal.io/content-hash.
  • providers/__init__.py has ProviderPlugin. It sets itself as the stream_provider of 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.py is 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.py is the conformance core, parametrised over providers. A provider declares what it lacks, like detects_divergent_retries or truncates, 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 lint and uv 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 in test_streams_internals.py. Workflow code can't reach a stream until #12, so nothing here runs inside a workflow.

  • Unit Tests
  • Staging
  • End to End Tests

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