Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
1260 commits
Select commit Hold shift + click to select a range
93f9ac1
Addressed the client's channel calls to a workflow's linked channel.
moetemp Oct 2, 2026
15ff000
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 2, 2026
55212a7
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moetemp Oct 2, 2026
09f28f9
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moetemp Oct 2, 2026
fb67a8f
Added a conformance case for a reader woken through its linked channel.
moetemp Oct 2, 2026
5157b03
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 2, 2026
2a18a2d
Pinned Core at the linked-channel head and regenerated the clients.
moetemp Oct 2, 2026
dd8ccd5
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moetemp Oct 2, 2026
401616e
Covered the Redis reader woken through its linked channel.
moetemp Oct 2, 2026
4448201
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moetemp Oct 2, 2026
52f9508
Gated the linked receive cases on the Core that delivers the job.
moetemp Oct 2, 2026
035d371
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 2, 2026
fce9be2
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 2, 2026
2f100a4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 2, 2026
de9228a
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moetemp Oct 2, 2026
394faeb
Flipped the linked-channel skip on the delivery Core pin.
moetemp Oct 2, 2026
10ca9cb
Matched the main chain's linked channel surface and the listener list.
moetemp Oct 2, 2026
c88c9d8
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moetemp Oct 2, 2026
505f030
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 2, 2026
77c4289
Expected not found from a poll on a closed workflow's linked channel.
moetemp Oct 2, 2026
270c760
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 2, 2026
5f338e8
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 2, 2026
1ce5766
Repinned Core to the external repairs head.
moetemp Oct 2, 2026
278a377
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moetemp Oct 2, 2026
9173526
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 2, 2026
223f727
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 2, 2026
8378092
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moetemp Oct 2, 2026
6bc691b
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moetemp Oct 2, 2026
83820c3
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moetemp Oct 2, 2026
8e026d4
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moetemp Oct 2, 2026
4a979b6
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moetemp Oct 2, 2026
9b65e1b
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moetemp Oct 2, 2026
d7af9eb
Pinned Core at the linked-channel union head and regenerated the clie…
moetemp Oct 2, 2026
cfd4914
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moetemp Oct 2, 2026
c051962
Added the Nexus operation that consumes a stream through its channel.
moetemp Oct 2, 2026
560e278
Tested the stream consumer against a channel server.
moetemp Oct 2, 2026
af0e4dd
Pinned Core with the describe field and the unsubscribe command.
moetemp Oct 2, 2026
b2baaad
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moetemp Oct 2, 2026
10f634b
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moetemp Oct 2, 2026
65df92f
Named the consumer's channel the way the server names a stream's.
moetemp Oct 2, 2026
5a85fb2
Added channel unsubscribe, subscriptions on describe and stream_channel.
moetemp Oct 2, 2026
86d4d50
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 2, 2026
f121ad6
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 2, 2026
7ee4c7b
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moetemp Oct 2, 2026
ce77f32
Pinned Core at the unsubscribe head and regenerated the clients.
moetemp Oct 2, 2026
4fcf9aa
Flipped the unsubscribe skip on the delivery Core pin.
moetemp Oct 2, 2026
f4a2e25
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moetemp Oct 2, 2026
f235b47
Followed a native stream through its channel on a live server.
moetemp Oct 2, 2026
3352512
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moetemp Oct 2, 2026
1f7cee0
Waited for each standalone append's notification before the next.
moetemp Oct 2, 2026
0acc115
Merge branch 'moe/AI-198-streams-provider-workflow-streams' of github…
moetemp Oct 2, 2026
b4e41f2
Pinned the linked channel's describe fields to the server's timing.
moetemp Oct 2, 2026
4eb1047
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 2, 2026
8dfcc08
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 2, 2026
da557f4
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moetemp Oct 2, 2026
68cdb9e
Repinned Core to the unsubscribe and channel-report head.
moetemp Oct 2, 2026
4795652
Reported the readers' channels instead of subscribing to them.
moetemp Oct 2, 2026
494d87c
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moetemp Oct 2, 2026
fb4caa8
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moetemp Oct 2, 2026
25e36fc
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moetemp Oct 2, 2026
22db179
Took the channel address from the client's stream_channel rule.
moetemp Oct 2, 2026
4b06358
Merge branch 'moe/AI-198-streams-provider-workflow-streams' of github…
moetemp Oct 2, 2026
790f833
Added conformance cases for a reader leaving its channel.
moetemp Oct 2, 2026
767475a
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 2, 2026
662eba9
Checked that the subscription waits for the completion that ends the …
moetemp Oct 2, 2026
b195d77
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 2, 2026
04f8025
Checked where the Redis reader's subscription lands.
moetemp Oct 2, 2026
8e72f58
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moetemp Oct 2, 2026
804f243
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moetemp Oct 2, 2026
90c972d
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moetemp Oct 2, 2026
13750be
Kept the client's ChannelAddress as the one type for both chains.
moetemp Oct 2, 2026
0a0a4c0
Pinned Core at the observability union head and regenerated the clients.
moetemp Oct 2, 2026
455c15d
Consumed a Redis stream through the channel its producer notifies.
moetemp Oct 2, 2026
0d1e659
Built the channel report test's instance the way the union's worker d…
moetemp Oct 2, 2026
6c117ad
Left the wake request's owner empty where the address carries none.
moetemp Oct 2, 2026
db14a33
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moetemp Oct 2, 2026
2a9234c
Pinned Core with a linked channel addressed by execution.
moetemp Oct 3, 2026
0013df7
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moetemp Oct 3, 2026
f84b3c4
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moetemp Oct 3, 2026
49faf69
Pinned Core at the protos that address a channel by execution.
moetemp Oct 3, 2026
a59c3a8
Addressed a linked channel by execution.
moetemp Oct 3, 2026
fbbaf0c
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moetemp Oct 3, 2026
58baa2a
Read the linked owner as an execution in the conformance case.
moetemp Oct 3, 2026
7e3c63f
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 3, 2026
7718b1e
Read the linked owner as an execution in the Redis replay case.
moetemp Oct 3, 2026
d6b39ad
Addressed a linked channel by execution on the client.
moetemp Oct 3, 2026
10dd23f
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 3, 2026
6c73610
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 3, 2026
ecaa55d
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 3, 2026
7516568
Pinned the delivery Core with a linked channel addressed by execution.
moetemp Oct 3, 2026
b110ed0
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 3, 2026
c85be89
Registered the stream consumer by execution.
moetemp Oct 3, 2026
7cd9b3a
Carried one execution field on the channel inputs.
moetemp Oct 3, 2026
4af0f93
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 3, 2026
0dce72e
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 3, 2026
94aa0a8
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 3, 2026
5d31c9e
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 3, 2026
f97af15
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moetemp Oct 3, 2026
56ab511
Named a standalone activity's stream channel by execution on the nati…
moetemp Oct 3, 2026
c2aba51
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moetemp Oct 3, 2026
8afd292
Merge commit 'c2aba51e9d88d7c27571504d9c451fb63c65e2fd' into moe/AI-1…
moetemp Oct 3, 2026
f7d6be0
Merge commit '5d31c9ee994f4175bb4e0175ae88f5ebfbd6ecc0' into moe/AI-1…
moetemp Oct 3, 2026
aa4fe82
Merge commit '7718b1ef273a5d80596a5fa53ee7e881cc4c8617' into moe/AI-1…
moetemp Oct 3, 2026
df93277
Pinned Core at the execution-addressed union head.
moetemp Oct 3, 2026
611d772
Kept one copy of the client's channel owner resolver.
moetemp Oct 3, 2026
ffde6ba
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moetemp Oct 3, 2026
97ddcd8
Pinned Core at the protos that reuse the old numbers for execution.
moetemp Oct 3, 2026
e9ae44d
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moetemp Oct 3, 2026
d4022f6
Pinned Core with the execution fields on their original numbers.
moetemp Oct 3, 2026
01d9773
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moetemp Oct 3, 2026
cdeab8a
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moetemp Oct 3, 2026
6a7f022
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moetemp Oct 3, 2026
44dfd4f
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moetemp Oct 3, 2026
7ef8805
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moetemp Oct 3, 2026
7657155
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 3, 2026
995c4b4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moetemp Oct 3, 2026
38aa617
Pinned the delivery Core with the execution fields on their original …
moetemp Oct 3, 2026
56bf209
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moetemp Oct 3, 2026
e473391
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moetemp Oct 3, 2026
0686ca4
Merge commit 'e4733918f4ea525562fd47960a0ac698838e5bfd' into moe/AI-1…
moetemp Oct 3, 2026
6012184
Merge commit '995c4b4bdf6bf38679770fc6252851083b99f555' into moe/AI-1…
moetemp Oct 3, 2026
30b6ad2
Merge commit '01d9773d9faeaa0798838c349cecefe90c6c721d' into moe/AI-1…
moetemp Oct 3, 2026
0ddd6a2
Pinned Core at the union head with the execution fields on their orig…
moetemp Oct 3, 2026
b7afe4e
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moetemp Oct 3, 2026
bfc3a5e
Merged the rebuilt union into the examples.
moetemp Oct 3, 2026
0a95052
Took the Nexus endpoint by name in the streams example.
moetemp Oct 3, 2026
e6898c5
Pinned Core at the external runtime's channel wake.
moetemp Oct 5, 2026
bbb4aaf
Listed the external-stream fixes the Core pin brings in.
moetemp Oct 5, 2026
2c40524
Move sdk-core submodule to 6e90e6d5 and fix bridge build
mfateev Aug 14, 2026
d8dbfdd
External streams: generated protos, record model, and the Redis fixture
mfateev Aug 15, 2026
bcc67e6
External streams: backend contract, annotation codec, codec, failure …
mfateev Aug 15, 2026
a6b44bc
External streams: parking contract, Redis provider, registry, produce…
mfateev Aug 15, 2026
1d4b22e
External streams: subscription manager and the Workflow-facing API
mfateev Aug 15, 2026
0615ea9
External streams: wire the feature into the Python Worker
mfateev Aug 16, 2026
aadf012
P13: the replay read path
mfateev Aug 16, 2026
ea590b0
P14: the producer wake-signal path
mfateev Aug 16, 2026
301af4d
P6b: publish()'s acknowledged-wake semantics
mfateev Aug 16, 2026
a6507e9
P15: the Continue-As-New cursor
mfateev Aug 16, 2026
4d0aea7
P21: multiple streams, merge, and same-stream subscriptions
mfateev Aug 16, 2026
8c12e9f
P20: the Worker shutdown wake sweep
mfateev Aug 16, 2026
bd8840b
P16b: the Milestone 2 required-test list
mfateev Aug 16, 2026
5fdab62
Fix the failure-taxonomy row for an annotation mismatch, and P20's mi…
mfateev Aug 16, 2026
d0bad85
P16a/P16b: make the milestone gates enforceable
mfateev Aug 16, 2026
7e0bf4a
Core bump for C8/C10, and three more Milestone 1 cases
mfateev Aug 16, 2026
e0b7957
Replay a stream history through the real Replayer
mfateev Aug 16, 2026
892e549
Deliver a record buffered while Workflow code was elsewhere
mfateev Aug 16, 2026
4ad9809
Core bump for C15b, and correct a replay finding made against a stale…
mfateev Aug 16, 2026
9f53650
The livelock does not exist: fix the test helper that caused it
mfateev Aug 17, 2026
74e5fc3
Bound delivery within one activation (ADR-026)
mfateev Aug 17, 2026
72d467d
Fix four defects the required-test list exposed
mfateev Aug 17, 2026
f87e790
Make an unparked wake's sender identity unique per sender
mfateev Aug 17, 2026
3ee3d05
Rollover cases: one test defect fixed, two reasons corrected
mfateev Aug 17, 2026
3c79510
Rollover cases 19 and 23: both were test defects
mfateev Aug 17, 2026
8ab01be
Milestone 1 gate: 53 of 55, and two honest gaps
mfateev Aug 17, 2026
3f68a8a
Case 29: reach the window it names, and state what actually blocks it
mfateev Aug 17, 2026
987c9f7
Fix four more defects the handoff case exposed, and take the gate to …
mfateev Aug 18, 2026
e878e21
Close the deadlock: report the quiescent snapshot even when commands …
mfateev Aug 18, 2026
01fb3fd
Close case 36: an empty stream parks, is evicted, and replays from it…
mfateev Aug 18, 2026
890af9d
Read the required-test lists from their own directory
mfateev Aug 18, 2026
0a157b5
Bind each wait to its own backend, and stop reporting unsent wakes as…
mfateev Aug 18, 2026
7458e1e
Make readiness reporting total, and stop dropping a stale answer
mfateev Aug 18, 2026
08b78f1
Make the Redis key layout injective, and stop a stream name widening …
mfateev Aug 18, 2026
8be6ec4
Make merge fair, and stop consuming a record before it decodes
mfateev Aug 18, 2026
1358c17
Reconcile an inherited park intent, and make parking match Core's wai…
mfateev Aug 18, 2026
f7c65f7
Bind a late subscription, and replay a segment in its recorded order
mfateev Aug 18, 2026
1c7ee60
Require a provider's key derivation to be injective
mfateev Aug 18, 2026
ae8a561
Decode off the Workflow thread, and connect the failure taxonomy
mfateev Aug 18, 2026
0dd5ba9
Ask the converter's public members, not a predicate marked for removal
mfateev Aug 18, 2026
34abed5
Finish close(): stop the watcher and take back the park intent
mfateev Aug 19, 2026
e37a63c
Guard the close() teardown, none of which was tested
mfateev Aug 19, 2026
99729cc
Close the seven follow-up findings, which were all failure paths
mfateev Aug 19, 2026
02eb1ed
Close the third review's six findings, and the seven the fixes introd…
mfateev Aug 19, 2026
90002e2
Settle the Run before asking the status probe to be repeatable
mfateev Aug 19, 2026
668c0c2
Re-announce buffered records on every completion, not only a budget stop
mfateev Aug 19, 2026
2620ee1
Close the fourth review's five findings, and the one its own fix left…
mfateev Aug 19, 2026
952ec2a
Bind the unknown-append recovery to its operation, stream and producer
mfateev Aug 19, 2026
a0ebb24
Keep one canonical operation across repeated unknown-append recoveries
mfateev Aug 20, 2026
a6e7ae2
Close out the agent handoff notes, and the one defect they left open
mfateev Aug 20, 2026
2de43e9
Close the fifth review's four findings
mfateev Aug 20, 2026
7d2f91a
Close the fifth review's remaining findings, including one in its own…
mfateev Aug 21, 2026
81b518d
Let the continuation writer be staged behind its reader
mfateev Aug 21, 2026
b4aaf04
Name the writer stage, and run a chain through it
mfateev Aug 21, 2026
f7fa7c7
Let a merged wait end for every member, not only its winner
mfateev Aug 21, 2026
381e177
State the merged wait's contract where it belongs, and cover its claim
mfateev Aug 21, 2026
bae10ec
Validate recorded stream state against History, not against configura…
mfateev Aug 22, 2026
60f451c
Draw the unparked wake counter from the sender, not from the Run
mfateev Aug 22, 2026
6c71094
Retry an owed park removal without waiting for another event
mfateev Aug 22, 2026
f4a0425
Close six findings in the owed-removal retry, one in its own fix
mfateev Aug 22, 2026
d66f29a
Carry the wake a retired park intent owes past its subscription
mfateev Aug 23, 2026
94a7a23
Stop three streaming tests failing on things they do not test
mfateev Aug 24, 2026
5897602
Coalesce external stream wake cycles
mfateev Aug 24, 2026
e1dba93
Correlate wake retries with task attempts
mfateev Aug 24, 2026
7f063c5
Stabilize sandboxed handoff fixtures
mfateev Aug 24, 2026
9dac2cc
Await streaming worker test shutdown
mfateev Aug 24, 2026
b1cc1a5
Keep pytest rewriting out of workflow sandboxes
mfateev Aug 24, 2026
afd9730
Format streaming issue fixes
mfateev Aug 24, 2026
776bed2
Document and type-check streaming support
mfateev Aug 24, 2026
e68425b
Point streaming docs at current design
mfateev Aug 25, 2026
d00168e
Use a single external stream backend
mfateev Aug 26, 2026
12f35b7
Expose external workflow stream APIs
mfateev Aug 26, 2026
5ed2941
Implement workflow-originated external streams
mfateev Aug 27, 2026
6857a4a
Update Core streaming documentation
mfateev Aug 27, 2026
3e36f0b
Moved the External Workflow Streams changelog entries to Unreleased.
moetemp Oct 5, 2026
b7f096e
Held the server-backed stream cases off the skipping clock.
moetemp Oct 3, 2026
eb27fe7
Made the external stream docstrings build under pydoctor.
moetemp Oct 3, 2026
bead929
Aligned the external stream input and output replay schedules.
moetemp Oct 3, 2026
e3c49e9
Kept a replay marker from earning a drain of its own.
moetemp Oct 3, 2026
3769016
Allowed external-stream publishes from the workflow constructor.
moetemp Oct 3, 2026
56ff0b8
Covered a publish from the workflow constructor in the conformance su…
moetemp Oct 5, 2026
aee1fe0
Answered a legacy query without the stream snapshot beside it.
moetemp Oct 3, 2026
c94d993
Pointed the legacy query id at Core's own constant.
moetemp Oct 3, 2026
2c83ee6
Woke parked external readers through the stream's channel.
moetemp Oct 3, 2026
69c9c20
Covered the external runtime's channel wake.
moetemp Oct 3, 2026
5ddbdba
Covered the channel wake in the conformance suite.
moetemp Oct 5, 2026
933d1e1
Let an external stream subscription name its start and yield offsets.
moetemp Oct 3, 2026
f857d24
Added the Redis provider over External Workflow Streams.
moetemp Oct 3, 2026
06a1ffc
Covered the Redis provider.
moetemp Oct 3, 2026
9cbac11
Said which provider carries a reader across continue-as-new.
moetemp Oct 3, 2026
398e932
Described the Redis provider's one log per topic in the changelog.
moetemp Oct 3, 2026
4f0065c
Trimmed the Redis streams by retention on every append.
moetemp Oct 3, 2026
8022baf
Covered the Redis retention trims.
moetemp Oct 3, 2026
86af499
Held an activity's own streams on the Redis provider.
moetemp Oct 3, 2026
544c021
Covered activity-owned streams on the Redis provider.
moetemp Oct 3, 2026
121eaed
Let an external stream subscription start at the tail or the newest N.
moetemp Oct 3, 2026
88718cb
Let a Redis workflow reader start at the tail or at the newest N reco…
moetemp Oct 3, 2026
1533839
Covered the Redis workflow reader's tail starts.
moetemp Oct 3, 2026
850526c
Bumped the Core pin to the suite-reconciled head.
moetemp Oct 5, 2026
f19765f
Merged moe/AI-198-o7-py-1-core-pin into moe/AI-198-o7-py-2-import.
moetemp Oct 5, 2026
d9a66ec
Merged moe/AI-198-o7-py-2-import into moe/AI-198-o7-py-3-tooling.
moetemp Oct 5, 2026
ee823cc
Merged moe/AI-198-o7-py-3-tooling into moe/AI-198-o7-py-4-replay-sche…
moetemp Oct 5, 2026
1c7bace
Merged moe/AI-198-o7-py-4-replay-schedule into moe/AI-198-o7-py-5-con…
moetemp Oct 5, 2026
7787e55
Merged moe/AI-198-o7-py-5-constructor-publish into moe/AI-198-o7-py-6…
moetemp Oct 5, 2026
ddda20c
Merged moe/AI-198-o7-py-6-parked-run-query into moe/AI-198-o7-py-7-ch…
moetemp Oct 5, 2026
2ad4207
Merged moe/AI-198-o7-py-7-channel-wake into moe/AI-198-o7-py-8-redis-…
moetemp Oct 5, 2026
3978468
Merged moe/AI-198-o7-py-8-redis-provider into moe/AI-198-o7-py-9-redi…
moetemp Oct 5, 2026
52203df
Merged moe/AI-198-o7-py-9-redis-trimming into moe/AI-198-o7-py-10-red…
moetemp Oct 5, 2026
ec87862
Merged moe/AI-198-o7-py-10-redis-activity-owners into moe/AI-198-o7-p…
moetemp Oct 5, 2026
be39a8c
Hosted standalone streams on the Redis provider.
moetemp Oct 3, 2026
4f3ff45
Covered standalone streams on the Redis provider.
moetemp Oct 3, 2026
2b9be04
Aborted a dead output stage when the Worker evicts its run.
moetemp Oct 3, 2026
c44bfde
Covered the dead stage abort on the Redis provider.
moetemp Oct 3, 2026
29aaa94
Merged the native main chain into the union.
moetemp Oct 5, 2026
c40987e
Merged the external chain into the union.
moetemp Oct 5, 2026
5d72e71
Merged the time-skipping unlock fix into the union.
moetemp Oct 5, 2026
02d4ecd
Skipped the native conformance setup without a server address.
moetemp Oct 3, 2026
477b78f
Ran the native provider cases in the union.
moetemp Oct 3, 2026
268215b
Consumed a Redis stream through the channel its producer notifies.
moetemp Oct 3, 2026
29cfd44
Ran the demo on each of the union's four providers.
moetemp Oct 3, 2026
9128bb0
Fixed the visitor generator's repeated-field check.
moetemp Oct 3, 2026
9320b44
Merged the v4 union into the examples.
moetemp Oct 5, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,8 @@ to include examples, links to docs, or any other relevant information.
through `ExternalOutputStreamClient`. Workflow output is staged outside
History and becomes readable only after its compact Workflow Task marker is
committed.
- Added `examples/streams`, one agent loop that runs unchanged on every
stream provider and on the Nexus front.

