Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
1186 commits
Select commit Hold shift + click to select a range
a86473e
Merge remote-tracking branch 'origin/moe/AI-198-py-10-redis-provider'…
moedash Sep 30, 2026
7868260
Regenerated the vendored stream service protos for the byte cap.
moedash Oct 1, 2026
966e6c6
Merge branch 'moe/AI-198-py-05-native-wire' into moe/AI-198-py-06-nat…
moedash Oct 1, 2026
20ae07a
Passed the byte cap through and trusted the server on a repeated create.
moedash Oct 1, 2026
6621c5a
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
e80180e
Reused the retention case's result variable instead of redeclaring it.
moedash Oct 1, 2026
a81a8c1
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
22c5e92
Merge remote-tracking branch 'origin/moe/AI-198-py-07-native-replay' …
moedash Oct 1, 2026
93fd4e2
Declared refusing an append past the byte cap as a capability.
moedash Oct 1, 2026
0306e24
Merge branch 'moe/AI-198-py-02-streams-package' into moe/AI-198-py-03…
moedash Oct 1, 2026
5627560
Merge branch 'moe/AI-198-py-03-workflow-runtime' into moe/AI-198-py-0…
moedash Oct 1, 2026
f05f6c3
Let the retention case accept a refused append and a coarse age sweep.
moedash Oct 1, 2026
cb6c773
Registered the provider's handlers before the first task's Updates run.
moedash Oct 1, 2026
9820335
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
303ed81
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Oct 1, 2026
fa5567a
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
moedash Oct 1, 2026
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
eee557c
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-nexu…
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
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
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
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
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
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
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
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
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
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
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
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
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
50 changes: 50 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,56 @@ jobs:
npx doctoc README.md
[[ -z $(git status --porcelain README.md) ]] || (git diff README.md; echo "README changed"; exit 1)

# The client-side (Redis) stream provider's own evidence. Its tests need a store
# this repo does not otherwise stand up, so without this job nothing that proves
# retention, cursor ownership, the staged commit or the paired producer write ever
# runs anywhere but a developer's machine.
streams-redis:
timeout-minutes: 30
runs-on: ubuntu-latest
services:
redis:
image: redis:8-alpine
ports:
- 6379:6379
options: >-
--health-cmd "redis-cli ping"
--health-interval 5s
--health-timeout 3s
--health-retries 10
steps:
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
submodules: recursive
- uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # stable
- uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5
with:
python-version: "3.13"
- uses: Swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2
with:
workspaces: temporalio/bridge -> target
key: streams-redis-${{ env.pythonLocation }}
- uses: arduino/setup-protoc@c65c819552d16ad3c9b72d9dfd5ba5237b9c906b # v3
with:
version: "23.x"
repo-token: ${{ secrets.GITHUB_TOKEN }}
- uses: astral-sh/setup-uv@cec208311dfd045dd5311c1add060b2062131d57 # v8
- run: uv tool install poethepoet
- run: uv sync --all-extras
- run: poe build-develop
# The dev server comes from the test environment, the store from the service
# above. Run serially: the cases measure real timing and share one Redis.
- run: uv run pytest tests/streams -p no:randomly -s
timeout-minutes: 20
env:
STREAMS_LIVE: redis
TEMPORAL_TEST_REDIS_URL: redis://127.0.0.1:6379
AI198_REDIS_URL: redis://127.0.0.1:6379
# Also without the store, so the gate itself keeps working and the memory
# provider's conformance run stays honest.
- run: uv run pytest tests/streams -p no:randomly -s
timeout-minutes: 10

# Verify the optional FIPS build: the Rust core must link aws-lc-fips-sys
# (aws-lc-rs FIPS mode) and must NOT link `ring` (the cargo-tree guard, ported
# from sdk-ruby PR #466's `fips_tree` guard); then run the test suite against the
Expand Down
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,6 @@ temporalio/bridge/temporal_sdk_bridge*
tags
/.claude
tmpclaude-*

