Skip to content

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

Closed
moetemp wants to merge 6 commits into
moe/AI-198-o7-py-7-channel-wakefrom
moe/AI-198-o7-py-8-redis-provider
Closed

moetemp wants to merge 6 commits into
moe/AI-198-o7-py-7-channel-wakefrom
moe/AI-198-o7-py-8-redis-provider

Conversation

@moetemp

@moetemp moetemp commented Oct 5, 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. Outside appends wake the reader through wake_transport, with counters that follow the entry id order. A retry is matched by the content hash of its plaintext body. It handles resets, refused wakes on a closing run, 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, 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.

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 and with -E against a local channel server. With a local Redis, tests/streams with STREAMS_LIVE=redis passes against that server with the linked channel kind on, and the Redis conformance and replay modules pass with it off, so both wake paths ran. The wake and leave conformance cases run on Redis here.

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.
@moetemp

moetemp commented Oct 6, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #106 in the v5 series (phase 1, Signal wake only). The branch stays as a pin.

@moetemp moetemp closed this Oct 6, 2026
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