Conversation
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.
The buffer holds a body, topic and offset rather than the wire message, and the cap moved to `workflow_read_stream_messages`.
The pinned Core carries the stream additions to the public API, so the command, enum, history and workflow service messages were stale. Only the vendored stream service protos still come from the server checkout.
CI runs `ruff check --select I` and `ruff format --check` over the whole tree, and `check-protos` reformats what it regenerates.
A Standalone Activity has no `workflow_id`, so opening its Workflow's stream now says that rather than sending `None` on the wire. `Optional` and `typing.Sequence` are deprecated spellings the linter rejects.
The changelog checkpoint requires an entry for any user-facing change.
This package has to load on a protobuf 3 runtime, and generated code refuses to load on a runtime older than the one it was built against. The protoc of that generation has no `--pyi_out`, so the stubs come from `mypy-protobuf`, which is what the other generated protos already use.
The stream client needs the replay slice lookahead fix that branch finished.
The new server_streams module gives WorkflowStream and WorkflowStreamClient a second definition, so pydoctor can no longer resolve the short names.
That head matches upstream's nexus model WIT for the workflow-id policies, which the current generator requires.
The generated module was missing the local-activity marker argument flag its own stub already declared, so the two disagreed about the message.
The append cursor names the last record and is None for a batch the server deduplicated, so the stream client reports what the server said about the append. Producer identity is resolved by the package, inbound frames carry no topic, and the inbound id comes from the shared helper.
A handle on a workflow's stream resolves the run once when it opens, so a follower is not redirected to a successor's empty stream after continue-as-new. The channel cache is keyed by the loop object and closed through the provider's close hook, and latest() reads a missing stream as empty but lets every other failure through.
A cancel landing inside the flusher's append unwound with the batch it had taken off the buffer. The client now asks it to stop and waits. The handle is pinned to the run it opened on, from the activity's own run id or one describe, and the channel cache is the shared one.
A range is recorded as consumed once and never resent, so a repeated, skipped or mis-sized delivery would hand the workflow duplicate or shifted bodies with nothing to say so. The worker also refuses a publish over the server's per-message and per-batch byte limits before the command exists, for the reason the count limit already gave.
The native activation job and the two native commands moved to the field numbers the combined Core tree uses, so both trees agree.
The activation job is now deliver_stream_records and the command append_stream_records, both carrying StreamRecord, so the payload visitor reaches each record's body and metadata.
…lishes. The delivered job and the publish command carry StreamRecord protos, and a task's publishes on one stream are held until the task completes so they become one command and one History event.
…rors. Every stub call translates the transport's failure into StreamNotFoundError or RPCError, and the client's payload codec is applied to bodies so the two sides of a namespace with a codec agree.
The workflow subscribes to the topic by name and publishes with the batched command; outside code appends and long-polls the same stream, following the run chain by cursor unless pinned.
This layer looks a successor up per topic and returns where to resume, so the activity handle's override, which never has a successor, takes the same arguments.
…-native-replay # Conflicts: # temporalio/streams/providers/native.py # tests/worker/test_workflow_stream_e2e.py
The replayer fetched recorded ranges over a bare channel to the client's host, so a client with TLS or an API key named the right server and could not authenticate. It now opens the shared channel from the client's connection and closes it by the same key. A range refused below the retention floor arrives typed now and is reported as the stream no longer holding it, as before.
…ran. The server re-supplies a replaying workflow's recorded ranges within a budget, and a range it cannot fit has to come from somewhere or the task cannot run. The worker now reads the missing offsets of a short DeliverStreamRecords job through the client's stream channel before the activation runs, with the range readers the replayer already had, moved where both can use them. The server still refuses an over-budget task outright and the pinned Core fails a short slice itself, so this half waits on a server flag and a Core pass-through.
The server now seeds the reset run's inherited stream with the range the re-run task had consumed, at the offsets it held, so that task decides on the same input twice over and the inherited stream's floor holds the seeded record. The reset case follows that contract from the srv chain at 217eed936.
…-native-replay # Conflicts: # temporalio/worker/_workflow.py
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 replays a workflow that read a server-side stream, and follows a workflow reset.
What changed?
Replayer(stream_client=)fetches every range the completed tasks recorded from the stream service and hands the records to Core with the history. A range the stream no longer holds fails that replay withStreamNotFoundError. Without a client, a consuming workflow's history is refused with both remedies in the message.HistoryPusher.push_historytakes the slices as serializedStreamSlicemessages.WorkflowHistory.stream_slicescarries the records.to_jsonandfrom_jsonwrite them as astreamSliceslist beside the events, andReplayer.fetch_stream_slices(client, history)fills them while the stream is retained. The exported file is then the whole replay input.temporalio/worker/_stream_ranges.pyholds the range readers andfill_short_stream_ranges, the worker's half of paged re-supply. Nothing triggers it yet, because the server and Core don't pass a short slice through.CHANGELOG.mdgets the entry for the native feature.This is the tip of the native chain. It ends with an empty merge of the replaced #1's head, so #6 and every pin naming that head stay valid.
Part of AI-198 (epic AI-37).
Why?
History records the offsets each task consumed, never the records, so a workflow that read a stream can't be replayed without this layer. A support case arrives as a history file, so the export has to carry its records. The sharp edge is a sticky task to a worker that evicted the run, which gets no records. The pinned Core fails that task before the workflow runs, so a legacy query goes back to the server and is retried non-sticky.
How did you test it?
The lint set from #8 ran on a cold
.mypy_cache.test_replayer.pyandtest_stream_resupply.pyrun without a server, andtest_workflow_stream_e2e.pyruns on a server from moedash/temporal#17. It covers the replayer with a client, tampered ranges, a deleted stream and no client, a query on a cold worker, a sticky query to an evicted run, a reset and a chain of two, a pinned read of a reset run fromBEGINNING, a codec on both halves, and an exported history with no client. The reset case needs a server that seeds the reset run's stream.