Conversation
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.
…-10-redis-provider
… 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
…ace' into moe/AI-198-py-10-redis-provider # Conflicts: # tests/streams/test_streams_conformance.py
2 of 3 tasks
…-10-redis-provider
…-10-redis-provider
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.
…-10-redis-provider
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.
…-10-redis-provider
…-10-redis-provider
…-10-redis-provider
…-10-redis-provider
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.
…-10-redis-provider
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. |
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 binds Redis to the stream interface as
RedisStreams.What changed?
temporalio/streams/providers/redis.pyholds 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.redis:<ms>-<seq>, and the olderredis:in:form still reads.retention=Noneopts out, and a staged batch at or abovemax_lenis 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.ENDor the newestN. 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.standalone/<stream id>with a policy hash, a seal, and count, age and byte trims.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,channelorsignal.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,basedpyrightandpydocstyleare clean.STREAMS_LIVE=redisrunstests/streamson 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 withchannel.linkedKindEnabledoff. 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, markedneeds_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 underwakes_by_notificationandneeds_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 thewakes_by_notificationandwakes_by_linked_notificationcapabilities. The external stream contrib tests pass on the same Redis.