Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
114 commits
Select commit Hold shift + click to select a range
75de955
Align external stream input and output replay schedules
mdashti Sep 6, 2026
c00339e
Preserve stream task retention with zero workflow cache
mdashti Sep 6, 2026
632e16f
Put the shared stream interface on the repaired external SDK.
moedash Sep 14, 2026
66f353c
Let a terminal record land on a consumer that is closing.
moedash Sep 6, 2026
f06271e
Added the temporalio.streams interface with a provider registry.
moedash Sep 15, 2026
3dd101f
Added the in-memory reference provider and the conformance tests.
moedash Sep 15, 2026
62be804
Added the shared demo loop.
moedash Sep 15, 2026
7964e49
Moved the client-side binding into the provider registry.
moedash Sep 15, 2026
af2483d
Added the prepare and drain hooks the Workflow Streams transport needs.
moedash Sep 16, 2026
0741345
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
fb74e79
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
295fe3f
Applied an external stream replay marker before the task's first drain.
moedash Sep 16, 2026
668e432
Sorted and formatted the streams package for ruff.
moedash Sep 16, 2026
0877fe4
Made the streams package pass the type and doc linters.
moedash Sep 16, 2026
29cc1cb
Made the Redis provider pass the linters.
moedash Sep 16, 2026
89c1a37
Gave the external stream test fakes the shapes the real path has.
moedash Sep 16, 2026
97849a8
Snapshotted the markers where the abandoned task begins.
moedash Sep 16, 2026
d963173
Added the changelog entry for the shared stream interface.
moedash Sep 16, 2026
5bf0094
Made the demo call the provider lifecycle hooks.
moedash Sep 16, 2026
105ad9b
Support google-genai 2.21 file downloads (#1865)
brianstrauch Sep 15, 2026
e69ad05
Fix tests with latest dependencies (#1798)
brianstrauch Aug 31, 2026
914ee40
Adopted main's Nexus system API generator script.
moedash Sep 17, 2026
14b5042
Loosened the status error's response annotation.
moedash Sep 17, 2026
c4bed86
Held the external stream suite off the time-skipping server.
moedash Sep 17, 2026
86dc631
Gave each mock tool call its own id.
moedash Sep 17, 2026
c27c620
Repinned Core to the finalized external repairs branch.
moedash Sep 17, 2026
6b34091
Regenerated the Nexus system API from the repinned Core.
moedash Sep 17, 2026
9437d75
Adopted main's payload visitor generator.
moedash Sep 17, 2026
e6007d3
Made the external stream docstrings build under pydoctor.
moedash Sep 17, 2026
a3f50cd
Woke outside readers across threads and ignored idle_timeout in memory.
moedash Sep 18, 2026
528b811
Named the last record in append's cursor and allowed None for it.
moedash Sep 18, 2026
217497f
Resolved producer identity once in the package.
moedash Sep 18, 2026
8495140
Required an explicit provider and exported instance.
moedash Sep 18, 2026
ad177b2
Typed the handles and records honestly.
moedash Sep 18, 2026
952adfb
Kept topics off inbound frames and skipped unreadable frames with a w…
moedash Sep 18, 2026
276cb15
Let a broken provider import surface instead of vanishing from the re…
moedash Sep 18, 2026
785402d
Added the workflow-side conformance tests.
moedash Sep 18, 2026
7325643
Replaced stale wording and renamed the demo types.
moedash Sep 18, 2026
0c7616c
Keyed stores by one helper that survives a colon in the workflow id.
moedash Sep 19, 2026
8368735
Added an async close hook next to prepare and drain.
moedash Sep 19, 2026
e3522ae
Parametrised the conformance suite over providers.
moedash Sep 19, 2026
31f6fc0
Let a provider setup own the workflow the conformance cases address.
moedash Sep 19, 2026
f507eef
Recorded whether a provider reports a dropped repeat.
moedash Sep 19, 2026
17bbcfa
Dropped the committed demo output and corrected two comments.
moedash Sep 19, 2026
f71a030
Adapted the Redis provider to the interface changes.
moedash Sep 19, 2026
6fb7abe
Registered the Redis provider in the conformance suite.
moedash Sep 19, 2026
05ac20d
Ran the interface loop and the replay query in one Redis live module.
moedash Sep 19, 2026
e1d1b22
Renamed the store key helper to inbound_stream_id for the provider br…
moedash Sep 19, 2026
8b9f110
Repinned Core to the reviewed moe/AI-198-external-core-repairs head.
moedash Sep 19, 2026
362b643
Added the StreamError family for stream conditions.
moedash Sep 21, 2026
600f6d5
Vendored the stream additions to the public API protos.
moedash Sep 21, 2026
4362c31
Split the stream provider in two and made the proto the record.
moedash Sep 21, 2026
e05e651
Carried the stream provider from the worker into the workflow runtime.
moedash Sep 21, 2026
d37f8ad
Moved the stream conformance suite onto the new surface.
moedash Sep 21, 2026
1ea907c
Moved the stream demo and changelog entry onto the new surface.
moedash Sep 21, 2026
2659a17
Rewrote the Redis provider with an input and an output stream per topic.
moedash Sep 21, 2026
88a6298
Kept the stream finish hook off evicted runs and collected coroutines.
moedash Sep 21, 2026
d11aa01
Typed the hook test's workflow input as awaitable.
moedash Sep 21, 2026
4995cf1
Added activity.stream_handle and Client.get_stream_handle.
moedash Sep 21, 2026
dc32947
Registered the demo's provider on the client and reported from the ac…
moedash Sep 21, 2026
a667640
Registered the Redis conformance provider on the client.
moedash Sep 21, 2026
be743fd
Added typed topic definitions.
moedash Sep 21, 2026
c37b351
Took topic definitions on the Redis handle.
moedash Sep 21, 2026
a9bda0f
Matched the Redis live module to typed topic definitions.
moedash Sep 21, 2026
62b57e3
Seeded a Redis workflow reader from a cursor.
moedash Sep 21, 2026
c7e30af
Passed a backend's integrity loss through the replay read.
moedash Sep 21, 2026
d87787d
Trimmed the Redis streams by retention on every append.
moedash Sep 21, 2026
4b6d572
Listed stream_provider in ClientConnectConfig.
moedash Sep 22, 2026
6ec392f
Answered a legacy query without the stream snapshot beside it.
moedash Sep 22, 2026
d3aeb56
Repinned Core to the repairs head that carries the stream protos.
moedash Sep 22, 2026
b7be7b4
Added the Redis provider with an input and an output stream per topic.
moedash Sep 25, 2026
e578e67
Seeded a Redis workflow reader from a cursor.
moedash Sep 25, 2026
84c77d4
Trimmed the Redis streams by retention on every append.
moedash Sep 25, 2026
2a5593c
Tied the series to the original branch head.
moedash Sep 25, 2026
66eaaff
Merged the external interface with the review-round fixes.
moedash Sep 25, 2026
bb2d7ca
Repaired the Redis provider's writes, refusals and error family.
moedash Sep 25, 2026
17a308e
Ran the Redis stream tests in CI against a service container.
moedash Sep 25, 2026
aaccdda
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Sep 26, 2026
0380e06
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Sep 29, 2026
5b9fbda
Held an activity's own streams on the Redis provider.
moedash Sep 30, 2026
aa67339
Tested activity-owned streams on the Redis provider.
moedash Sep 30, 2026
2c19822
Sent a refused wake again along the chain instead of failing the sender.
moedash Sep 30, 2026
bf99f74
Trimmed Redis stream keys to seven days unless told otherwise.
moedash Sep 30, 2026
b51c3e3
Keyed an activity's Redis streams by the run its execution is in.
moedash Sep 30, 2026
84e268c
Kept one Redis log per topic instead of an input and an output key.
moedash Sep 30, 2026
1acece8
Carried a Redis consumer through a reset and dropped a wake a closing…
moedash Sep 30, 2026
708ed5c
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Sep 30, 2026
10d56a7
Let a Redis reader start at the tail or at the newest N records.
moedash Sep 30, 2026
5a22233
Hosted standalone streams on Redis and matched retries by the plainte…
moedash Sep 30, 2026
4ee6ed5
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Sep 30, 2026
cbe1ec4
Declared the Redis standalone stream's byte bound and open-stream age…
moedash Sep 30, 2026
67f80ea
Took an aborted batch's entries out of the shared log.
moedash Oct 1, 2026
78aeddd
Aborted only History-rejected output stages at eviction.
moedash Oct 1, 2026
bd7bc95
Merge remote-tracking branch 'origin/moe/AI-198-py-09-external-interf…
moedash Oct 1, 2026
78454f4
Woke Redis readers with the entry id as the wake position.
moedash Oct 1, 2026
fccdaae
Merge remote-tracking branch 'origin/moe/AI-198-py-09-external-interf…
moedash Oct 1, 2026
db43568
Merge remote-tracking branch 'origin/moe/AI-198-py-09-external-interf…
moedash Oct 1, 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
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
5157b03
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
401616e
Covered the Redis reader woken through its linked channel.
moedash Oct 2, 2026
505f030
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
9173526
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
moedash Oct 2, 2026
767475a
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
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
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
01d9773
Merge branch 'moe/AI-198-py-09-external-interface' into moe/AI-198-py…
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
6 changes: 5 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,11 @@ to include examples, links to docs, or any other relevant information.
`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.
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.

- Added experimental External Workflow Streams in
`temporalio.contrib.external_workflow_streams`. Workflow stream payloads are
Expand Down
26 changes: 17 additions & 9 deletions streams_demo/provider_setup.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
"""Pick the provider for a demo run from the environment.

``STREAMS_PROVIDER`` names a provider; this base tree carries only ``memory``,
and each provider branch adds its own name here. The demo needs a Temporal
server to run the workflow either way; ``TEMPORAL_ADDRESS`` points at it.
``STREAMS_PROVIDER`` names a provider; this tree carries ``memory`` and
``redis``. The demo needs a Temporal server to run the workflow either way;
``TEMPORAL_ADDRESS`` points at it, and the Redis demo reads its store from
``AI198_REDIS_URL`` and ``AI198_REDIS_PREFIX``.
"""

from __future__ import annotations
Expand All @@ -21,16 +22,23 @@

async def open() -> tuple[str, ProviderPlugin]:
"""The server to connect to and the provider the worker and the client share."""
if NAME != "memory":
raise SystemExit(f"this tree carries no stream provider named {NAME!r}")
return os.environ.get("TEMPORAL_ADDRESS", "localhost:7233"), MemoryStreams()
address = os.environ.get("TEMPORAL_ADDRESS", "localhost:7233")
if NAME == "memory":
return address, MemoryStreams()
if NAME == "redis":
from temporalio.streams.providers.redis import RedisStreams

return address, RedisStreams(
url=os.environ.get("AI198_REDIS_URL", "redis://127.0.0.1:6379"),
key_prefix=os.environ.get("AI198_REDIS_PREFIX", "ai198-contract"),
)
raise SystemExit(f"this tree carries no stream provider named {NAME!r}")


async def close(provider: ProviderPlugin) -> None:
"""Let go of whatever :func:`open` acquired.

The memory provider holds no connection, so this is its ``close()`` and
nothing more. A provider branch that opens one closes it the same way, so
the demo's teardown reads the same on every provider.
The memory provider holds no connection and the Redis provider closes
the client it opened, so the demo's teardown reads the same on both.
"""
await provider.close()
2 changes: 2 additions & 0 deletions temporalio/contrib/external_workflow_streams/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@
Offset,
OffsetComparator,
RecordKind,
StartAtTail,
StreamRecord,
)
from temporalio.contrib.external_workflow_streams._wake import WakeRequest
Expand Down Expand Up @@ -128,6 +129,7 @@
"ExternalStreamProducerTopic",
"ExternalStreamSubscription",
"ExternalStreamTopic",
"StartAtTail",
"IdempotencyKey",
"Offset",
"OffsetComparator",
Expand Down
46 changes: 44 additions & 2 deletions temporalio/contrib/external_workflow_streams/_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,9 @@
classify_read_failure,
)
from temporalio.contrib.external_workflow_streams._record import (
Cursor,
Offset,
StartAtTail,
StreamRecord,
)
from temporalio.contrib.external_workflow_streams._wake import channel_for
Expand Down Expand Up @@ -95,6 +98,8 @@ def register(
wait_id: int,
stream_key: StreamKey,
idle_timeout: timedelta,
start_cursor: Cursor | None = None,
start_at_tail: StartAtTail | None = None,
) -> None:
"""Registers a wait with the Worker's subscription manager.

Expand Down Expand Up @@ -275,7 +280,12 @@ class ExternalStreamTopic(Generic[AnyType]):
value_type: type[AnyType] | None
options: ExternalStreamOptions

def subscribe(self) -> ExternalStreamSubscription[AnyType]:
def subscribe(
self,
*,
start_cursor: Cursor | None = None,
start_at_tail: StartAtTail | None = None,
) -> ExternalStreamSubscription[AnyType]:
"""Starts a new subscription and returns its async iterator.

Each call is an **independent** subscription with its own ``wait_id``,
Expand All @@ -286,6 +296,18 @@ def subscribe(self) -> ExternalStreamSubscription[AnyType]:
the same hazard class as timers and activities: inserting, removing, or
reordering a ``subscribe()`` call renumbers every later wait in the Run
and must be gated behind ``workflow.patched()``.

Args:
start_cursor: The boundary the subscription begins after. ``None``
resumes where the predecessor Run committed this wait, or at
``BEGINNING`` on a first execution. A boundary the Workflow
names is recorded in the marker's header like the restored one,
so it must be derived deterministically: replay names it again.
start_at_tail: Start at the stream's tail instead, or at its newest
``last`` records. The Worker resolves the boundary against the
store after this Workflow Task and records it with the
subscription, so replay starts where the live run did rather
than asking the store again. Exclusive with ``start_cursor``.
"""
state = _run_state()
if state.runtime is None:
Expand All @@ -297,6 +319,13 @@ def subscribe(self) -> ExternalStreamSubscription[AnyType]:
state.next_wait_id += 1

stream_key = state.runtime.stream_key(self.name)
# Only a named boundary travels; the runtime derives the default itself
# from the predecessor Run's continuation.
start: dict[str, Any] = (
{} if start_cursor is None else {"start_cursor": start_cursor}
)
if start_at_tail is not None:
start["start_at_tail"] = start_at_tail
state.runtime.register(
wait_id=wait_id,
stream_key=stream_key,
Expand All @@ -307,6 +336,7 @@ def subscribe(self) -> ExternalStreamSubscription[AnyType]:
# `with_options` was given, so no configured value can ever reach
# the reduction and every set parks after one second.
idle_timeout=self.options.idle_timeout,
**start,
)
# The channel the stream's writers notify. Asked of the SDK's object for
# the Run rather than of the stream runtime, because the answer is part
Expand Down Expand Up @@ -611,6 +641,16 @@ def __aiter__(self) -> AsyncIterator[AnyType]:
return self._iterate()

async def _iterate(self) -> AsyncIterator[AnyType]:
async for _offset, value in self.records():
yield value

async def records(self) -> AsyncIterator[tuple[Offset, AnyType]]:
"""Each value with the provider offset it was read from.

Same delivery, same consumption, same commit order as iterating values.
A reader that has to name where it got to, or hand a position to
something outside the Workflow, cannot do it from the values alone.
"""
while not self._finished:
# Re-filled before *every* record rather than once per batch. A
# record buffered while Workflow code was doing something else -- a
Expand Down Expand Up @@ -640,8 +680,10 @@ async def _iterate(self) -> AsyncIterator[AnyType]:
# becomes a value or raises, and neither outcome can leave the
# ready list half-consumed.
value = self._decode(record)
offset = record.offset
self._commit(record)
yield value
assert offset is not None
yield offset, value
continue
await self._await_readiness()

Expand Down
17 changes: 17 additions & 0 deletions temporalio/contrib/external_workflow_streams/_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -275,6 +275,23 @@ def compare_offsets(self, left: Offset, right: Offset) -> int:
inside a validation loop.
"""

async def tail_cursor(self, key: StreamKey, *, before_last: int = 0) -> Cursor:
"""The boundary the newest ``before_last`` records begin after.

With ``before_last=0`` it is the boundary after the newest record, so a
read from it sees only what is appended after this call. ``BEGINNING``
when the stream holds fewer records than asked for. Resolved on the
Worker, never on the Workflow thread, and recorded with the
subscription so replay reads it from the marker instead of asking
again. A provider that cannot answer leaves this as it is, and a
subscription that asks for a tail start fails when the Worker resolves
it.
"""
raise NotImplementedError(
f"{type(self).__name__} cannot resolve a start {before_last} records "
f"before the tail of {key}; subscribe with a cursor instead"
)

# --- waking -------------------------------------------------------------

wake_transport: WakeTransport = "auto"
Expand Down
22 changes: 22 additions & 0 deletions temporalio/contrib/external_workflow_streams/_record.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,29 @@
from dataclasses import dataclass, field
from typing import Final, Protocol


@dataclass(frozen=True)
class StartAtTail:
"""A subscription start the Worker resolves against the store: the tail.

``last`` is how many of the newest records the read begins with; zero
means the boundary after the newest record, so the read sees only what
is appended after it was opened. The Workflow thread cannot ask the store
where the tail is, so the runtime records the request and the Worker
resolves it before the watcher starts, writing the boundary into the
marker beside the subscription as it does a cursor the Workflow named.
"""

last: int = 0

def __post_init__(self) -> None:
"""Refuse a negative count."""
if self.last < 0:
raise ValueError(f"last must not be negative, got {self.last}")


__all__ = [
"StartAtTail",
"AFTER",
"BEGINNING",
"Cursor",
Expand Down
4 changes: 3 additions & 1 deletion temporalio/contrib/external_workflow_streams/_replay.py
Original file line number Diff line number Diff line change
Expand Up @@ -252,7 +252,9 @@ async def _read_range(
"""
try:
return await backend.read_range(key, run.first_offset, run.last_offset)
except StreamStorageError:
except (StreamStorageError, StreamIntegrityError):
# A backend that can tell the range is gone reports the loss itself;
# wrapping it would file a permanent loss under a transient failure.
raise
except Exception as err:
# Not integrity loss: nothing has been shown to be missing, only
Expand Down
Loading
Loading