Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
212 commits
Select commit Hold shift + click to select a range
ec20038
Move sdk-core submodule to 6e90e6d5 and fix bridge build
mfateev Aug 14, 2026
8cdaa2e
Point the Core submodule at the task fork branch
mfateev Aug 15, 2026
245ccba
External streams: generated protos, record model, and the Redis fixture
mfateev Aug 15, 2026
e179ab6
External streams: backend contract, annotation codec, codec, failure …
mfateev Aug 15, 2026
42ff459
External streams: parking contract, Redis provider, registry, produce…
mfateev Aug 15, 2026
f9b199e
External streams: subscription manager and the Workflow-facing API
mfateev Aug 15, 2026
e533438
External streams: regenerate protos for the marker terminal reasons
mfateev Aug 16, 2026
4c6ddc5
External streams: wire the feature into the Python Worker
mfateev Aug 16, 2026
4f95dbe
P13: the replay read path
mfateev Aug 16, 2026
f9ee870
P14: the producer wake-signal path
mfateev Aug 16, 2026
65e48fa
P6b: publish()'s acknowledged-wake semantics
mfateev Aug 16, 2026
6fbd4dc
P15: the Continue-As-New cursor
mfateev Aug 16, 2026
c25123a
P21: multiple streams, merge, and same-stream subscriptions
mfateev Aug 16, 2026
f202383
P20: the Worker shutdown wake sweep
mfateev Aug 16, 2026
ead7644
P16b: the Milestone 2 required-test list
mfateev Aug 16, 2026
4d8bdcf
Fix the failure-taxonomy row for an annotation mismatch, and P20's mi…
mfateev Aug 16, 2026
3cc7c24
P16a/P16b: make the milestone gates enforceable
mfateev Aug 16, 2026
d6ed30a
Core bump for C8/C10, and three more Milestone 1 cases
mfateev Aug 16, 2026
6bfe3cf
Replay a stream history through the real Replayer
mfateev Aug 16, 2026
8e58be1
Deliver a record buffered while Workflow code was elsewhere
mfateev Aug 16, 2026
d191bbd
Core bump for C15b, and correct a replay finding made against a stale…
mfateev Aug 16, 2026
b30b3df
The livelock does not exist: fix the test helper that caused it
mfateev Aug 17, 2026
7313e00
Bound delivery within one activation (ADR-026)
mfateev Aug 17, 2026
faa1c84
Fix four defects the required-test list exposed
mfateev Aug 17, 2026
838810f
Make an unparked wake's sender identity unique per sender
mfateev Aug 17, 2026
c70ec80
Rollover cases: one test defect fixed, two reasons corrected
mfateev Aug 17, 2026
08f1206
Core bump: the replay flag fix
mfateev Aug 17, 2026
ca4008d
Rollover cases 19 and 23: both were test defects
mfateev Aug 17, 2026
ac12d27
Milestone 1 gate: 53 of 55, and two honest gaps
mfateev Aug 17, 2026
0ce689d
Case 29: reach the window it names, and state what actually blocks it
mfateev Aug 17, 2026
c301f17
Fix four more defects the handoff case exposed, and take the gate to …
mfateev Aug 18, 2026
93f0306
Close the deadlock: report the quiescent snapshot even when commands …
mfateev Aug 18, 2026
8d5ee1f
Close case 36: an empty stream parks, is evicted, and replays from it…
mfateev Aug 18, 2026
62b3ff1
Read the required-test lists from their own directory
mfateev Aug 18, 2026
b87e094
Core bump: the review guide
mfateev Aug 18, 2026
e6d4cd9
Core bump: the review guide as a usable root
mfateev Aug 18, 2026
a9c03ec
Bind each wait to its own backend, and stop reporting unsent wakes as…
mfateev Aug 18, 2026
0c92c99
Make readiness reporting total, and stop dropping a stale answer
mfateev Aug 18, 2026
ee4fbf8
Make the Redis key layout injective, and stop a stream name widening …
mfateev Aug 18, 2026
03969db
Make merge fair, and stop consuming a record before it decodes
mfateev Aug 18, 2026
0442dc4
Reconcile an inherited park intent, and make parking match Core's wai…
mfateev Aug 18, 2026
1e369ec
Bind a late subscription, and replay a segment in its recorded order
mfateev Aug 18, 2026
7e040b7
Require a provider's key derivation to be injective
mfateev Aug 18, 2026
8ac891f
Decode off the Workflow thread, and connect the failure taxonomy
mfateev Aug 18, 2026
18a93e5
Ask the converter's public members, not a predicate marked for removal
mfateev Aug 18, 2026
1e94c30
Vendor the documents the review guide points at
mfateev Aug 18, 2026
88c3578
Finish close(): stop the watcher and take back the park intent
mfateev Aug 19, 2026
8c2ff85
Guard the close() teardown, none of which was tested
mfateev Aug 19, 2026
cbb832a
Point sdk-core at the close() documentation
mfateev Aug 19, 2026
2763a25
Close the seven follow-up findings, which were all failure paths
mfateev Aug 19, 2026
5a88733
Point sdk-core at the follow-up review documentation
mfateev Aug 19, 2026
8abb8eb
Close the third review's six findings, and the seven the fixes introd…
mfateev Aug 19, 2026
bf33670
Point sdk-core at the third review documentation
mfateev Aug 19, 2026
bff3f44
Settle the Run before asking the status probe to be repeatable
mfateev Aug 19, 2026
e3c0571
Re-announce buffered records on every completion, not only a budget stop
mfateev Aug 19, 2026
0836355
Close the fourth review's five findings, and the one its own fix left…
mfateev Aug 19, 2026
88e3af7
Bind the unknown-append recovery to its operation, stream and producer
mfateev Aug 19, 2026
4827835
Keep one canonical operation across repeated unknown-append recoveries
mfateev Aug 20, 2026
6d7d193
Point sdk-core at the repeated-recovery documentation
mfateev Aug 20, 2026
065e7f5
Close out the agent handoff notes, and the one defect they left open
mfateev Aug 20, 2026
5ab10ac
Close the fifth review's four findings
mfateev Aug 20, 2026
51e829f
Close the fifth review's remaining findings, including one in its own…
mfateev Aug 21, 2026
0076f3d
Point sdk-core at the fifth review's documentation
mfateev Aug 21, 2026
3992263
Let the continuation writer be staged behind its reader
mfateev Aug 21, 2026
081a949
Point sdk-core at the staged-writer decision
mfateev Aug 21, 2026
6c5e324
Name the writer stage, and run a chain through it
mfateev Aug 21, 2026
5cb918b
Let a merged wait end for every member, not only its winner
mfateev Aug 21, 2026
d73575f
State the merged wait's contract where it belongs, and cover its claim
mfateev Aug 21, 2026
8a4680d
Point sdk-core at the merged-wait blocked-set statement
mfateev Aug 21, 2026
601ca54
Validate recorded stream state against History, not against configura…
mfateev Aug 22, 2026
f8aca38
Draw the unparked wake counter from the sender, not from the Run
mfateev Aug 22, 2026
72d5b8d
Retry an owed park removal without waiting for another event
mfateev Aug 22, 2026
f36a9a0
Point sdk-core at the autonomous owed-removal decision
mfateev Aug 22, 2026
e03eb25
Close six findings in the owed-removal retry, one in its own fix
mfateev Aug 22, 2026
e73ee2c
Point sdk-core at the bound, owned and announced cleanup
mfateev Aug 22, 2026
7392fbd
Carry the wake a retired park intent owes past its subscription
mfateev Aug 23, 2026
85aa5c5
Point sdk-core at the buffered Workflow Task admission fix
mfateev Aug 23, 2026
d9ccd4d
Point sdk-core at the open-issues register
mfateev Aug 24, 2026
50de639
Stop three streaming tests failing on things they do not test
mfateev Aug 24, 2026
82b3c7d
Coalesce external stream wake cycles
mfateev Aug 24, 2026
86144a0
Correlate wake retries with task attempts
mfateev Aug 24, 2026
163019c
Point sdk-core at defensive WFT admission
mfateev Aug 24, 2026
4eca5b4
Stabilize sandboxed handoff fixtures
mfateev Aug 24, 2026
2994f4f
Await streaming worker test shutdown
mfateev Aug 24, 2026
bff2571
Keep pytest rewriting out of workflow sandboxes
mfateev Aug 24, 2026
552e4f0
Format streaming issue fixes
mfateev Aug 24, 2026
380fe87
Document and type-check streaming support
mfateev Aug 24, 2026
d40a076
Record final Core streaming validation
mfateev Aug 24, 2026
7e9c821
Point streaming docs at current design
mfateev Aug 25, 2026
858669c
Use a single external stream backend
mfateev Aug 26, 2026
4a49e4f
Expose external workflow stream APIs
mfateev Aug 26, 2026
2e3e38c
Update external streams design proposal
mfateev Aug 26, 2026
1d7b073
Implement workflow-originated external streams
mfateev Aug 27, 2026
c5f2d8e
Update SDK Core streaming proposal
mfateev Aug 27, 2026
d921fe7
Update Core streaming documentation
mfateev Aug 27, 2026
2efc738
Shared one time-skipping unlock between concurrent result waiters.
moedash Sep 29, 2026
370bfa9
Waited for the target workflow to close before the update that must f…
moedash Sep 29, 2026
6728823
Backported the upstream fixes for the latest dependency set.
moedash Sep 25, 2026
34c5e90
Adopted main's generator scripts and regenerated their output.
moedash Sep 25, 2026
1cff589
Repinned Core to upstream main and brought the bridge up to it.
moedash Oct 3, 2026
89871a2
Pinned the channel Core and regenerated the channel protos.
moedash Oct 3, 2026
02d494d
Added Execution and ExecutionType to temporalio.common.
moedash Oct 3, 2026
f8bf154
Added the notification channel calls to the client.
moedash Oct 3, 2026
7838fa1
Covered the client channel calls.
moedash Oct 3, 2026
34cf5f4
Added channel subscriptions to the workflow runtime.
moedash Oct 3, 2026
f04fe11
Listed a run's channel subscriptions on its description.
moedash Oct 3, 2026
ef6ca30
Covered the workflow channel surface.
moedash Oct 3, 2026
8607599
Pinned the Core that vendors the stream record envelope.
moedash Oct 3, 2026
e9cbce5
Pinned Core at the channel machines on the external base.
moedash Oct 3, 2026
798c1cb
Added the temporalio.streams interface types.
moedash Oct 3, 2026
b558185
Carried a stream provider on client, worker and replayer config.
moedash Oct 3, 2026
3dd3479
Covered the stream interface types.
moedash Oct 3, 2026
bc887e4
Added Execution and ExecutionType to temporalio.common.
moedash Oct 3, 2026
bc77a84
Added MemoryStreams, the in-memory reference provider.
moedash Oct 3, 2026
a6af595
Added the provider conformance suite and ran it on memory.
moedash Oct 3, 2026
8a7e2a6
Added the notification channel calls to the client.
moedash Oct 3, 2026
f3f6dfb
Covered the client channel calls.
moedash Oct 3, 2026
fdd81d3
Added workflow.stream_reader and workflow.stream_writer.
moedash Oct 3, 2026
732e0d2
Ran the stream provider's lifecycle hooks around the workflow.
moedash Oct 3, 2026
7bc1c11
Covered the workflow stream runtime.
moedash Oct 3, 2026
af1ee40
Added channel subscriptions to the workflow runtime.
moedash Oct 3, 2026
0db172e
Listed a run's channel subscriptions on its description.
moedash Oct 3, 2026
ed12a9f
Covered the workflow channel surface.
moedash Oct 3, 2026
d220153
Added the stream accessors for activities and clients.
moedash Oct 3, 2026
d917032
Covered the stream accessors.
moedash Oct 3, 2026
67efbda
Added stream_channel to name the channel a stream notifies.
moedash Oct 3, 2026
5460ca2
Added a changelog entry for the stream interface.
moedash Oct 3, 2026
c4ca2fd
Added a streams demo that runs the agent loop on a provider.
moedash Oct 3, 2026
536c905
Opened hooks in contrib.workflow_streams for a stream provider.
moedash Oct 3, 2026
eb5ebfb
Added the Workflow Streams stream provider.
moedash Oct 3, 2026
99a55aa
Ran the conformance suite on the Workflow Streams provider.
moedash Oct 3, 2026
9193dd7
Added the Temporal Streams Nexus contract and its generated bindings.
moedash Oct 3, 2026
b8a7a3c
Added the Nexus front for stream providers.
moedash Oct 3, 2026
43629ca
Covered the Nexus front.
moedash Oct 3, 2026
19cbdd7
Added a Nexus operation that consumes a stream through its channel.
moedash Oct 3, 2026
1646132
Added a standalone aiohttp host for the stream consumer operation.
moedash Oct 3, 2026
97afa38
Covered the Nexus stream consumer operation.
moedash Oct 3, 2026
3f93403
Warned when a producer's attempt goes backwards on a read.
moedash Oct 3, 2026
bbc5b45
Refused a stream publish from a query handler at the call.
moedash Oct 3, 2026
898c3af
Pinned the stream Core and regenerated the stream api protos.
moedash Oct 3, 2026
2683dd2
Held the server-backed stream cases off the skipping clock.
moedash Oct 3, 2026
ce19b35
Vendored the server's stream service protos.
moedash Oct 3, 2026
6f2f9ad
Added a client for the server-side stream service.
moedash Oct 3, 2026
5220369
Covered the stream service client.
moedash Oct 3, 2026
0b2172f
Made the external stream docstrings build under pydoctor.
moedash Oct 3, 2026
ba5a02a
Shared, retried and codec-aware connections in the stream client.
moedash Oct 3, 2026
2fa3b3c
Added the native stream commands to the workflow runtime.
moedash Oct 3, 2026
6643bae
Added contrib.server_streams over a native stream.
moedash Oct 3, 2026
414e3bf
Added the native stream provider.
moedash Oct 3, 2026
c6a50b8
Replayed workflows that read native streams.
moedash Oct 3, 2026
bedc0e9
Covered replay of native stream reads.
moedash Oct 3, 2026
93f03ea
Added a changelog entry for server-side streams.
moedash Oct 3, 2026
27385ea
Removed a demo run's output from the streams demo.
moedash Oct 3, 2026
f631aed
Merged moe/AI-198-if-py-6-demo into moe/AI-198-if-py-7-workflow-strea…
moedash Oct 3, 2026
f6ab743
Merged moe/AI-198-if-py-7-workflow-streams-provider into moe/AI-198-i…
moedash Oct 3, 2026
8f9dafe
Merged moe/AI-198-if-py-8-nexus-front into moe/AI-198-if-py-9-nexus-c…
moedash Oct 3, 2026
432b7bd
Merged moe/AI-198-if-py-9-nexus-consumer into moe/AI-198-if-py-10-bac…
moedash Oct 3, 2026
a6c758a
Merged moe/AI-198-if-py-10-backwards-attempt-warning into moe/AI-198-…
moedash Oct 3, 2026
79dc5f6
Merged moe/AI-198-if-py-11-query-publish-refusal into moe/AI-198-st-p…
moedash Oct 3, 2026
548d574
Merged moe/AI-198-st-py-1-core-pin into moe/AI-198-st-py-2-native-wire.
moedash Oct 3, 2026
000bd20
Merged moe/AI-198-st-py-2-native-wire into moe/AI-198-st-py-3-native-…
moedash Oct 3, 2026
bbb5159
Merged moe/AI-198-st-py-3-native-provider into moe/AI-198-st-py-4-nat…
moedash Oct 3, 2026
e602e41
Pinned Core at the external-stream repairs.
moedash Oct 3, 2026
28ad8af
Aligned the external stream input and output replay schedules.
moedash Oct 3, 2026
166a6ff
Kept a replay marker from earning a drain of its own.
moedash Oct 3, 2026
b8b3bbb
Allowed external-stream publishes from the workflow constructor.
moedash Oct 3, 2026
c49846c
Pinned Core at the external runtime's channel wake.
moedash Oct 3, 2026
73ab2c6
Woke parked external readers through the stream's channel.
moedash Oct 3, 2026
9dfded4
Covered the external runtime's channel wake.
moedash Oct 3, 2026
0aa1927
Answered a legacy query without the stream snapshot beside it.
moedash Oct 3, 2026
0c63317
Pointed the legacy query id at Core's own constant.
moedash Oct 3, 2026
bce9b70
Added the temporalio.streams interface on the external base.
moedash Oct 3, 2026
579cd92
Covered the stream interface with the conformance suite.
moedash Oct 3, 2026
ac46576
Added the stream interface demo.
moedash Oct 3, 2026
6250dfe
Let an external stream subscription name its start and yield offsets.
moedash Oct 3, 2026
640e594
Added the Redis provider over External Workflow Streams.
moedash Oct 3, 2026
69213bc
Covered the Redis provider.
moedash Oct 3, 2026
44b0b87
Trimmed the Redis streams by retention on every append.
moedash Oct 3, 2026
eff9b2c
Covered the Redis retention trims.
moedash Oct 3, 2026
fb396fe
Held an activity's own streams on the Redis provider.
moedash Oct 3, 2026
a445ad9
Covered activity-owned streams on the Redis provider.
moedash Oct 3, 2026
2896240
Let an external stream subscription start at the tail or the newest N.
moedash Oct 3, 2026
da2588b
Let a Redis workflow reader start at the tail or at the newest N reco…
moedash Oct 3, 2026
e4389c4
Covered the Redis workflow reader's tail starts.
moedash Oct 3, 2026
102cd7e
Hosted standalone streams on the Redis provider.
moedash Oct 3, 2026
3a82166
Covered standalone streams on the Redis provider.
moedash Oct 3, 2026
bc2786c
Aborted a dead output stage when the Worker evicts its run.
moedash Oct 3, 2026
28ad9d9
Covered the dead stage abort on the Redis provider.
moedash Oct 3, 2026
13703e7
Said which provider carries a reader across continue-as-new.
moedash Oct 3, 2026
d2b79e9
Merged the stream_reader docstring fix forward.
moedash Oct 3, 2026
79a50f7
Merged the stream_reader docstring fix forward.
moedash Oct 3, 2026
86bd821
Merged the stream_reader docstring fix forward.
moedash Oct 3, 2026
aade62e
Merged the stream_reader docstring fix forward.
moedash Oct 3, 2026
95d49f6
Merged the stream_reader docstring fix forward.
moedash Oct 3, 2026
f1b65f9
Described the Redis provider's one log per topic in the changelog.
moedash Oct 3, 2026
936f443
Merged the native main chain into the union.
moedash Oct 3, 2026
2fbcd31
Merged the changelog wording fix forward.
moedash Oct 3, 2026
e8a395a
Merged the changelog wording fix forward.
moedash Oct 3, 2026
0d854a9
Merged the changelog wording fix forward.
moedash Oct 3, 2026
71bce94
Merged the changelog wording fix forward.
moedash Oct 3, 2026
aa6c5e0
Merged the changelog wording fix forward.
moedash Oct 3, 2026
7dbff0c
Merged the external chain into the union.
moedash Oct 3, 2026
a2af60d
Merged the time-skipping unlock fix into the union.
moedash Oct 3, 2026
0953f1c
Built the external channel tests' instances the way the worker does.
moedash Oct 3, 2026
cb21ff0
Numbered the Redis producer's records from one.
moedash Oct 3, 2026
adec146
Let a Redis stream handle default to the default topic.
moedash Oct 3, 2026
37e0692
Skipped the native conformance setup without a server address.
moedash Oct 3, 2026
d278a73
Read the constructor-publish case while the worker still polls.
moedash Oct 3, 2026
b4c2d2a
Ran the native provider cases in the union.
moedash Oct 3, 2026
0a9b446
Consumed a Redis stream through the channel its producer notifies.
moedash Oct 3, 2026
75718d3
Ran the demo on each of the union's four providers.
moedash Oct 3, 2026
992452f
Fixed the visitor generator's repeated-field check.
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-*/
2 changes: 1 addition & 1 deletion .gitmodules
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
[submodule "sdk-core"]
path = temporalio/bridge/sdk-core
url = https://github.com/temporalio/sdk-rust.git
url = https://github.com/moedash/sdk-rust.git
109 changes: 109 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,99 @@ 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**: notification channels. Requires a server that serves notification channels.
- `Client.notify_channel`, `poll_channel`, `describe_channel`, `register_channel_listener` and
`unregister_channel_listener` reach a named channel on the server. The channel is either
independent or linked to an execution, which `execution=` (`temporalio.common.Execution`) or
the `workflow_id=` shorthand names.
- A workflow subscribes with `workflow.subscribe_channel(name)`, reads a channel linked to it with
`workflow.linked_channel(name)` and ends a subscription with `unsubscribe()`. Notifications
arrive with the workflow's tasks. `WorkflowExecutionDescription.channel_subscriptions` lists
the channels a run listens on.
- **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, 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, and
`temporalio.streams.providers.redis.RedisStreams` serves the same interface
over External Workflow Streams, with one Redis log per topic that the workflow
and outside readers share.
- **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.
- **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>`.
- **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, 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.nexus.stream_consumer_operation` builds an
asynchronous Nexus operation that consumes a stream. It registers a callback listener on the
stream's notification channel, reads on each delivery and completes when the stream closes.
`temporalio.streams.providers.nexus_consumer_service` hosts it on its own and needs the
`streams-nexus` extra.
- `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 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.

### Changed

Expand All @@ -52,6 +145,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 Down
15 changes: 13 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ cloud-run-worker-otel = [
aioboto3 = ["aioboto3>=10.4.0", "types-aioboto3[s3]>=10.4.0"]
google-genai = ["google-genai>=2.21.0,<3.0.0"]
strands-agents = ["strands-agents>=1.51.0"]
streams-nexus = ["aiohttp>=3.9,<4"]

[project.urls]
Homepage = "https://github.com/temporalio/sdk-python"
Expand Down Expand Up @@ -95,6 +96,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 +127,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 +203,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 +251,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 +278,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