Conversation
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.
This was referenced Oct 3, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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?
subscribe(start_cursor=)seeds a workflow reader from a cursor, andExternalStreamSubscription.records()yields each value with its offset. The test fakes now place records the way a provider does, sincerecords()needs an offset on each one.temporalio.streams.providers.redis.RedisStreams:wake_transport, with counters that follow the entry id order.ENDandlast=.StreamUnsupportedError. Later PRs fill those in.stream_readerdocstring now says which provider carries a reader across continue-as-new.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 -
poe lintis clean at each commit. The external stream suite passes on the dev server. With a local Redis, the streams suite withSTREAMS_LIVE=redispasses against a channel server, once with the linked channel kind on and once with it off, so both wake paths ran.