Skip to content

Added the Redis stream provider over External Workflow Streams. - #106

Open
moetemp wants to merge 6 commits into
moe/AI-198-s1-py-16-parked-run-queryfrom
moe/AI-198-s1-py-17-redis-provider
Open

moetemp wants to merge 6 commits into
moe/AI-198-s1-py-16-parked-run-queryfrom
moe/AI-198-s1-py-17-redis-provider

Conversation

@moetemp

@moetemp moetemp commented Oct 6, 2026

Copy link
Copy Markdown
Owner

This PR adds RedisStreams, a stream provider over External Workflow Streams that keeps a workflow's streams in a Redis the customer runs.

What changed?

The transport gains two things a provider needs. subscribe(start_cursor=) seeds a workflow reader from a cursor, and ExternalStreamSubscription.records() yields each value with its offset. The test fakes now place records the way a provider does, since records() needs an offset on each one.

temporalio.streams.providers.redis.RedisStreams keeps one Redis log per topic, shared by the workflow and outside readers. A workflow's publishes are staged with its task and promoted once the completion marker proves the task was accepted. An outside append wakes a parked reader with the transport's reserved Signal. A wake the server refuses because the run is closing is sent again for a short window, then dropped, since the run's next park checks the log again. A NOT_FOUND says the chain has ended, with no describe needed. A retry is matched by the content hash of its plaintext body. It handles resets and an aborted stage's entries leaving the log. Outside reads start at END and last=. Activity streams, standalone streams and workflow tail starts raise StreamUnsupportedError for now, and the next PRs fill them in.

The provider speaks the shared wire from the start. Its producer numbers records from one, since zero says a producer doesn't number its records, and its handle takes the default topic when a call names none.

The stream_reader docstring now says which provider carries a reader across continue-as-new. It used to say nothing crosses continue-as-new, which isn't true for Redis. The changelog gets the Redis clause on the interface entry and a line for records(). The conformance suite registers the Redis case under STREAMS_LIVE=redis, a live module covers it against a server and a Redis, streams_demo/provider_setup.py gains a redis choice, and CI gets a Redis service.

Part of AI-198 (epic AI-37).

Why?

A Redis store is the provider customers asked for first. The workflow's writes go through the transport's staged commit, so an outside reader never sees a batch from a task History rejected. Replay reads the ranges the markers recorded.

The earlier version took a wake_transport argument that picked between a channel notification and the Signal, with wake counters that followed the entry id order. The notification channel is paused per the 2026-10-05 design review, so the argument and the counters are gone and the Signal is the only wake. The channel wake tests and the conformance wake and leave cases go with them.

How did you test it?

Link to a test plan if any -

  • Unit Tests
  • Staging
  • End to End Tests

The lint set is clean. The external stream suite and tests/streams pass on the dev server the fixtures start. With a local Redis, the Redis lane (tests/streams with STREAMS_LIVE=redis, the conformance case and the replay module included) passes against a stock dev server.

A provider over the transport needs both: a reader seeded from a cursor starts
after it, and records() hands back the offset each value was read from. The
test fakes now place records the way a provider does, since records() needs an
offset on every one.
RedisStreams serves the stream interface from a store the customer runs: one
log per topic, a workflow's publishes staged with its task and promoted by the
completion marker, outside appends that wake the reader over the channel or the
Signal, and a content-hash match for retries. CI gets a Redis service.
The conformance suite runs the Redis case when STREAMS_LIVE=redis, and the live
module checks the staged commit, replay, resets and the channel wake against a
server and a Redis.
The stream_reader docstring said nothing crosses continue-as-new, which the
Redis provider contradicts: its transport resumes a successor's reader where
the predecessor committed.
The entry still said each topic was an input and an output stream, which the
provider stopped doing when it moved to one shared log per topic.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant