Conversation
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.
This was referenced Sep 25, 2026
# 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.
…9-external-interface
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)
(cherry picked from commit 2e42984)
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)
…tream. (cherry picked from commit 06cfa01)
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)
…s' into moe/AI-198-py-09-external-interface
…s' into moe/AI-198-py-09-external-interface
…s' into moe/AI-198-py-09-external-interface
2 of 3 tasks
…9-external-interface
…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.
…9-external-interface
…9-external-interface
…9-external-interface
…9-external-interface
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.
…9-external-interface
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.
…9-external-interface
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 places the
temporalio.streamsinterface on the repaired external base.What changed?
temporalio/streamsholds the record, typed topics, cursors,StreamRef, theStreamErrorfamily, the two provider protocols andProviderPlugin.MemoryStreamsis the memory provider, the reference implementation.workflow.stream_reader()andworkflow.stream_writer()reach the provider the waypayload_converter()does. The worker brackets the workflow function with the provider's lifecycle hooks.activity.stream_handle(),Client.get_stream_handle()andstream_providerinClientConnectConfiggive every context the same way in.scope=andactivity_id=reach an activity's own streams by a static rule.BEGINNING,ENDor the lastNrecords, on both accessors.client.create_stream(),close()andStreamRefcover standalone streams, andMemoryStreamshosts them.temporalio/streams/_body.pyruns bodies through the data converter and stamps a plaintext content hash before the codec.@workflow.initconstructor, 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 thewakes_by_notificationandwakes_by_linked_notificationcapabilities 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 theneeds_channel_serverandneeds_linked_servermarkers, so each runs only against a server named with-E host:portthat serves its kind.FINISH, closed atFINISH, 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 carrywakes_by_notificationandneeds_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
runstarts, 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 lintis clean.tests/streamspasses on the memory provider withTEMPORAL_TEST_REDIS_URLset. 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 withchannel.linkedKindEnabledoff.tests/contrib/external_workflow_streamspasses against a dev server.