Skip to content

Added the stream examples, one per path and one composed. - #7

Open
moedash wants to merge 1136 commits into
moe/AI-198-st-py-unionfrom
moe/AI-198-streams-examples
Open

moedash wants to merge 1136 commits into
moe/AI-198-st-py-unionfrom
moe/AI-198-streams-examples

Conversation

@moedash

@moedash moedash commented Sep 15, 2026 •

Copy link
Copy Markdown
Owner

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 from latest().
  • 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.py and run.py compose 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.py runs them all on one provider.
  • s2 runs the three standalone alternatives on the handle, alt 3 with a StreamRef passed into the workflow start. s8 (b) returns a StreamRef from a Nexus operation.
  • README.md maps each path to its call. _setup.py is the one place a store is named, and every example defines its topics once with streams.topic().
  • Outside examples/: a changelog line, examples under the type checker in pyproject.toml, and the Nexus endpoint setup in streams_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 -

  • Unit Tests
  • Staging
  • End to End Tests

poe lint passes with examples/ under the type checker. Against a local server built from the stream server PRs, the three path examples and run.py exit 0 on workflow_streams, native and redis, and run.py nexus exits 0 through an endpoint on that server. The June suite exits 0 on native, workflow_streams, memory and redis. A provider that can't serve a scenario says so: s2 and s5 (b) and (c) on workflow_streams, s1 (d) off memory, and s7 on memory. The worker and consumer that die in s6 and s8 do so on purpose.

@moedash
moedash force-pushed the moe/AI-198-streams-examples branch 5 times, most recently from ea477d7 to 472df4c Compare September 16, 2026 16:46
@moedash
moedash force-pushed the moe/AI-198-streams-examples branch 4 times, most recently from 38deab6 to 58d7657 Compare September 17, 2026 04:32
@moedash
moedash force-pushed the moe/AI-198-streams-all branch from 549bbdb to 307dedb Compare September 17, 2026 04:37
@moedash
moedash force-pushed the moe/AI-198-streams-examples branch 11 times, most recently from 10900be to b9956a3 Compare September 21, 2026 17:44
@moedash moedash changed the title Added the worked example that runs one agent on every provider. Added the stream examples, one per path and one composed. Sep 21, 2026
@moedash
moedash force-pushed the moe/AI-198-streams-examples branch from b9956a3 to 535b561 Compare September 21, 2026 18:04
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.
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:13
…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.
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
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.
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