Skip to content

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

Open
moedash wants to merge 5 commits into
moe/AI-198-if-pyext-8-external-interfacefrom
moe/AI-198-if-pyext-9-redis-provider
Open

moedash wants to merge 5 commits into
moe/AI-198-if-pyext-8-external-interfacefrom
moe/AI-198-if-pyext-9-redis-provider

Conversation

@moedash

@moedash moedash commented Oct 3, 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:
    • 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.
    • Resets, refused wakes on a closing run, and an aborted stage's entries leaving the log.
    • Outside reads at END and last=.
    • Activity streams, standalone streams and workflow tail starts raise StreamUnsupportedError. Later PRs fill those in.
  • The stream_reader docstring now says which provider carries a reader across continue-as-new.
  • 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

poe lint is clean at each commit. The external stream suite passes on the dev server. With a local Redis, the streams suite with STREAMS_LIVE=redis passes against a channel server, once with the linked channel kind on and once with it off, so both wake paths ran.

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