Conversation
moedash
force-pushed
the
moe/AI-198-streams-all
branch
from
September 15, 2026 23:51
d8f85b9 to
1771757
Compare
2 of 3 tasks
moedash
force-pushed
the
moe/AI-198-streams-all
branch
from
September 17, 2026 04:37
549bbdb to
307dedb
Compare
moedash
added a commit
that referenced
this pull request
Sep 22, 2026
Takes #6 at `3cf17128`: `WorkflowHistory.stream_slices` and `Replayer.fetch_stream_slices` from the native branch. No conflicts, nothing under `examples/` changed.
This was referenced Sep 25, 2026
… into moe/AI-198-streams-all # Conflicts: # temporalio/activity.py # temporalio/client/_client.py # temporalio/streams/__init__.py # temporalio/streams/_provider.py # temporalio/streams/_ref.py # temporalio/streams/providers/redis.py # temporalio/workflow/_streams.py # tests/streams/conftest.py # tests/streams/test_redis_provider.py # tests/streams/test_redis_replay.py # tests/streams/test_stream_accessors.py # tests/streams/test_streams_conformance.py # tests/streams/test_streams_internals.py # tests/streams/test_streams_workflow.py
The server's StreamLifecycle gained max_bytes beside max_items, and StreamState reports held_bytes, so the standalone policy can bound bytes. Generated with the same protobuf 3 toolchain as before, so the modules still load on the oldest runtime this package supports.
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.
…into moe/AI-198-streams-all
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.
…-workflow-runtime
A store that refuses the append crossing its byte cap is checked with a cap that holds one stored record, hash metadata included, and not two. The age subcase waits for the older record to leave BEGINNING and accepts the newer one aging out as well, since a store reclaims on its own schedule.
The SDK handles a task's Signals and Updates ahead of the workflow function, so a publish or poll Update that arrived with the first task was rejected and, on a worker with no cache, on every rebuilt instance too. The start hook now runs when the instance's loop first turns, and Workflow Streams binds the shipped stream object lazily so a workflow's own stays.
…AI-198-streams-provider-nexus
… moe/AI-198-streams-provider-workflow-streams
…s' into moe/AI-198-streams-all # Conflicts: # temporalio/worker/_workflow.py
The provider hosts no standalone streams, so there is no byte cap to refuse an append past; the flag says so next to the other three.
…AI-198-streams-provider-nexus
…s' into moe/AI-198-streams-all
A completion the server rejects, such as a closing task with a wake Signal buffered under it, leaves the batch it staged with a token no marker will name; the re-run stages the batch again, and the dead stage stayed in the log once a reader or the worker aborted it. On the one-log layout those entries count against every window and bound, so the abort now deletes them and the stage's hashes stand in for them, and the worker settles a run's pending stages against History when it evicts the run rather than leaving them for a reader to find.
…vider-workflow-streams
…AI-198-py-05-native-wire
…AI-198-streams-provider-nexus
…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.
…-native-replay # Conflicts: # temporalio/worker/_workflow.py
…98-streams-all # Conflicts: # temporalio/bridge/sdk-core # tests/streams/conftest.py # tests/streams/test_activity_streams.py # tests/streams/test_streams_conformance.py
…98-streams-all # Conflicts: # tests/streams/conftest.py # tests/streams/test_nexus_consumer.py
…98-streams-all # Conflicts: # temporalio/api/workflowservice/v1/request_response_pb2.py # temporalio/bridge/sdk-core # temporalio/client/_channel.py # temporalio/client/_client.py # temporalio/client/_impl.py # temporalio/client/_interceptor.py # temporalio/contrib/external_workflow_streams/_wake.py # temporalio/worker/_workflow.py # tests/contrib/external_workflow_streams/test_wake.py
The union Core addresses a linked channel by execution in both lineages, so the api and bridge protos regenerate from one tree again.
Both lineages added the same helper in different places of the module, and the type checkers flag the shadowed one.
The execution fields take the numbers of the workflow-addressed shape, which was never released, and nothing is reserved. Only the two generated descriptors change.
…9-external-interface
The previous shape was never released, so the five channel requests and both `linked_to` fields reuse the numbers `workflow_execution` and the old `linked_to` had, with nothing reserved. Names and types are unchanged.
…-10-redis-provider
…-workflow-runtime
…vider-workflow-streams
…AI-198-py-05-native-wire # Conflicts: # temporalio/api/workflowservice/v1/request_response_pb2.py # temporalio/bridge/sdk-core
…AI-198-streams-provider-nexus
…numbers. The core-3 head carries the renumbered api: the five channel requests and both `linked_to` fields sit on the numbers the workflow fields had, with nothing reserved. Names and types are unchanged.
…98-streams-all # Conflicts: # temporalio/bridge/sdk-core
…98-streams-all # Conflicts: # temporalio/api/workflowservice/v1/request_response_pb2.py # temporalio/bridge/sdk-core
…inal numbers. The previous shape was never released, so the fields keep the numbers the workflow fields had and nothing is reserved.
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 merges both streaming lineages onto
main, so one workflow runs on every provider.What changed?
This is an aggregate PR, labelled
skip-changelog. The content is reviewed in the layers. This PR is the merge resolution, the Core repin and the union conformance run. It merges:result()waiters.What the merge adds of its own:
temporalio/bridge/sdk-corepinsmoe/AI-198-core-all, the union Core with the notification channel in both kinds, the unsubscribe command, the external chain's channel report and a linked channel addressed by execution. The api and bridge protos are regenerated from it.CHANGELOG.md's unreleased block is rewritten by hand from what the tip exports. Repeated merges had spliced and duplicated bullets.temporalio.common.Executionfor the owner. The client'sChannelAddress(channel, execution)is the one address type, withNonefor no owner, the main chain'sExecutionis the one owner type, and the external module re-exports the address. The client's helper that resolvesworkflow_id=into anExecutionis kept once as well.What the union carries from the observability round, each from its layer:
WorkflowExecutionDescription.channel_subscriptions: the channels a run stands on, read fromdescribe().ChannelSubscription.unsubscribe(): records the unsubscribe command and drops the handle before the next task.temporalio.client.stream_channel(ref): the channel a native stream notifies on every append and on its close.stream_consumer_operation: a Nexus operation that consumes a stream through its channel, with a standalone HTTP host.WorkflowStreamChannels: Core subscribes and unsubscribes from it, so no command rides the task that opens a reader.The notification channel comes in from both chains in its two kinds. An independent channel is addressed by name: a workflow listens with
workflow.subscribe_channel(), which issues the subscribe command, and the client'snotify_channel,describe_channel,poll_channeland the listener calls address it without a workflow. A channel linked to an execution, a workflow or a standalone activity, is addressed byexecution=on the same client calls, atemporalio.common.Executionof a type, a business id and an optional run id, withworkflow_id=kept as the short form for a workflow. The owner reads it withworkflow.linked_channel(), which issues no command, since the owner is the listener by construction, and a notification carrieslinked_toas anExecutionso the instance routes it to the handle of its kind. The Redis provider notifies the linked channel of a workflow-owned stream and the independent channel of a standalone or activity-owned one, ahead of the Signal, and the Signal stays as the fallback on a server without channels. The live channel cases are gated on a server named with-E, since the dev server the suite starts has neither kind.Capability split on this head: the memory provider, the native provider and the Redis provider hold activity-owned and standalone streams. Workflow Streams refuses a standalone activity's streams and standalone streams. The Nexus front refuses activity-owned streams. Each refusal is a
StreamUnsupportedError.Part of AI-198 (epic AI-37).
Why?
The swap claim needs one branch where every provider coexists. A workflow file stays the same, and only the provider in
Client.connect(plugins=[...])changes. It looks large againstmainbecause it includes both lineages whole.How did you test it?
poe lintpasses.tests/streamsruns withSTREAMS_LIVE=redis,STREAMS_LIVE=nativeandSTREAMS_LIVE=nexus, in separate runs, against a server from moedash/temporal#15, the layer with the unsubscribe and the channel a native stream notifies. The Redis run was repeated in the server's other switch position,channel.linkedKindEnabledoff, where the external-stream and Redis channel suites run on the independent kind and the retention cases run again. The Nexus consumer file runs there with nothing skipped, the native case and the Redis case included. The memory suite and the Signal fallback run on the stock dev server, and the external stream contrib tests pass on a local Redis in both switch positions. One follow-up stays: a pending Nexus operation task is reported at interpreter exit after the consumer file, harmless.tests/worker/test_workflow.pylast passed at the previous head, where the suite was stable under--workflow-environment time-skippingwith #18 merged. moedash/temporal-agent-harness#2 pins this head, and its four lanes pass. #7 runs the June scenarios on top of it.