Repository navigation
Conversation
The test server holds one lock per in-flight task and counts one unlock per request, so a second waiter let the clock skip while another workflow still had a task in flight.
…ail. The last update's result does not mean the target workflow has closed. An update that lands before the closure stays pending until pytest's timeout instead of failing, which hung test_temporal_operation_update_workflow on two Linux lanes. Ported from temporalio#1810.
The series builds its Core layers on current upstream sdk-rust main, so the bridge moves to the 1.0 crates and the protos follow upstream's api drift.
The pinned Core carries the notification channel api, the bridge job and commands, and the five channel RPCs, so the regen brings them into Python. The submodule now names the fork, since upstream does not have the commit.
A linked channel names its owner as an execution, a workflow or a standalone activity, so the client needs a typed value for it.
The five calls reach an independent channel or one linked to an execution. Polls return workflow.Notification, so the workflow side shares the type later.
The unit case fakes the service. The live cases skip unless -E names a server that serves channels.
A workflow subscribes by command or listens on its linked channel, and Core hands the notifications over as a job. The handle routes them by kind and ends with unsubscribe().
Describe is how an operator sees what a run listens on and what is pending for it.
The unit cases drive the instance with hand-built activations. The live cases skip unless -E names a channel server.
The interface's record on the wire is temporal.api.stream.v1.StreamRecord, and the regen from this pin brings that module in.
One record, topic, cursor, error family and provider protocol set that every provider and every context shares.
A provider is registered once as a plugin and the contexts that ask for a stream find it in their config.
Wire round trips, supersession, topic keys, cursors, body encoding and refs, with no provider involved.
It holds streams in process memory, so the interface has a provider the conformance suite and unit tests can run without a server.
The suite is parametrised by provider, and capability markers skip what a provider does not claim. Every later provider registers here.
Workflow code reads and writes a stream through the provider's workflow half, made once per instance so its state dies with the instance.
The start hook runs on the workflow's loop before the first task's handlers, so a provider's handlers exist for an Update that arrives with that task. The finish hook skips eviction.
The hooks, the reader's start positions and a workflow reading and writing through the memory provider.
activity.stream_handle and client.get_stream_handle open a stream by ref or by topic, and create_stream makes a standalone one, so every context reaches a stream the same way.
Activity and client handles on memory, and the conformance cases for standalone streams.
The server derives the channel from the stream's identity, and the client derives the same one, so a listener needs no lookup.
One agent loop, read, decide, publish and an activity, that runs unchanged on every provider, so a provider branch only swaps provider_setup.
The provider adopts the shipped log and its handlers rather than copying them, and names the missing-query failure by a constant.
It serves the stream interface over the shipped Workflow Streams transport, so a workflow's History holds the records and no other store is needed.
Plus the provider's own cases for publish dedupe, continue-as-new, reset and activity streams, on the stock dev server.
The contract is a nexgen yaml, so any language nexgen targets gets the same append and read operations. gen-protos regenerates the bindings.
NexusStreams appends and reads through one endpoint, and TemporalStreamsHandler serves it in front of any storage provider, so callers never name the store.
The handler and caller cases on the dev server's Nexus endpoints, and the front as a conformance case.
The shared activity suite runs on Redis when STREAMS_LIVE=redis, and the Redis module checks the key scheme, the run keying and when a read ends.
The workflow thread cannot see where the store's tail is, so the transport records the request, the Worker resolves it against the store after the task that opened the subscription, and the marker records the boundary so replay reads it from History.
…rds. stream_reader(after=END) and stream_reader(last=N) on the Redis provider now ask the transport for a tail start instead of being refused.
A reader at the newest N records and one at END each start where the live run resolved, and replay starts there too.
…structor-publish.
…-parked-run-query.
…is-activity-owners.
…y-11-redis-start-positions.
A stream with an id and no owner gets its own keys: a hash for the policy and the seal, and a log per topic. One append script reads the policy, refuses a sealed stream, and trims by count, age and a per-topic byte total.
The Redis case now declares that it hosts standalone streams, bounds their bytes and trims an open stream by age, so the conformance suite holds it to all three.
A completion the server rejects leaves the batch it staged under a token no marker will name. When History already says the task failed, the Worker aborts that stage at eviction instead of leaving it for a reader to find. A stage History committed stays pending for the client's repair.
A batch staged for a rejected task leaves the shared log once the Worker evicts the run, so an outside reader never meets it.
Core is pinned at the union head, with the bridge protos regenerated from it in this commit so the merge builds. The changelog keeps each chain's entries once.
Without TEMPORAL_ADDRESS the cases reached the environment's own server, which has no stream service, and failed on the wire.
The union carries the native provider, so the cases gated on it no longer skip.
The Redis producer notifies the stream's channel itself, so one live case needs no notify by hand. The native case opens its own client with the native provider on it.
Only the union carries the memory, Workflow Streams, native and Redis providers together, so the demo's provider switch and its notes live here.
A getattr default is evaluated whether or not the attribute is there, and protobuf 6 descriptors have is_repeated but no label, so the generator needed a protobuf 3 interpreter.
2 of 3 tasks
Owner
Author
|
Paused per the 2026-10-05 design review (the decision record is on the proposal page, section 14). Resumes in phase 2 with the CHASM notification work. 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 merges the Option 7 chain, Max's External Workflow Streams with the Redis provider, and the native chain into one SDK that carries every stream provider.
What changed?
The branch starts at
main. A--no-ffmerge brings in the native chain's top, #87: the channel, the stream interface, the memory, Workflow Streams and Nexus providers, the Nexus consumer and native streams. A second merge brings in the Option 7 chain's top, #83: Max's External Workflow Streams, the fixes on top of it, the channel as its wake and the Redis provider. Both chains share the channel and the interface below the fork, so the merge keeps one copy of each by construction.Core moves to the union Core in moetemp/sdk-rust#36, which speaks both prototypes, and the bridge protos are regenerated from it. That pin is the merge's own resolution, since neither side's pin fits the merged tree. The changelog auto-merged into the wrong place with no conflict (the Redis clause of the interface entry landed under a native entry), so it was rebuilt by hand inside the merge.
A third merge brings in #18, the time-skipping unlock fix, since the union also runs on the time-skipping server. #18 stays its own PR.
Then come the commits where the two chains meet. The native conformance setup skips without
TEMPORAL_ADDRESS.PINNED_LAYER_HOSTS_NATIVE_STREAMSis on, so the cases that need the native provider run here. The Nexus consumer gets a Redis case, where the producer notifies the stream's channel itself, and a native case on its own client. The demo runs on all four providers and its README says how. The visitor generator no longer readslabelon a field, so it runs on protobuf 6.One docstring differs from the earlier union on purpose.
workflow.stream_readerkeeps the wording from #78: a reader in a successor run starts a new subscription, and whether it picks up where the predecessor left off is up to the provider. The earlier union tookmain's copy ofworkflow/_streams.py, which still said nothing crosses continue-as-new. That isn't true for Redis.Part of AI-198 (epic AI-37).
Why?
Each chain is reviewable alone, but only together do they give a worker every provider at once. That's what the examples, the harness and the June scenarios run on. Most of the reconciliation the earlier union carried now happens inside the Option 7 chain, so what's left here is only what needs both chains.
How did you test it?
Link to a test plan if any -
The bridge builds, the lint set is clean, and a regen from the pinned Core changes nothing. On the dev server the fixtures start,
tests/streams, the external stream suite, the replayer and stream client modules and the worker suite pass. Against a local server built from the stream server PRs, with the linked channel kind on,tests/streamspasses withSTREAMS_LIVEset to native, redis and nexus, and so do the external suite and the native stream modules. With the linked kind off, only the cases that need a linked channel fail, as expected. In the worker suite against that server, the unfinished-handler termination cases flake under-n 4on its SQLite store and pass serially. The hand-written files equal the earlier union's apart from the docstring above.