Repository navigation
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 5, 2026
Owner
Author
|
Replaced by #106 in the v5 series (phase 1, Signal wake only). The branch stays as a pin. |
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?
The transport gains two things a provider needs.
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.RedisStreamskeeps 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 throughwake_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 atENDandlast=. Activity streams, standalone streams and workflow tail starts raiseStreamUnsupportedErrorfor 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_readerdocstring 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 forrecords(). The conformance suite registers the Redis case underSTREAMS_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 -
The lint set is clean. The external stream suite and
tests/streamspass on the dev server the fixtures start and with-Eagainst a local channel server. With a local Redis,tests/streamswithSTREAMS_LIVE=redispasses 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.