Skip to content

Bound Redis to the stream interface as RedisStreams. - #17

Closed
moedash wants to merge 114 commits into
moe/AI-198-py-09-external-interfacefrom
moe/AI-198-py-10-redis-provider
Closed

moedash wants to merge 114 commits into
moe/AI-198-py-09-external-interfacefrom
moe/AI-198-py-10-redis-provider

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR binds Redis to the stream interface as RedisStreams.

What changed?

  • temporalio/streams/providers/redis.py holds both halves. A workflow publish is synchronous and stages with its task, and the completion marker promotes it. An outside append is visible at once.
  • One log per topic. Outside appends and the workflow's promoted batches share it, and the workflow's own entries are dropped from its reads. Cursors are redis:<ms>-<seq>, and the older redis:in: form still reads.
  • An outside append writes its record and its idempotency hash in one Lua script. Repeats compare the plaintext content hash, so a retry through a nonce codec still matches.
  • Retention defaults to seven days, trimmed on every append. retention=None opts out, and a staged batch at or above max_len is refused where it's staged. A batch whose completion the server rejected is aborted at eviction and leaves the log. Eviction never recommits a stage whose marker is in History.
  • A reader starts at END or the newest N. Inside a workflow the worker resolves the tail after the subscribing task and records it in the marker, so replay never asks the log again. A reset run re-reads its inherited ranges and runs the reset-point task again from the log.
  • Activity-owned streams are keyed by the run their execution is in. Standalone streams live under standalone/<stream id> with a policy hash, a seal, and count, age and byte trims.
  • An outside append notifies the stream's channel with the appended entry id as its position, and the counter derives from the id by one rule, ms * 2^20 + min(seq, 2^20 - 1). A wake without a position, the worker's shutdown sweep, takes the same rule applied to the current time with the largest sequence, so it is not folded under the channel's latest. The channel is addressed to the workflow that owns the stream, named as an execution: on a server with linked channels it lives in the reader's own state and nothing subscribes, and on one with independent channels Core subscribes the reader's run when the task that opened the read ends and unsubscribes it when the reader closes, so that task stays retained. RedisStreams(wake_transport=) sets the transport for producers and workers alike: auto, channel or signal.
  • RedisStreams(client=) takes a Redis client the caller owns. CI runs the live cases against a Redis service container.

The tip ends with an empty merge of the replaced #2's head, so #6 and every pin naming it stay valid.

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

Why?

This is the layer where the external transport's shape shows: a staged publish, a notification, and retention without a consumer floor. The two layers below it carry no Redis. Deriving the counter from the entry id lets a producer and a worker order the same notifications the same way, and the server's fold only works if every counter on a channel follows that one rule.

How did you test it?

ruff, mypy, pyright, basedpyright and pydocstyle are clean. STREAMS_LIVE=redis runs tests/streams on the stock dev server, which takes the Signal fallback and never sees the subscribe command, and against the server of moedash/temporal PR #17 named with -E host:port, which takes the channel, in both modes: linked by default and independent with channel.linkedKindEnabled off. One conformance case skips, since Redis refuses a trimmed cursor on the first read, and its own module covers that floor. The live channel case appends past the idle timeout and finds no Signal event in the consumer's History, the subscribed event right after a marker, and the entry id on the notified scheduled event. Its linked variant, marked needs_linked_server, finds no subscribed event and the owner on every notification. Each skips on a server without its kind. The shared cases for a reader leaving its channel run on Redis under wakes_by_notification and needs_unsubscribe_server. The reset cases find their reset point from the markers, since the task that opens the reader differs between servers. The shared channel cases run on Redis through the wakes_by_notification and wakes_by_linked_notification capabilities. The external stream contrib tests pass on the same Redis.

  • Unit Tests
  • Staging
  • End to End Tests

mdashti and others added 30 commits September 5, 2026 20:49
The server refuses the wake while the execution is closing, and at that moment
its status is still the running one, so a single look read "running" and turned
the ordinary ending into an error.
A provider that answers outside readers through handlers on the workflow has
to register them before the first task completes, and one that parks a
long-poll update against the run has to let go before the workflow returns.
Both are no-ops on the storage providers, so workflow code calls them
unconditionally and stays portable.
Porting the agent harness onto the interface surfaced three things the small
example never needed. A reader resumes after the record a cursor names, so
it can store the last cursor it handled without advancing an opaque token. A
consumer reports its latest position, so a client can follow a turn it is
about to start without the workflow reporting a position. And a producer can
append onto a topic of the stream the workflow publishes, so an activity's
live output lands next to the workflow's own records for one outside reader.
Porting the agent harness onto the interface surfaced three things the small
example never needed. A reader resumes after the record a cursor names, so
it can store the last cursor it handled without advancing an opaque token. A
consumer reports its latest position, so a client can follow a turn it is
about to start without the workflow reporting a position. And a producer can
append onto a topic of the stream the workflow publishes, so an activity's
live output lands next to the workflow's own records for one outside reader.
A replay marker carries what its Workflow Task published, and installing it
resets the publish records. The activation applied it with the non-query
jobs, after the signal and update set had already run and published, so a
task whose publishes a signal woke replayed as zero records against a
manifest of several, and any query against the completed run failed. The
marker now goes in with the first set that drains, after that set's other
jobs and before the drain.
CI runs `ruff check --select I` and `ruff format --check` over the whole
tree, and the new files were written by hand.
`prepare` and `drain` are hooks only a transport that parks something
against the running workflow needs, so they move to their own protocol
rather than making every provider carry two empty methods. The demo's
teardown called a `close` that no provider setup defined.
The provider imports nothing from `workflow`, and its producer, consumer
and `latest` need the docstrings the doc linter asks every public method
for.
A record only reaches a reader after a provider placed it, so the fake
runtime hands out placed records and the reader can name where each one
came from. The subscription manager captures the loop it is built on,
which is the Worker's, so its fixture has to build it on one.
Every Workflow Task that closes commits its own marker, and empty input
activations now close one too, so a baseline read three steps earlier
counts markers that belong to tasks this case says nothing about.
The changelog checkpoint requires an entry for any user-facing change.
A transport that serves outside readers through handlers on the workflow
registers them in `prepare`, so without the call an outside reader finds
no handler and the demo's follow fails against that provider.
* Support Google ADK 2.7 streaming import

