Skip to content

Added one stream interface with its four providers and the Nexus front. - #6

Closed
moedash wants to merge 1186 commits into
mainfrom
moe/AI-198-streams-all
Closed

moedash wants to merge 1186 commits into
mainfrom
moe/AI-198-streams-all

Conversation

@moedash

@moedash moedash commented Sep 15, 2026 •

Copy link
Copy Markdown
Owner

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:

What the merge adds of its own:

  • temporalio/bridge/sdk-core pins moe/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.
  • The constructor-publish conformance case reads while the worker still polls. Here it also runs on Workflow Streams, which answers reads from the workflow.
  • Both lineages lifted the same channel surface, so the merge keeps one copy of each method on the workflow instance. The copy kept is the one that also feeds the external-stream runtime and shares the once-per-channel command with it. The unsubscribe sits next to it and clears that once-per-run set, so a later subscribe reaches the server again.
  • Both lineages had grown an address type for a channel and its owner, and both added temporalio.common.Execution for the owner. The client's ChannelAddress(channel, execution) is the one address type, with None for no owner, the main chain's Execution is the one owner type, and the external module re-exports the address. The client's helper that resolves workflow_id= into an Execution is kept once as well.
  • One live case consumes a Redis stream from a Nexus operation through the channel the Redis producer notifies itself, next to the memory-store case where the test notifies by hand.

What the union carries from the observability round, each from its layer:

  • WorkflowExecutionDescription.channel_subscriptions: the channels a run stands on, read from describe().
  • 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.
  • The external reader's channel report, 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's notify_channel, describe_channel, poll_channel and the listener calls address it without a workflow. A channel linked to an execution, a workflow or a standalone activity, is addressed by execution= on the same client calls, a temporalio.common.Execution of a type, a business id and an optional run id, with workflow_id= kept as the short form for a workflow. The owner reads it with workflow.linked_channel(), which issues no command, since the owner is the listener by construction, and a notification carries linked_to as an Execution so 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 against main because it includes both lineages whole.

How did you test it?

poe lint passes. tests/streams runs with STREAMS_LIVE=redis, STREAMS_LIVE=native and STREAMS_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.linkedKindEnabled off, 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.py last passed at the previous head, where the suite was stable under --workflow-environment time-skipping with #18 merged. moedash/temporal-agent-harness#2 pins this head, and its four lanes pass. #7 runs the June scenarios on top of it.

  • Unit Tests
  • Staging
  • End to End Tests

@moedash
moedash force-pushed the moe/AI-198-streams-all branch from d8f85b9 to 1771757 Compare September 15, 2026 23:51
@moedash
moedash force-pushed the moe/AI-198-streams-all branch from 549bbdb to 307dedb Compare September 17, 2026 04:37
moedash added a commit that referenced this pull request Sep 21, 2026
Takes #6 at `91a52c71`: the native branch's `Replayer(stream_client=)`, the
reset following, the Redis retention window and the seeded reader, on the
Core pinned to sdk-rust #4's head `f32dee03`. No conflicts.
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.
… 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.
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.
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.
… 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.
@moedash moedash added the skip-changelog Changelog entry rides another PR label Oct 1, 2026
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.
moedash added 28 commits October 2, 2026 18:10
…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.
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.
…AI-198-py-05-native-wire

# Conflicts:
#	temporalio/api/workflowservice/v1/request_response_pb2.py
#	temporalio/bridge/sdk-core
…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.
@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