### Changed

Expand Down
1 change: 1 addition & 0 deletions examples/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
"""Worked examples that run against a Temporal server."""
57 changes: 57 additions & 0 deletions examples/streams/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
# Streams, by path

One provider, registered once on the client. Workers built from that client
inherit it, and every context asks for its stream the same way. A topic is
defined once, with the type its records carry, and every context refers to
that definition, so no call names a type again.

```python
client = await Client.connect("localhost:7233", plugins=[provider])

INPUTS = streams.topic("inputs", Token)
DECISIONS = streams.topic("decisions", Decision)
```

| Path | Who | Call | Example file |
|---|---|---|---|
| A: the workflow publishes | workflow code | `workflow.stream_writer(PROGRESS).publish(Progress(...))`, then `.finish()` | `path_a_publish.py` |
| A: a backend follows | any process with a client | `stream = client.get_stream_handle(workflow_id)`, then `stream.read(topic=PROGRESS, after=await stream.latest(topic=PROGRESS))` | `path_a_publish.py` |
| B: an Activity produces | activity code | `await activity.stream_handle().producer(topic=INPUTS).append(Token(...))` | `path_b_produce.py` |
| B: a backend produces | any process with a client | `client.get_stream_handle(workflow_id).producer(topic=NOTES, producer_id=..., attempt=...)` | `path_b_produce.py` |
| B: a backend consumes | any process with a client | `client.get_stream_handle(workflow_id).read(topic=INPUTS)` | `path_b_produce.py` |
| C: the workflow consumes | workflow code | `async for record in workflow.stream_reader(COMMANDS)` | `path_c_consume.py` |

