Conversation
moedash
force-pushed
the
moe/AI-198-streams-examples
branch
5 times, most recently
from
September 16, 2026 16:46
ea477d7 to
472df4c
Compare
2 of 3 tasks
moedash
force-pushed
the
moe/AI-198-streams-examples
branch
4 times, most recently
from
September 17, 2026 04:32
38deab6 to
58d7657
Compare
moedash
force-pushed
the
moe/AI-198-streams-all
branch
from
September 17, 2026 04:37
549bbdb to
307dedb
Compare
moedash
force-pushed
the
moe/AI-198-streams-examples
branch
11 times, most recently
from
September 21, 2026 17:44
10900be to
b9956a3
Compare
moedash
force-pushed
the
moe/AI-198-streams-examples
branch
from
September 21, 2026 18:04
b9956a3 to
535b561
Compare
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
…/AI-198-streams-examples
…s' into moe/AI-198-streams-all
…/AI-198-streams-examples
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.
…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.
The union was rebuilt as one PR on main, so the examples move onto it. The tree is the union's plus the examples and their docs.
moedash
changed the base branch from
moe/AI-198-streams-all
to
moe/AI-198-st-py-union
October 3, 2026 12:37
NexusStreams resolves the endpoint by name through the client, so the id the docs asked for failed the lookup. The docs also give --http as a full URL.
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 adds the stream examples, one per path, one composed agent, and the June scenarios.
What changed?
examples/streams/path_a_publish.py: a workflow publishes progress and a backend follows it fromlatest().path_b_produce.py: an Activity appends to its workflow's topic, a backend produces and consumes, and the consumer handles a supersession when the first attempt fails.path_c_consume.py: a workflow consumes a topic fed from outside and runs an Activity per record with the cache off.agent.pyandrun.pycompose the three on every provider and behind the Nexus front. The agent races its reader against the generator, so an exhausted generator fails the run.examples/streams/june_scenarios/has one demo per scenario in Roey's June notes. Each file opens with the notes' heading, a status and why.run.pyruns them all on one provider.s2runs the three standalone alternatives on the handle, alt 3 with aStreamRefpassed into the workflow start.s8(b) returns aStreamReffrom a Nexus operation.README.mdmaps each path to its call._setup.pyis the one place a store is named, and every example defines its topics once withstreams.topic().examples/: a changelog line,examplesunder the type checker inpyproject.toml, and the Nexus endpoint setup instreams_demo/README.md.Part of AI-198 (epic AI-37).
Why?
It puts the portability claim in code a reviewer can run, with the same accessor in every context. The June suite answers Roey's scenarios one by one. It says plainly which are implemented, which are emulated and which are unsupported by design. It sits on the union, since it needs every provider at once.
How did you test it?
Link to a test plan if any -
poe lintpasses withexamples/under the type checker. Against a local server built from the stream server PRs, the three path examples andrun.pyexit 0 onworkflow_streams,nativeandredis, andrun.py nexusexits 0 through an endpoint on that server. The June suite exits 0 onnative,workflow_streams,memoryandredis. A provider that can't serve a scenario says so:s2ands5(b) and (c) onworkflow_streams,s1(d) offmemory, ands7onmemory. The worker and consumer that die ins6ands8do so on purpose.