Repository navigation
Conversation
moetemp
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
moetemp
force-pushed
the
moe/AI-198-streams-examples
branch
4 times, most recently
from
September 17, 2026 04:32
38deab6 to
58d7657
Compare
moetemp
force-pushed
the
moe/AI-198-streams-all
branch
from
September 17, 2026 04:37
549bbdb to
307dedb
Compare
moetemp
force-pushed
the
moe/AI-198-streams-examples
branch
11 times, most recently
from
September 21, 2026 17:44
10900be to
b9956a3
Compare
moetemp
force-pushed
the
moe/AI-198-streams-examples
branch
from
September 21, 2026 18:04
b9956a3 to
535b561
Compare
2 of 3 tasks
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.
…AI-198-streams-provider-nexus
…9-external-interface
…e-wire # Conflicts: # temporalio/api/workflow/v1/message_pb2.py # temporalio/api/workflowservice/v1/request_response_pb2.py # temporalio/bridge/sdk-core
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.
…structor-publish.
…-parked-run-query.
…is-activity-owners.
…y-11-redis-start-positions.
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
Owner
Author
|
Reopened as #89 to keep the numbers in review order. Same branch and content. |
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.