A plain string names a topic decided at runtime, with `result_type=` on the
call; the examples never need one.

`agent.py` and `run.py` compose all three paths in one agent, on every provider
and behind the Nexus front. `_setup.py` is the one place a store is named.
`june_scenarios/` maps every scenario in Roey's June design notes onto this
surface, one file per scenario family, with a status on each.

## Running

Each example takes the provider's name and runs against a dev server:

```sh
python -m examples.streams.path_a_publish workflow_streams
python -m examples.streams.path_a_publish native --address 127.0.0.1:7333
python -m examples.streams.path_a_publish redis --redis redis://127.0.0.1:6379
```

Swap `path_a_publish` for `path_b_produce`, `path_c_consume` or `run`. The
`native` provider needs a server built from the stream-carrying branch; the
`redis` provider needs a Redis to point at. `run.py` also takes `nexus`, with
an endpoint routed to the handler worker's task queue.

## Why the workflow's verbs differ

An Activity and a backend hold the same `StreamHandle`, with the same verbs,
because both act on the store at once: a producer's records are visible as
soon as the store accepts them, and a read follows the store live. Workflow
code gets two verbs of its own because its semantics differ. `publish` is
buffered and commits with the Workflow Task, so no reader can see a record
from a task that failed, and a `stream_reader` is an observation the SDK
records, so replay re-supplies the same records in the same order. That is
why `publish` is a plain call and the reader is an async iterator, and why
neither takes a client.
1 change: 1 addition & 0 deletions examples/streams/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
"""One agent loop run on every stream provider."""
60 changes: 60 additions & 0 deletions examples/streams/_setup.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
"""Provider selection for the examples: the one place a store is named.

Every example takes the provider's name on the command line, builds it here,
and registers it once on the client. Nothing else in the examples names a
store: workers built from the client inherit the provider, and each context
asks for its stream through ``workflow.stream_reader`` or
``workflow.stream_writer``, ``activity.stream_handle()`` and
``client.get_stream_handle()``.
"""

