Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 6 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,12 @@ 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, with one Redis log per topic that the workflow
and outside readers share.
- `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()
34 changes: 32 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,8 @@
classify_read_failure,
)
from temporalio.contrib.external_workflow_streams._record import (
Cursor,
Offset,
StreamRecord,
)
from temporalio.contrib.external_workflow_streams._wake import channel_for
Expand Down Expand Up @@ -95,6 +97,7 @@ def register(
wait_id: int,
stream_key: StreamKey,
idle_timeout: timedelta,
start_cursor: Cursor | None = None,
) -> None:
"""Registers a wait with the Worker's subscription manager.

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

def subscribe(self) -> ExternalStreamSubscription[AnyType]:
def subscribe(
self, *, start_cursor: Cursor | 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 +291,13 @@ 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.
"""
state = _run_state()
if state.runtime is None:
Expand All @@ -297,6 +309,11 @@ 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}
)
state.runtime.register(
wait_id=wait_id,
stream_key=stream_key,
Expand All @@ -307,6 +324,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 +629,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 +668,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
Loading
Loading