* Fix tests with latest dependencies

* Apply suggestion from @brianstrauch
The branch's own copy shells out to a binary name the nexgen crate does not install, so generation died before it could diff anything.
Under the latest dependency set openai carries its own httpx as httpx2, so pyright will not accept the httpx.Response the test builds. Going through the class object was tried first and did not silence it.
The same suite on the same Core is 673 passed and 0 failed against a dev server, and fails under time skipping. Which cases fail drifts between fourteen and seventeen across runs, so the suite is held as a whole rather than by a list of names that was never stable.
A script that drives two tool calls handed both invocations the same completed id, which the Agents SDK now rejects outright rather than tolerating. This diverges from upstream, which still hands out a constant there; unique ids are strictly more correct for a builder handing out completed call ids, so the right long-term home is upstream rather than here.
That branch now matches upstream's nexus model WIT for the workflow-id policies, which is what the current generator needs.
The generator needs its matching support file and the upstream lint exclusions, so both came across with it. One test asserted the old model's field names; the only consumers of that model outside the generated package are a generated re-export in temporalio.workflow and three package-level helpers, none of which touch those fields.
The generator renames the service module, and the visitor script still looked for the old filename, so the generation sequence died after the model was written.
pydoctor resolves single-backtick prose as symbol references, so shell words, tuples and file names were read as missing targets. Two simple tables also lost their column alignment because pydoctor rewrites role references before parsing, so one is now a literal block.
A workflow's publish runs on the workflow thread, so waking a consumer parked on another loop needs call_soon_threadsafe. The memory provider also stops reading idle_timeout as its poll period, which gave the parameter a second meaning.
A workflow reader took the beginning or a cursor only. The tail is where the
log is when the worker looks, which the workflow thread cannot see, so the
transport records the request, the worker resolves it against the log after the
task that opened the subscription and writes the entry into the marker beside
it, and replay reads it from History. Outside, the tail is resolved on the
generator's first step, as the retention check already is.
…xt hash.

A stream with an id and no owner gets its own key scheme: a hash for the policy
and the seal, a log per topic, and an append script that reads the policy,
refuses a sealed stream and trims by count, age and a per-topic byte total. Every
record now carries the hash of its converted body and the append scripts compare
that hash, so a retry through a codec that differs on every call still matches.
… trim.

The conformance suite now holds a provider to the standalone capabilities it
declares. The Redis append script keeps a byte total per topic and trims by age
on every append, so both hold while the stream is open.
A completion the server rejects, such as a closing task with a wake Signal
buffered under it, leaves the batch it staged with a token no marker will name;
the re-run stages the batch again, and the dead stage stayed in the log once a
reader or the worker aborted it. On the one-log layout those entries count
against every window and bound, so the abort now deletes them and the stage's
hashes stand in for them, and the worker settles a run's pending stages against
History when it evicts the run rather than leaving them for a reader to find.
The eviction-time reconcile also recommitted a stage whose marker is authoritative but whose
one promotion had failed, which took the repair away from the cold client and broke the
post-report commit failure contract. Eviction now applies an abort decision only, and a stage
History has committed stays pending for the client.
…ace' into moe/AI-198-py-10-redis-provider

# Conflicts:
#	temporalio/contrib/external_workflow_streams/_backend.py
The counter derives from the id so producers and workers order wakes alike, and the server's NOT_FOUND ends a wake without a describe.
…ace' into moe/AI-198-py-10-redis-provider

# Conflicts:
#	tests/streams/test_streams_conformance.py
moedash added 17 commits October 1, 2026 18:18
The live case asks for the channel outright and checks the subscribed event, the notified scheduled event with the entry id as position, and no Signal. The shared conformance case runs on Redis through the wakes_by_notification capability.
A task index names a different task on a server with channels, where the task that opens the reader carries the subscribe command and ends at once.
A wake without a position takes the current time with the largest sequence, so it orders above every entry appended before now and is not folded under the channel's latest.
The provider passes the owner through the shared channel address, so nothing changes in it beyond the transport's documentation; the live channel case gains a linked variant and the case on the removed transport goes.
The channel case finds the subscribed event after a marker, so the task that
opened the reader stayed retained until it left. The provider case's note
says the same.
…-10-redis-provider

# Conflicts:
#	tests/streams/test_streams_conformance.py
The notified scheduled event names the owner as a
`temporal.api.common.v1.Execution`, so the case checks the type and the
business id through the public helper.
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36, #37, #38, #39, #40, #41, #42, #43, #44, #45, #46, #47, #48, #49, #50, #51, #52, #53, #54 and #55.

Same content, split into 37 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

@moedash moedash closed this Oct 3, 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.

3 participants