from __future__ import annotations

import argparse

from temporalio.client import Client
from temporalio.streams.providers import ProviderPlugin

PROVIDERS = ("workflow_streams", "native", "redis")


def parser(
description: str, providers: tuple[str, ...] = PROVIDERS
) -> argparse.ArgumentParser:
"""The flags every example shares."""
parser = argparse.ArgumentParser(description=description)
parser.add_argument("provider", choices=providers)
parser.add_argument("--address", default="localhost:7233")
parser.add_argument("--redis", default="redis://127.0.0.1:6379")
return parser


def make_provider(name: str, args: argparse.Namespace) -> ProviderPlugin:
"""The whole difference between the stores: one constructor call."""
if name == "workflow_streams":
from temporalio.streams.providers.workflow_streams import (
WorkflowStreamsProvider,
)

return WorkflowStreamsProvider()
if name == "redis":
from temporalio.streams.providers.redis import RedisStreams

return RedisStreams(url=args.redis)
if name == "native":
from temporalio.streams.providers.native import NativeStreams

return NativeStreams()
if name == "memory":
# Only for examples that keep a warm cache: this provider is not
# replay-safe, which is why PROVIDERS leaves it out.
from temporalio.streams.providers.memory import MemoryStreams

return MemoryStreams()
raise SystemExit(f"unknown provider {name}")


