Skip to content

Added the stream interface on the repaired external base. - #16

Closed
moedash wants to merge 39 commits into
moe/AI-198-py-08-external-repairsfrom
moe/AI-198-py-09-external-interface
Closed

moedash wants to merge 39 commits into
moe/AI-198-py-08-external-repairsfrom
moe/AI-198-py-09-external-interface

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR places the temporalio.streams interface on the repaired external base.

What changed?

  • temporalio/streams holds the record, typed topics, cursors, StreamRef, the StreamError family, the two provider protocols and ProviderPlugin. MemoryStreams is the memory provider, the reference implementation.
  • workflow.stream_reader() and workflow.stream_writer() reach the provider the way payload_converter() does. The worker brackets the workflow function with the provider's lifecycle hooks.
  • activity.stream_handle(), Client.get_stream_handle() and stream_provider in ClientConnectConfig give every context the same way in. scope= and activity_id= reach an activity's own streams by a static rule.
  • A read starts at BEGINNING, END or the last N records, on both accessors.
  • client.create_stream(), close() and StreamRef cover standalone streams, and MemoryStreams hosts them.
  • temporalio/streams/_body.py runs bodies through the data converter and stamps a plaintext content hash before the codec.
  • A publish from a query handler is refused at the call. A producer attempt that goes backwards is logged as a reordering, not passed on as current data.
  • The conformance suite runs on memory. It has cases for each read start, for a publish from the @workflow.init constructor, and for a reader woken through the notification channel, one through the independent channel and one through the channel linked to the reader's workflow, behind the wakes_by_notification and wakes_by_linked_notification capabilities the memory provider lacks because its reader polls on a timer. The independent case checks that the subscribed event follows a marker, so the task that opened the reader stayed retained until it left. The channel cases carry the needs_channel_server and needs_linked_server markers, so each runs only against a server named with -E host:port that serves its kind.
  • Three cases for a reader leaving its channel: closed short of FINISH, closed at FINISH, and opened and closed inside the task that opens another. The first two find one unsubscribed event, naming the subscribed one, right after the marker of the leaving completion and before the workflow's own command, and no notification after it. The third finds a subscription for the reader that stayed and none for the one that came and went. They carry wakes_by_notification and needs_unsubscribe_server, so they run on a storage provider against a server that accepts the unsubscribe command.

This layer has no Redis. #17 implements the provider against this surface.

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

Why?

This is the interface chain (#10 to #13) placed on Max's base instead of main. The trees differ, so the diff isn't identical, but the surface and the reasoning are the same. Reading #10 to #13 first turns this layer into a diff against them. The constructor case keeps every provider honest about a publish before run starts, which the harness relies on. The channel case keeps a storage provider honest about how an outside append reaches a parked reader.

How did you test it?

uv run poe lint is clean.

tests/streams passes on the memory provider with TEMPORAL_TEST_REDIS_URL set. Two cases are expected failures on memory: they're the replay rules a storage provider turns into passes. The channel cases and the leaving cases skip on memory and without a server of their kind. On Redis, through #17, they run against the server of moedash/temporal PR #17 in both modes: the linked case passes with the linked kind on, and the independent channel case and the three leaving cases pass with channel.linkedKindEnabled off. tests/contrib/external_workflow_streams passes against a dev server.

  • Unit Tests
  • Staging
  • End to End Tests

One stream interface a workflow reads, decides on and writes, with the
provider registered once as a plugin on the client and every context asking
for its stream the same way. This is the same surface as the interface chain,
placed on Max's external base so the Redis provider above it has something to
implement.
A parked Run has no open Workflow Task, so the server dispatches a query on a
task of its own and Core allows nothing beside the answer there. The instance
still reported the registered wait set, Core refused the completion, and the
query timed out.
# Conflicts:
#	temporalio/worker/_workflow_instance.py
A query is answered on a task that carries the answer alone, so the record
could only ever be dropped on the way out. The call says so instead.
Attempts only rise on one producer, so a lower one means the store reordered
two generations and the older half would render as the current answer.
Mirrors the accessor layer of the interface chain on the external base, so
the Redis provider above it can hold a stream an activity owns. A standalone
activity now reaches its own stream, which used to raise, and scope="activity"
gives a workflow's activity its own streams without changing what a bare call
reaches. The rule never probes, because a stream is created by its first write
and probing would split attempts across two streams.
Covers all three arms of the rule, a retry writing to the same stream, and
the misaddressed calls, on the memory provider. A storage provider adds
itself to the setups behind its own STREAMS_LIVE gate.
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.

(cherry picked from commit ec46a74)
A workflow could only start from BEGINNING or a cursor, so it had no way to follow from now. The provider resolves the start outside the workflow and records it, so replay reproduces it; last= is passed only when given, which keeps older providers working.

(cherry picked from commit a3edfbb)
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.

(cherry picked from commit 50b278b)
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.

(cherry picked from commit ca0515c)
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.

(cherry picked from commit 5d600c4)
A StreamRef in place of the workflow id opens the stream it names on whatever provider the client carries, and its topic becomes the handle's default. A standalone stream is created with client.create_stream and reached by its stream_id.

(cherry picked from commit 61a58f7)
A standalone stream lives in the provider with its policy and a sealed flag. The policy trims on append, a seal wakes parked readers so their reads end, and a stream id that was never created is refused at the call, so the conformance suite can run the standalone cases without a server.

(cherry picked from commit a95c913)
The conformance suite opens a ref, creates and seals a standalone stream and checks its retention policy on every provider that hosts one. The accessor tests carry a ref through a workflow argument and result and open it from the activity and the client.

(cherry picked from commit 5287f6e)
This chain has no default topic, so a ref may name the owner alone and the
handle opened from a ref is a public `RefHandle` whose calls may leave the
topic out; the two accessors return it for a ref. The memory topic drops a
parked waiter on cancellation, which the ported provider test checks, and the
plugin-registry refusal test targets a registry this chain does not have.
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.

(cherry picked from commit 440a18d)
…e's client.

A provider that cannot bound a standalone stream by bytes is now expected to refuse the policy, and one that trims by age only after close skips that check. The converter cases derive their client from the provider's own, so a live provider is not asked to reach the fixture's dev server.

(cherry picked from commit f5a7e1b)
moedash added 13 commits October 1, 2026 20:02
…9-external-interface

# Conflicts:
#	temporalio/workflow/__init__.py
#	temporalio/workflow/_context.py
Gated by a wakes_by_notification capability, which the memory provider lacks since its reader polls on a timer.
The marker skips it unless -E names a server, the way the main chain gates its live channel cases; the memory provider skips it by capability as well.
A reader closed short of FINISH and one closed at FINISH both end the run's
subscription on the completion that leaves, after the progress marker, and a
reader opened and closed inside one task never subscribes. The cases need a
server that accepts the unsubscribe command, so they carry a marker of their
own.
…task.

The channel case now finds the subscribed event after a marker, which is
where Core puts it, so the task that opened the reader stayed retained.
The notified scheduled event names the owner as a
`temporal.api.common.v1.Execution`, so the case checks the type and the
business id through the public helper.
@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

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant