Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
1136 commits
Select commit Hold shift + click to select a range
8ed2523
Declared the byte-cap refusal absent on Workflow Streams.
moedash Oct 1, 2026
ee0ab1b
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
ddec3a8
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
eee557c
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 1, 2026
50ac082
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
67f80ea
Took an aborted batch's entries out of the shared log.
moedash Oct 1, 2026
e5fb141
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 1, 2026
89647e4
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
a7804c2
Declared the reference's optional members nullable on the wire.
moedash Oct 1, 2026
b52e331
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 1, 2026
ee631c1
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moedash Oct 1, 2026
78aeddd
Aborted only History-rejected output stages at eviction.
moedash Oct 1, 2026
efcc65f
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 1, 2026
669e994
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
c502b09
Regenerated the vendored protos and the bridge client for the wake.
moedash Oct 1, 2026
968e32c
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moedash Oct 1, 2026
70de9eb
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 1, 2026
bf31846
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 1, 2026
b473638
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 1, 2026
369b792
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
85f3618
Repinned Core for the wake protos.
moedash Oct 1, 2026
d33b6af
Sent external-stream wakes over the wake call, with a Signal fallback.
moedash Oct 1, 2026
485fb08
Merge remote-tracking branch 'origin/moe/AI-198-py-08-external-repair…
moedash Oct 1, 2026
bd7bc95
Merge remote-tracking branch 'origin/moe/AI-198-py-09-external-interf…
moedash Oct 1, 2026
41b05d8
Repinned Core for the wake protos.
moedash Oct 1, 2026
3c2040a
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 1, 2026
176ad2a
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
a9049ba
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moedash Oct 1, 2026
eb0310c
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 1, 2026
d42234a
Merge remote-tracking branch 'origin/moe/fix-time-skipping-unlock-ref…
moedash Oct 1, 2026
78454f4
Woke Redis readers with the entry id as the wake position.
moedash Oct 1, 2026
9bbe9ad
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 1, 2026
6c337fd
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
a9fcc6d
Kept the streams docstrings free of links to later layers.
moedash Oct 1, 2026
8de8b62
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 1, 2026
e69e3e0
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 1, 2026
bb7188f
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 1, 2026
355cbcf
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
0b9bf57
Pointed the Nexus provider docstring at the call that opens a reference.
moedash Oct 1, 2026
17ef4dd
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Oct 1, 2026
8d1d9a0
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 1, 2026
6979d7a
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
19216f9
Repinned Core to the delivery head with the C bridge wake dispatch.
moedash Oct 1, 2026
20330e8
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 1, 2026
a254ad9
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
a32694a
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moedash Oct 1, 2026
06a5924
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 1, 2026
c19be0b
Marked a wake() mention as literal so the docs build resolves.
moedash Oct 1, 2026
bdea7b3
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
b3c3e1f
Kept the producer's doc comment free of a pydoctor link.
moedash Oct 1, 2026
48874dd
Merge remote-tracking branch 'origin/moe/AI-198-py-08-external-repair…
moedash Oct 1, 2026
fccdaae
Merge remote-tracking branch 'origin/moe/AI-198-py-09-external-interf…
moedash Oct 1, 2026
b252c4d
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 1, 2026
d32b88f
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
74e0b75
Allowed external-stream publishes from the workflow constructor.
moedash Oct 1, 2026
94a2d5d
Merge remote-tracking branch 'origin/moe/AI-198-py-08-external-repair…
moedash Oct 1, 2026
41ddb61
Added a conformance case for a publish from the constructor.
moedash Oct 1, 2026
db43568
Merge remote-tracking branch 'origin/moe/AI-198-py-09-external-interf…
moedash Oct 1, 2026
bbc098f
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 1, 2026
9131330
Read the constructor-publish case while the worker still polls.
moedash Oct 1, 2026
a790cc0
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 1, 2026
f57fd80
Satisfied basedpyright on the wake transport and the new external-str…
moedash Oct 2, 2026
bdb7d94
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
72bb7ed
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
a7c8a42
Satisfied basedpyright on the Redis provider's wake transport guard.
moedash Oct 2, 2026
0eb2927
Pinned the Core protos with the notification channel and regenerated …
moedash Oct 2, 2026
28f6d36
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moedash Oct 2, 2026
11496a1
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 2, 2026
aab2644
Added channel subscriptions to the workflow runtime.
moedash Oct 2, 2026
5f28a14
Added the notification channel calls to the client.
moedash Oct 2, 2026
1580fc9
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
d7e7492
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
d751e93
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
f8a99a9
Made the live channel case fail fast on a Core without the command.
moedash Oct 2, 2026
09d8ab0
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
d91ec39
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
83c2792
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
8226574
Skipped the workflow channel case where the pinned Core lacks the com…
moedash Oct 2, 2026
651c32e
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
6a77ad9
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
f7b3e31
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
1b4b23f
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
7e427c5
Pinned Core at the channel-round delivery head and regenerated the cl…
moedash Oct 2, 2026
89eafa3
Mapped the producer conflict under its failed-precondition status.
moedash Oct 2, 2026
cf87ec6
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 2, 2026
c5b259c
Expected the typed cursor error from a read below the floor.
moedash Oct 2, 2026
781a4b1
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
4d049cd
Checked the untouched channel and the retained position in the poller…
moedash Oct 2, 2026
5017569
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
b680133
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
dc3cab3
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
c053c1a
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
afb734a
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 2, 2026
41b898d
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
9c19f23
Pinned Core with the notification channel and regenerated the clients.
moedash Oct 2, 2026
2ce4cd9
Notified the stream's channel ahead of the wake call and the Signal.
moedash Oct 2, 2026
500f494
Subscribed external-stream readers to their streams' channels.
moedash Oct 2, 2026
65b96d0
Covered the channel transport, the subscription and the notified reader.
moedash Oct 2, 2026
bb7e6fb
Added the notification channel surface to the workflow and client.
moedash Oct 2, 2026
8c40b07
Derived a positionless wake's counter from the backend's clock rule.
moedash Oct 2, 2026
1e84b51
Skipped the retained-first-task cases on servers with channels.
moedash Oct 2, 2026
fa7ae83
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
c0a0d18
Added a conformance case for a reader woken through the channel.
moedash Oct 2, 2026
58c87ba
Marked the channel conformance case for a channel server.
moedash Oct 2, 2026
c2b506f
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
98b7349
Accepted the channel transport on the Redis provider.
moedash Oct 2, 2026
a10ced0
Covered the Redis reader woken through the channel.
moedash Oct 2, 2026
9869278
Chose the reset point by the task that consumed the record.
moedash Oct 2, 2026
3bf0c94
Derived the Redis sweep counter from the entry id rule.
moedash Oct 2, 2026
12fad6c
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moedash Oct 2, 2026
f55d4de
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 2, 2026
3f2642f
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 2, 2026
c9daebc
Kept one copy of the channel methods both chains lifted.
moedash Oct 2, 2026
0ad4b3e
Pinned Core at the channel-round union head and regenerated the clients.
moedash Oct 2, 2026
a4392cc
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moedash Oct 2, 2026
3aaa09d
Pinned Core without the point-to-point wake and with the linked chann…
moedash Oct 2, 2026
12fa935
Removed the point-to-point wake transport.
moedash Oct 2, 2026
b4a4216
Pinned Core without the wake and with the linked channel contract.
moedash Oct 2, 2026
c8fb5c1
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moedash Oct 2, 2026
4a9b49e
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 2, 2026
741c97b
Added the channel linked to a workflow to the runtime and the client.
moedash Oct 2, 2026
99df1e6
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
26384f2
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
c14a88d
Listened on the linked channel for the streams a run owns.
moedash Oct 2, 2026
93f9ac1
Addressed the client's channel calls to a workflow's linked channel.
moedash Oct 2, 2026
15ff000
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
55212a7
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
09f28f9
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
fb67a8f
Added a conformance case for a reader woken through its linked channel.
moedash Oct 2, 2026
5157b03
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
2a18a2d
Pinned Core at the linked-channel head and regenerated the clients.
moedash Oct 2, 2026
dd8ccd5
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 2, 2026
401616e
Covered the Redis reader woken through its linked channel.
moedash Oct 2, 2026
4448201
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
52f9508
Gated the linked receive cases on the Core that delivers the job.
moedash Oct 2, 2026
035d371
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
fce9be2
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
2f100a4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
de9228a
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
394faeb
Flipped the linked-channel skip on the delivery Core pin.
moedash Oct 2, 2026
10ca9cb
Matched the main chain's linked channel surface and the listener list.
moedash Oct 2, 2026
c88c9d8
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
505f030
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
77c4289
Expected not found from a poll on a closed workflow's linked channel.
moedash Oct 2, 2026
270c760
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
5f338e8
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
1ce5766
Repinned Core to the external repairs head.
moedash Oct 2, 2026
278a377
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
9173526
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
223f727
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
8378092
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
6bc691b
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 2, 2026
83820c3
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
8e026d4
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moedash Oct 2, 2026
4a979b6
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 2, 2026
9b65e1b
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 2, 2026
d7af9eb
Pinned Core at the linked-channel union head and regenerated the clie…
moedash Oct 2, 2026
cfd4914
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moedash Oct 2, 2026
c051962
Added the Nexus operation that consumes a stream through its channel.
moedash Oct 2, 2026
560e278
Tested the stream consumer against a channel server.
moedash Oct 2, 2026
af0e4dd
Pinned Core with the describe field and the unsubscribe command.
moedash Oct 2, 2026
b2baaad
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 2, 2026
10f634b
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moedash Oct 2, 2026
65df92f
Named the consumer's channel the way the server names a stream's.
moedash Oct 2, 2026
5a85fb2
Added channel unsubscribe, subscriptions on describe and stream_channel.
moedash Oct 2, 2026
86d4d50
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
f121ad6
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
7ee4c7b
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
ce77f32
Pinned Core at the unsubscribe head and regenerated the clients.
moedash Oct 2, 2026
4fcf9aa
Flipped the unsubscribe skip on the delivery Core pin.
moedash Oct 2, 2026
f4a2e25
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 2, 2026
f235b47
Followed a native stream through its channel on a live server.
moedash Oct 2, 2026
3352512
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
1f7cee0
Waited for each standalone append's notification before the next.
moedash Oct 2, 2026
0acc115
Merge branch 'moe/AI-198-streams-provider-workflow-streams' of github…
moedash Oct 2, 2026
b4e41f2
Pinned the linked channel's describe fields to the server's timing.
moedash Oct 2, 2026
4eb1047
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 2, 2026
8dfcc08
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 2, 2026
da557f4
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
68cdb9e
Repinned Core to the unsubscribe and channel-report head.
moedash Oct 2, 2026
4795652
Reported the readers' channels instead of subscribing to them.
moedash Oct 2, 2026
494d87c
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 2, 2026
fb4caa8
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 2, 2026
25e36fc
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
22db179
Took the channel address from the client's stream_channel rule.
moedash Oct 2, 2026
4b06358
Merge branch 'moe/AI-198-streams-provider-workflow-streams' of github…
moedash Oct 2, 2026
790f833
Added conformance cases for a reader leaving its channel.
moedash Oct 2, 2026
767475a
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
662eba9
Checked that the subscription waits for the completion that ends the …
moedash Oct 2, 2026
b195d77
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
04f8025
Checked where the Redis reader's subscription lands.
moedash Oct 2, 2026
8e72f58
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moedash Oct 2, 2026
804f243
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 2, 2026
90c972d
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Oct 2, 2026
13750be
Kept the client's ChannelAddress as the one type for both chains.
moedash Oct 2, 2026
0a0a4c0
Pinned Core at the observability union head and regenerated the clients.
moedash Oct 2, 2026
455c15d
Consumed a Redis stream through the channel its producer notifies.
moedash Oct 2, 2026
0d1e659
Built the channel report test's instance the way the union's worker d…
moedash Oct 2, 2026
6c117ad
Left the wake request's owner empty where the address carries none.
moedash Oct 2, 2026
db14a33
Merge remote-tracking branch 'origin/moe/AI-198-streams-all' into moe…
moedash Oct 2, 2026
2a9234c
Pinned Core with a linked channel addressed by execution.
moedash Oct 3, 2026
0013df7
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moedash Oct 3, 2026
f84b3c4
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 3, 2026
49faf69
Pinned Core at the protos that address a channel by execution.
moedash Oct 3, 2026
a59c3a8
Addressed a linked channel by execution.
moedash Oct 3, 2026
fbbaf0c
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 3, 2026
58baa2a
Read the linked owner as an execution in the conformance case.
moedash Oct 3, 2026
7e3c63f
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 3, 2026
7718b1e
Read the linked owner as an execution in the Redis replay case.
moedash Oct 3, 2026
d6b39ad
Addressed a linked channel by execution on the client.
moedash Oct 3, 2026
10dd23f
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 3, 2026
6c73610
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 3, 2026
ecaa55d
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
7516568
Pinned the delivery Core with a linked channel addressed by execution.
moedash Oct 3, 2026
b110ed0
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
c85be89
Registered the stream consumer by execution.
moedash Oct 3, 2026
7cd9b3a
Carried one execution field on the channel inputs.
moedash Oct 3, 2026
4af0f93
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 3, 2026
0dce72e
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 3, 2026
94aa0a8
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
5d31c9e
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
f97af15
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 3, 2026
56ab511
Named a standalone activity's stream channel by execution on the nati…
moedash Oct 3, 2026
c2aba51
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 3, 2026
8afd292
Merge commit 'c2aba51e9d88d7c27571504d9c451fb63c65e2fd' into moe/AI-1…
moedash Oct 3, 2026
f7d6be0
Merge commit '5d31c9ee994f4175bb4e0175ae88f5ebfbd6ecc0' into moe/AI-1…
moedash Oct 3, 2026
aa4fe82
Merge commit '7718b1ef273a5d80596a5fa53ee7e881cc4c8617' into moe/AI-1…
moedash Oct 3, 2026
df93277
Pinned Core at the execution-addressed union head.
moedash Oct 3, 2026
611d772
Kept one copy of the client's channel owner resolver.
moedash Oct 3, 2026
ffde6ba
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moedash Oct 3, 2026
97ddcd8
Pinned Core at the protos that reuse the old numbers for execution.
moedash Oct 3, 2026
e9ae44d
Merge branch 'moe/AI-198-py-08-external-repairs' into moe/AI-198-py-0…
moedash Oct 3, 2026
d4022f6
Pinned Core with the execution fields on their original numbers.
moedash Oct 3, 2026
01d9773
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 3, 2026
cdeab8a
Merge branch 'moe/AI-198-py-01-protos' into moe/AI-198-py-02-streams-…
moedash Oct 3, 2026
6a7f022
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 3, 2026
44dfd4f
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 3, 2026
7ef8805
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-streams-pro…
moedash Oct 3, 2026
7657155
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
995c4b4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
38aa617
Pinned the delivery Core with the execution fields on their original …
moedash Oct 3, 2026
56bf209
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 3, 2026
e473391
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 3, 2026
0686ca4
Merge commit 'e4733918f4ea525562fd47960a0ac698838e5bfd' into moe/AI-1…
moedash Oct 3, 2026
6012184
Merge commit '995c4b4bdf6bf38679770fc6252851083b99f555' into moe/AI-1…
moedash Oct 3, 2026
30b6ad2
Merge commit '01d9773d9faeaa0798838c349cecefe90c6c721d' into moe/AI-1…
moedash Oct 3, 2026
0ddd6a2
Pinned Core at the union head with the execution fields on their orig…
moedash Oct 3, 2026
b7afe4e
Merge branch 'moe/AI-198-streams-all' into moe/AI-198-streams-examples
moedash Oct 3, 2026
bfc3a5e
Merged the rebuilt union into the examples.
moedash Oct 3, 2026
0a95052
Took the Nexus endpoint by name in the streams example.
moedash Oct 3, 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