async def connect(args: argparse.Namespace) -> tuple[Client, ProviderPlugin]:
"""A client with the provider registered on it, and the provider to close later."""
provider = make_provider(args.provider, args)
return await Client.connect(args.address, plugins=[provider]), provider
125 changes: 125 additions & 0 deletions examples/streams/agent.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
"""The workflow and activities. Identical on every provider.

Nothing here names a store, a transport, or an option. The two topics are
defined once, with the types their records carry, and the workflow, the
Activity and the backend in ``run.py`` all refer to them. The loop reads its
``inputs`` topic, decides, publishes the decision, and runs an ordinary
activity in the same workflow task, which is the shape the design doc calls
Paths A, B and C together. The Activity that streams model output asks its
context for its own workflow's stream, the way workflow code asks its
runtime, so the file is the same whichever provider the process registered.
"""

from __future__ import annotations

import asyncio
from dataclasses import dataclass
from datetime import timedelta

from temporalio import activity, streams, workflow
from temporalio.common import RetryPolicy
from temporalio.streams import RecordKind


@dataclass
class Token:
"""One piece of model output."""

n: int


@dataclass
class Decision:
"""What the workflow decided about a token, or which attempt it retracted."""

echo: int | None = None
retracting_attempt: int | None = None


INPUTS = streams.topic("inputs", Token)
DECISIONS = streams.topic("decisions", Decision)