# Demo run output, written per provider; the numbers live in the doc.
streams_demo/results-*/
3 changes: 2 additions & 1 deletion .gitmodules
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
[submodule "sdk-core"]
path = temporalio/bridge/sdk-core
url = https://github.com/temporalio/sdk-rust.git
url = https://github.com/moedash/sdk-rust.git
branch = moe/AI-198-core-all
129 changes: 129 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,97 @@ to include examples, links to docs, or any other relevant information.
worker-side factories registered with `StrandsPlugin(sandboxes=...)`.

- Added the `temporalio.contrib.gcp.cloud_run.id` module with the `CloudRunIdPlugin` client plugin to set the worker identity on Cloud Run.
- **Experimental**: `temporalio.streams` defines one stream interface a workflow
can read, decide on, and write. A provider is registered once as a plugin,
`Client.connect(plugins=[provider])`, and workers built from that client
inherit it; each context then asks for its stream the same way:
`workflow.stream_reader()` and `workflow.stream_writer()` in workflow code,
`activity.stream_handle()` in an activity, and `client.get_stream_handle()`
anywhere a client is held. A topic is a typed definition,
`streams.topic("inputs", Token)`, shared by workflow, activity and client
code; a plain string names a topic decided at runtime. The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider, and
code; a plain string names a topic decided at runtime, and a call that names
no topic addresses the default topic, `streams.DEFAULT_TOPIC` (`"output"`,
the server's default stream name). The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider. A stream is
handed to another process as a `streams.StreamRef`, plain data naming the
owner and the topic, which `client.get_stream_handle(ref)` and
`activity.stream_handle(ref)` open; `client.create_stream(stream_id, ...)`
creates a standalone stream with a retention policy, and its handle's
`close()` seals it. A provider runs record bodies through the client's data
converter, so a payload codec and external storage apply to them.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against.
- **Experimental**: `temporalio.streams.providers.redis.RedisStreams` serves the
stream interface over External Workflow Streams, holding one topic as an
input and an output stream.
- **Experimental**: `temporalio.streams.providers.workflow_streams.WorkflowStreamsProvider`
serves the stream interface over the shipped Workflow Streams transport as a
worker plugin, so a workflow reads and publishes through
`temporalio.contrib.workflow_streams` without naming it. Records are the
`StreamRecord` proto inside the shipped item payload, and a handle without a
run id follows continue-as-new run by run.
- **Experimental**: `temporalio.streams.providers.nexus.NexusStreams` puts one
Nexus endpoint in front of a storage provider, so a caller reaches a stream
through the endpoint and never names the store, and
`TemporalStreamsHandler` serves that endpoint by fronting the provider's own
handles. Its contract is defined in `temporal_streams.nexusrpc.yaml` and the
bindings are generated from it; a record crosses as the serialized
`StreamRecord` proto. Configure the front with `data_converter=` to run a
payload codec on the caller side, so records are encoded before they leave
the process.
- **Experimental**: server-side streams. A workflow publishes to a stream it
owns with a command the server applies in its Workflow Task's commit, and
reads the ranges the server delivers on its Workflow Tasks, through
`temporalio.workflow.append_stream_records`, `subscribe_stream` and
`read_stream_records`. `temporalio.client_stream` and
`temporalio.contrib.server_streams` reach the same stream from outside a
workflow, and `temporalio.streams.providers.native.NativeStreams` puts it
behind the shared stream interface with one owned stream per topic. Requires
a server that serves the stream service. `Replayer(stream_client=)` replays a
workflow that read such a stream while the server still holds it: History
records only the offsets each task consumed, so the replayer fetches the
records from the stream service and hands them to the replay with the
history. A range the stream no longer holds fails the replay with
`StreamNotFoundError`. A handle without a run id follows a workflow reset as
it follows a continue-as-new, reading the reset run from the floor its stream
reports, and the replayer fetches the ranges recorded before a reset point
from the run the workflow was reset from. For offline replay,
`Replayer.fetch_stream_slices(client, history)` attaches the records to a
`WorkflowHistory` while the stream is retained, `to_json()` and `from_json()`
carry them as `streamSlices` beside the events, and a history that carries
them replays with no server.
- Added experimental External Workflow Streams in
`temporalio.contrib.external_workflow_streams`. Workflow stream payloads are
stored in a configured external backend instead of Temporal History, with a
Redis Streams provider included. Workflows subscribe with `external_stream`,
external processes publish with `ExternalStreamProducer`, and Workers are
configured with `external_stream_backend`.
- Added the output direction for External Workflow Streams.
Workflows publish with `external_output_stream`, Activities and external
processes use `ExternalOutputStreamProducer`, and external consumers resume
through `ExternalOutputStreamClient`. Workflow output is staged outside
History and becomes readable only after its compact Workflow Task marker is
committed.
- `ExternalStreamSubscription.records()` yields each value with the provider
offset it was read from, for a reader that has to name where it got to.
- Added the `temporalio.contrib.gcp.cloud_run.id` module with the `CloudRunIdPlugin` client plugin to set the worker identity on Cloud Run.
`StreamRecord` proto, and both operations address a stream by a `StreamRef`
naming its owner (a workflow, an activity or a standalone stream) and topic,
which the handler maps onto the store's accessor for that owner. Configure
the front with `data_converter=` to run a payload codec on the caller side,
so records are encoded before they leave the process.
- **Experimental**: `temporalio.streams.providers.workflow_streams.WorkflowStreamsProvider`
serves the stream interface over the shipped Workflow Streams transport as a
worker plugin, so a workflow reads and publishes through
`temporalio.contrib.workflow_streams` without naming it. Records are the
`StreamRecord` proto inside the shipped item payload, and a handle without a
run id follows continue-as-new run by run and a reset into the run reset to.
An outside publish is an Update that answers with the batch's position and
refuses a conflicting repeat, falling back to the shipped Signal on a
workflow whose worker predates it. A workflow's activity keeps its own
streams in the workflow's log under `activity/<id>/<name>`.

### Changed

Expand All @@ -52,6 +143,22 @@ to include examples, links to docs, or any other relevant information.

### Fixed

- Preserve empty activations in the shared External Workflow Streams input and
output replay schedule. Workflows that read input, publish decisions, and
schedule Activities now reproduce that schedule during replay. Inconsistent
prerelease markers are rejected explicitly rather than guessing where omitted
activations belonged.
- Resume external input waits when cold replay encounters a wake in an already
loaded History page, including Workers with workflow caching disabled.
- Avoid an unnecessary output replacement Workflow Task after stream input has
resumed the Workflow and it is waiting on an Activity or timer.
- Keep an incomplete retained external stream task alive when workflow caching
is disabled; evict it after its normal task boundary instead of repeatedly
interrupting input readiness with shutdown markers.
- `WorkflowEnvironment.start_time_skipping()`: concurrent `WorkflowHandle.result()` waiters now
share one time-skipping unlock. The test server holds one lock per in-flight task, so a second
unlock let the clock jump while another workflow still had a task in flight, and that task then
timed out.
- `GoogleAdkPlugin` now passes the optional `anthropic`, `litellm`, and `openai` SDKs through
the workflow sandbox.
- `contrib.deepagents`: prevent duplicate input messages after continue-as-new.
Expand All @@ -72,6 +179,28 @@ to include examples, links to docs, or any other relevant information.
### Added

#### Standalone Activity operator commands
- **Experimental**: `temporalio.streams` defines one stream interface a workflow
can read, decide on, and write. A provider is registered once as a plugin,
`Client.connect(plugins=[provider])`, and workers built from that client
inherit it; each context then asks for its stream the same way:
`workflow.stream_reader()` and `workflow.stream_writer()` in workflow code,
`activity.stream_handle()` in an activity, and `client.get_stream_handle()`
anywhere a client is held. A topic is a typed definition,
`streams.topic("inputs", Token)`, shared by workflow, activity and client
code; a plain string names a topic decided at runtime. The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider. A stream is
handed to another process as a `streams.StreamRef`, plain data naming the
owner and, when it has one, the topic, which `client.get_stream_handle(ref)`
and `activity.stream_handle(ref)` open; `client.create_stream(stream_id, ...)`
creates a standalone stream with a retention policy, and its handle's
`close()` seals it. A provider runs record bodies through the client's data
converter, so a payload codec and external storage apply to them.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against, and
`temporalio.streams.providers.redis.RedisStreams` serves the same interface
over External Workflow Streams, one topic as an input and an output stream.
- `ExternalStreamSubscription.records()` yields each value with the provider
offset it was read from, for a reader that has to name where it got to.

- `ActivityHandle` now supports operator commands for standalone activities: `pause`,
`unpause`, `update_options` and `restore_original_options`.
Expand Down
14 changes: 12 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ dev = [
"pytest-rerunfailures>=16.1",
"pytest-xdist>=3.6,<4",
"moto[s3,server]>=5",
"redis>=5,<9",
"langgraph>=1.1.0",
"langsmith>=0.7.34,<0.9",
"deepagents>=0.6.12,<0.7; python_version >= '3.11'",
Expand Down Expand Up @@ -125,16 +126,19 @@ format = [
]
gen-docs = "uv run scripts/gen_docs.py"
gen-nexus-system-api = "uv run scripts/gen_nexus_system_api.py"
gen-streams-nexus-api = "uv run scripts/gen_streams_nexus_api.py"
gen-protos = [
{ cmd = "uv run scripts/gen_protos.py" },
{ ref = "gen-nexus-system-api" },
{ ref = "gen-streams-nexus-api" },
{ cmd = "uv run scripts/gen_payload_visitor.py" },
{ cmd = "uv run scripts/gen_bridge_client.py" },
{ ref = "format" },
]
gen-protos-docker = [
{ cmd = "uv run scripts/gen_protos_docker.py" },
{ ref = "gen-nexus-system-api" },
{ ref = "gen-streams-nexus-api" },
{ cmd = "uv run scripts/gen_payload_visitor.py" },
{ cmd = "uv run scripts/gen_bridge_client.py" },
{ ref = "format" },
Expand Down Expand Up @@ -198,16 +202,21 @@ exclude = [
'temporalio/api',
'temporalio/bridge/proto',
'temporalio/nexus/system/workflow_service',
'temporalio/streams/providers/_nexus_generated',
]

[[tool.mypy.overrides]]
module = "temporalio.nexus.system.workflow_service.*"
ignore_errors = true

[[tool.mypy.overrides]]
module = "temporalio.streams.providers._nexus_generated.*"
ignore_errors = true

[tool.pydocstyle]
convention = "google"
# https://github.com/PyCQA/pydocstyle/issues/363#issuecomment-625563088
match_dir = "^(?!(docs|scripts|tests|api|proto|system|\\.)).*"
match_dir = "^(?!(docs|scripts|tests|api|proto|system|_nexus_generated|\\.)).*"
add_ignore = [
# We like to wrap at a certain number of chars, even long summary sentences.
# https://github.com/PyCQA/pydocstyle/issues/184
Expand Down Expand Up @@ -241,6 +250,7 @@ privacy = [
"HIDDEN:temporalio.worker.workflow_sandbox.importer",
"HIDDEN:temporalio.worker.workflow_sandbox.in_sandbox",
"HIDDEN:**.*_pb2*",
"HIDDEN:temporalio.streams.providers._nexus_generated._definitions",
]
project-name = "Temporal Python"
sidebar-expand-depth = 2
Expand All @@ -267,7 +277,7 @@ reportUnnecessaryIsInstance = "none"
reportUnnecessaryTypeIgnoreComment = "none"
reportUnusedCallResult = "none"
reportUnknownLambdaType = "none"
include = ["temporalio", "tests"]
include = ["temporalio", "tests", "streams_demo"]
exclude = [
# Exclude auto generated files
"temporalio/api",
Expand Down
11 changes: 4 additions & 7 deletions scripts/gen_payload_visitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,13 +63,10 @@ def name_for(desc: Descriptor) -> str:


def field_is_repeated(field: FieldDescriptor) -> bool:
return bool(
getattr(
field,
"is_repeated",
getattr(field, "label") == FieldDescriptor.LABEL_REPEATED,
)
)
is_repeated = getattr(field, "is_repeated", None)
if is_repeated is not None:
return bool(is_repeated)
return getattr(field, "label") == FieldDescriptor.LABEL_REPEATED


def emit_loop(
Expand Down
Loading
Loading