Skip to content

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

Closed
moetemp wants to merge 1260 commits into
moe/AI-198-st-py-unionfrom
moe/AI-198-streams-examples
Closed

moetemp wants to merge 1260 commits into
moe/AI-198-st-py-unionfrom
moe/AI-198-streams-examples

Conversation

@moetemp

@moetemp moetemp 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.

@moetemp
moetemp force-pushed the moe/AI-198-streams-examples branch 5 times, most recently from ea477d7 to 472df4c Compare September 16, 2026 16:46
@moetemp
moetemp force-pushed the moe/AI-198-streams-examples branch 4 times, most recently from 38deab6 to 58d7657 Compare September 17, 2026 04:32
@moetemp
moetemp force-pushed the moe/AI-198-streams-all branch from 549bbdb to 307dedb Compare September 17, 2026 04:37
@moetemp
moetemp force-pushed the moe/AI-198-streams-examples branch 11 times, most recently from 10900be to b9956a3 Compare September 21, 2026 17:44
@moetemp moetemp 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
@moetemp
moetemp force-pushed the moe/AI-198-streams-examples branch from b9956a3 to 535b561 Compare September 21, 2026 18:04
The five calls take workflow_id and run_id, the description reports the kind and the owner, and a notification names the workflow it is linked to.
…e-wire

# Conflicts:
#	temporalio/api/workflow/v1/message_pb2.py
#	temporalio/api/workflowservice/v1/request_response_pb2.py
#	temporalio/bridge/sdk-core
moetemp added 28 commits October 5, 2026 00:49
The shared activity suite runs on Redis when STREAMS_LIVE=redis, and the Redis
module checks the key scheme, the run keying and when a read ends.
The workflow thread cannot see where the store's tail is, so the transport
records the request, the Worker resolves it against the store after the task
that opened the subscription, and the marker records the boundary so replay
reads it from History.
…rds.

stream_reader(after=END) and stream_reader(last=N) on the Redis provider now
ask the transport for a tail start instead of being refused.
A reader at the newest N records and one at END each start where the live run
resolved, and replay starts there too.
A stream with an id and no owner gets its own keys: a hash for the policy and
the seal, and a log per topic. One append script reads the policy, refuses a
sealed stream, and trims by count, age and a per-topic byte total.
The Redis case now declares that it hosts standalone streams, bounds their
bytes and trims an open stream by age, so the conformance suite holds it to
all three.
A completion the server rejects leaves the batch it staged under a token no
marker will name. When History already says the task failed, the Worker aborts
that stage at eviction instead of leaving it for a reader to find. A stage
History committed stays pending for the client's repair.
A batch staged for a rejected task leaves the shared log once the Worker
evicts the run, so an outside reader never meets it.
Core is pinned at the union head, with the bridge protos regenerated from it
in this commit so the merge builds. The changelog keeps each chain's entries
once.
Without TEMPORAL_ADDRESS the cases reached the environment's own server,
which has no stream service, and failed on the wire.
The union carries the native provider, so the cases gated on it no longer
skip.
The Redis producer notifies the stream's channel itself, so one live case
needs no notify by hand. The native case opens its own client with the native
provider on it.
Only the union carries the memory, Workflow Streams, native and Redis providers
together, so the demo's provider switch and its notes live here.
A getattr default is evaluated whether or not the attribute is there, and
protobuf 6 descriptors have is_repeated but no label, so the generator needed
a protobuf 3 interpreter.
# Conflicts:
#	streams_demo/README.md
#	temporalio/bridge/proto/workflow_activation/workflow_activation_pb2.py
#	temporalio/bridge/proto/workflow_commands/workflow_commands_pb2.py
#	temporalio/bridge/sdk-core
@moetemp moetemp closed this Oct 5, 2026
@moetemp

moetemp commented Oct 5, 2026

Copy link
Copy Markdown
Owner Author

Reopened as #89 to keep the numbers in review order. Same branch and content.

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.

2 participants