@activity.defn
async def generate(count: int) -> None:
"""Stream model output onto this workflow's ``inputs`` topic.

No workflow id and no run id: the handle is this Activity's own
workflow, pinned to its run. The producer carries the Activity's own id
and attempt, so a retry deduplicates and a new attempt is reported to
readers as a supersession.
"""
model = activity.stream_handle().producer(topic=INPUTS)
for n in range(count):
await model.append(Token(n))
await model.finish()


@activity.defn
async def record_decision(decision: Decision) -> str:
"""An ordinary activity, run from the same task that read and published."""
return f"recorded {decision.echo}"


@workflow.defn
class Agent:
"""Reads ``inputs``, publishes a decision each time, ends on FINISH."""

@workflow.run
async def run(self, count: int) -> int:
"""Decide on at most ``count`` inputs, then return how many landed."""
decisions = workflow.stream_writer(DECISIONS)

generating = workflow.start_activity(
generate,
count,
start_to_close_timeout=timedelta(minutes=1),
# Bounded, so a generator that cannot finish gives up instead of
# retrying forever while every attempt streams from the start.
retry_policy=RetryPolicy(maximum_attempts=3),
)
consuming = asyncio.create_task(self._consume(count, decisions))

# Raced rather than awaited in turn: an attempt that fails writes no
# FINISH, so a generator that exhausts its attempts leaves the reader
# waiting forever. Its failure ends the run with its cause instead.
done, _ = await workflow.wait(
[consuming, generating], return_when=asyncio.FIRST_COMPLETED
)
if generating in done and consuming not in done:
try:
await generating
except BaseException:
consuming.cancel()
raise
seen = await consuming
await generating

decisions.finish()
return seen

async def _consume(
self, count: int, decisions: workflow.StreamWriter[Decision]
) -> int:
seen = 0
async for record in workflow.stream_reader(INPUTS):
if record.kind is RecordKind.FINISH:
break
if record.kind is RecordKind.SUPERSEDED:
assert record.supersession is not None
decisions.publish(
Decision(retracting_attempt=record.supersession.previous_attempt)
)
continue
assert record.value is not None
seen += 1
decision = Decision(echo=record.value.n)
decisions.publish(decision)
await workflow.execute_activity(
record_decision,
decision,
start_to_close_timeout=timedelta(minutes=1),
)
if seen >= count:
break
return seen
53 changes: 53 additions & 0 deletions examples/streams/june_scenarios/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
# Roey's June scenarios, on the shipped surface

Every scenario in Roey's Notion page "Streaming Design Discussion Prep Notes"
(June 2, under "Streaming Links") mapped onto `temporalio.streams` as it
ships on this branch. Each file opens with his scenario heading, a status,
and one sentence why. Where his sketch uses a call shape we do not have, the
docstring shows his shape in one line and the code uses ours. Nothing here
reaches into private SDK code or adds a feature.

| Roey's scenario | File | Status | Note |
|---|---|---|---|
| Client starts and consume stream: primary and named | `s1_client_consumes.py` | implemented | Default topic with no name, a typed topic, `last=N`, `after=END`, `BEGINNING` on a moved floor (memory only) |
| Client starts and consume stream: standalone alt 1, 2, 3 | `s2_standalone_streams.py` | alts 1, 2 and 3 implemented | Alt 1 as a read that parks until `create_stream` and an append land (native; memory and Redis answer `StreamNotFoundError`); alt 2 as `client.create_stream`, a policy floor, `close()` and `StreamClosedError`; alt 3 as a stream created first and passed into the workflow start as a `StreamRef`, opened in the activity with `activity.stream_handle(ref)`. A start that commits the stream with the workflow remains the design question. Workflow Streams declines |
| Workflow as Producer: as named handle | `s3_workflow_producer.py` | implemented | His turn loop with continue-as-new; the client follows the chain live |
| Workflow as Producer: as return type | `s4_workflow_as_generator.py` | emulated | Default-topic publishes plus `FINISH`, result from the workflow; the generator signature is sugar not built |
| Activity as Producer: as named handle | `s5_activity_producers.py` | implemented | Workflow topic (Path B), `scope="activity"`, standalone activity; all three on native, memory and Redis, the last two skipped on Workflow Streams |
| Activity as Producer: as return type | `s6_activity_as_generator.py` | emulated | Appends plus a heartbeat checkpoint; the retry resumes and readers see `SUPERSEDED` |
| Workflow as Consumer | `s7_workflow_consumer.py` | implemented; foreign stream unsupported | Own inbound topic across continue-as-new, handing over per batch with one producer per batch and carrying a checkpoint; runs on native, Workflow Streams and Redis, memory skips by design; reading a foreign stream from a workflow is rule 5 |
| Client as Consumer over Standalone Nexus | `s8_nexus_consumers.py` (a) | implemented | Activity reads through the `NexusStreams` front and resumes from a heartbeat cursor |
| Nexus operation handler | `s8_nexus_consumers.py` (b) | implemented | The operation returns `temporalio.streams.StreamRef`, taken from the producing workflow's handle, and the client opens it with `get_stream_handle(ref)` on a client whose provider is the front; a stream type of its own in the operation IDL is the nexgen follow-on |
| Workflow as Consumer over Nexus | `s8_nexus_consumers.py` docstring | unsupported by design | A workflow's reads ride its Workflow Task and never cross Nexus |

## Running

Each file runs on its own and takes the provider's name, the same way the
examples one directory up do. `run.py` runs them all in order:

```sh
python -m examples.streams.june_scenarios.run native --address 127.0.0.1:7433 --http http://127.0.0.1:7343
python -m examples.streams.june_scenarios.run workflow_streams --address 127.0.0.1:7433 --http http://127.0.0.1:7343
python -m examples.streams.june_scenarios.run memory --address 127.0.0.1:7433 --http http://127.0.0.1:7343
python -m examples.streams.june_scenarios.run redis --address 127.0.0.1:7433 --http http://127.0.0.1:7343 --redis redis://127.0.0.1:6379
python -m examples.streams.june_scenarios.s2_standalone_streams native --address 127.0.0.1:7433
```

`native` needs a server built from the stream-carrying branch, and `s2`
alt 1's read that parks until the stream is created needs one built from
its current head. `s5` (b) and (c) need a server with standalone activities
and activity-owned streams, and `s8` needs the server's Nexus HTTP ingress
(`--http`, default `http://127.0.0.1:7243`, `7343` on the server above);
the stream-carrying server has all of them, so the commands above point
every provider at it. `s8` creates and deletes its own Nexus endpoint.
`memory` is offered here, not in the parent examples, because it is not
replay-safe; these scenarios keep a warm cache. `memory`, `native` and
`redis` hold standalone streams, so `s2` runs on all three; `redis` runs
with `--redis` naming a local Redis. Every scenario ran green on all four
providers against that server, apart from the refusals below.

A scenario a provider cannot serve says so in its output and moves on:
`s2` on `workflow_streams`, which keeps a stream inside a workflow's log,
`s5` (b) and (c) on `workflow_streams`, `s1` (d) on anything but `memory`,
and `s7` on `memory`, which keeps one topic across a chain rather than one
per run.
1 change: 1 addition & 0 deletions examples/streams/june_scenarios/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
"""Roey's June streaming scenarios, each mapped onto the shipped streams surface."""
38 changes: 38 additions & 0 deletions examples/streams/june_scenarios/_common.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
"""What every scenario shares: its flags, its ids and one way to print.

The store is still named in one place, ``examples.streams._setup``. The
scenarios add the memory provider to the choices because two of them need an
activity-owned stream or truncation, and memory is the one in-process store
that has both.
"""

from __future__ import annotations

import argparse
import uuid

from examples.streams import _setup

PROVIDERS = (*_setup.PROVIDERS, "memory")


def parser(description: str) -> argparse.ArgumentParser:
"""The example flags, plus the memory provider and the Nexus ingress."""
parser = _setup.parser(description, PROVIDERS)
parser.add_argument(
"--http",
default="http://127.0.0.1:7243",
help="the server's Nexus HTTP ingress, for s8",
)
return parser


def ids(prefix: str) -> tuple[str, str]:
"""A fresh workflow id and its own task queue, so reruns never collide."""
workflow_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
return workflow_id, f"tq-{workflow_id}"


def banner(title: str, provider: str) -> None:
"""Head each scenario's output, so a run of all of them reads in sections."""
print(f"\n== {title} [{provider}]")
Loading
Loading