Conversation
Advances the vendored Core from upstream's pin 8cf682b7 to 6e90e6d5, the
head both task worktrees are based on, in preparation for the External
Workflow Streams work.
Core 6e90e6d5 added a `format` field to `Logger::Console`, so the bridge
no longer compiled. Set it to `None`, which preserves existing behavior:
the format then resolves from the pretty-logs env var exactly as before.
Also adds `redis>=5,<9` to the dev dependency group for the Redis Streams
backend that the streaming feature will use in tests.
Verified: cargo check, maturin develop, and the bridge/converter/common
unit tests all pass.
Note: building this tree requires libprotobuf-dev, which supplies
google/protobuf/{duration,timestamp}.proto for prost-wkt-types.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The vendored Core submodule tracked temporalio/sdk-rust with its only remote named `origin`, which is the name the global pre-push hook reserves for the fork. Repoint it at mfateev/sdk-core on task/python-sdk-streaming so Core work for this task has somewhere to land. The submodule commit moves 6e90e6d5 -> 67a15f15. That range touches arch_docs/streaming-poc-docs/ only; no Core source changes, so the vendored tree is byte-identical and the build is unaffected. Recording `branch` in .gitmodules lets `git submodule update --remote` follow the task branch.
Stage 1's remainder (X2) and stage 2's Python half (P1), plus the Python protos regenerated for the Core-side C1 commit this bumps the submodule to. X2 -- a Redis fixture giving each test its own key prefix over the one shared server, so the suite stays parallel-safe under pytest-xdist. The prefix embeds the xdist worker id; teardown scans and deletes under it whether or not the test passed. P1 -- the record model. Offset is an opaque, serializable token that deliberately does *not* implement `<`: offsets are totally ordered, but by their provider's rule, and lexical comparison silently misorders Redis ids the moment the millisecond component changes width. Cursor is a position boundary (BEGINNING | AFTER(offset)), never the identity of a record, so a consumer parked at the tail never has to name an id nobody has written. Data and control records share one offset sequence; only the yield differs. The package exports nothing yet -- the public API lands with Milestone 1. Protos were regenerated with the pinned toolchain CI uses (Python 3.10, protobuf<4, grpcio-tools 1.48.x) so the generated code stays readable under the protobuf 3.x test mode. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…taxonomy Stage 3's Python half: P2, P4, P5, P18. P2 -- the StreamBackend contract, and the conformance suite that is its real deliverable. Every check names the obligation it encodes, and the suite's own tests run a single check against a stub broken in exactly one way: an exclusive range read, a lexical offset comparator, an inclusive watch, idempotency on the key alone, a range read that fills a gap from beyond the range, and immutability undeclared or denied. A suite that only ever passes is not evidence of anything, so each of those must fail, and fail with the message naming the obligation. P5 -- the annotation codec. Core accumulates by byte concatenation and never parses, so an annotation is a sequence of self-delimiting frames and the concatenation of a Workflow Task's deltas *is* the annotation, with no reassembly step that could disagree with Core's. Both run endpoints are recorded rather than a start plus a count, because backend offsets are ordered but not dense. A golden-file test pins the exact bytes; a size test asserts encoded *bytes* stay flat as a single-stream batch grows three orders of magnitude, which is what catches a per-record field being added -- a run-count assertion would not. P4 -- payload encoding through the Workflow's DataConverter, keeping the Payload envelope so a consumer decodes what the producer wrote rather than being told out of band. P18 -- StreamIntegrityError and StreamDecodeError with separate counters, and the mechanical classification rule in one place: if the range validated, the bytes are the bytes that were written, so anything later is a decode failure. They stay plain exceptions rather than FailureErrors -- integrity loss blocks a Workflow, and a blocked Workflow can be resumed after repair where a failed one cannot. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…r, bridge Stage 4's Python half plus P3b: P2b, P3, P3b, P6a, P6, P17, P7. P2b -- park intents keyed (stream key, wait_id), never by stream alone; a stream-keyed intent looks correct until a Workflow subscribes to one stream twice, after which only one of the two can ever be woken. Claims are leased, because a producer that crashes between claiming a generation and signaling would otherwise strand it: every other producer sees it claimed and concludes the wake is handled. A provider that cannot lease declares so and grants every claim -- duplicate Signals are harmless, lost ones are not. The suite proves it fails a stream-keyed stub and a never-expiring one. P3/P3b -- the Redis provider, passing both conformance suites *unmodified*. If a check had to be relaxed for Redis it would be encoding this provider's behaviour rather than the contract. XRANGE serves the inclusive replay read and XREAD the exclusive watch; they are different commands, not one with a flag, and a test asserts they disagree exactly as the contract says. Append is a Lua script so a crash between XADD and the idempotency write cannot leave a record no retry recognises as its own. P17 -- the named-backend registry, and the single enforcement point for the immutability precondition. "Forgot" and "cannot" get the same rejection because both break the same thing. The sandbox refuses a direct provider import, since a provider reached from inside would be a second, unregistered instance with its own connection and no watcher owning it. P6a/P6 -- the producer side. Every binding input is explicit because none can be inferred: activity.Info has no first execution Run ID, so the Workflow passes the chain key in and the producer verifies it by describing before its first append. An Activity derives its session id from its own identity and *not* its attempt, which is what makes a retry re-append rather than duplicate; a plain process must supply one, since a random default looks correct until the first retry. P7 -- the readiness call and the read-only status probe, synchronous and acknowledged. Synchronous rather than a Python awaitable on purpose: a future from future_into_py is bound to the loop that created it, and there is one watcher per subscription with no shared loop. The blocking form is documented as unsafe to call from the loop that owns the Worker -- that loop drives the poll the answer comes from -- and the async wrapper runs it off-loop. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Stage 5's Python half: P8 and P9, plus the Core bump for C7/C14a/C12a. P8 -- the per-Worker subscription manager. activate() is synchronous, runs on a thread-pool executor under a 2-second deadlock timeout, and drives a custom deterministic event loop, so any design in which _apply reads the backend is wrong by construction rather than merely slow. The manager therefore prefetches into a bounded thread-safe buffer and reports readiness to Core only once a record is *buffered* -- readiness for an unbuffered record would produce an activation whose drain must block, which is the hazard the buffer exists to remove. A test drives a provider slower than the deadlock timeout and asserts it delays the report rather than the caller, and another drives the drain from a real second thread against a provider that raises if touched off the manager's loop. Backpressure is the buffer bound: a full buffer stops prefetch, drops nothing, blocks nobody. The three cursors stay apart -- committed advances only on marker commit, delivery on hand-off, prefetch on buffering -- so eviction can discard both speculative ones and restart from committed, which is exactly why "no cursor advances unless the marker commits" is safe to state. Watchers survive Workflow Task completion and are torn down only on cancellation, eviction, or shutdown; `NoOpenWorkflowTask` in particular *keeps* its watcher, since that is the normal window between tasks. P9 -- the Workflow-facing API. Workflow code names a backend rather than holding one, and gets an opaque runtime handle across the sandbox boundary. wait_id comes from a per-Run counter in subscribe() call order, held on the Workflow instance rather than in a module global so an evicted Run takes it with it instead of handing its ids to the next one. Two subscriptions to one stream name are two independent waits, which is the only way broadcast delivery from separate cursors is expressible. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Adds PARK_REASON_TASK_COMPLETED, so a marker written by a Workflow Task that consumed records and produced no commands says which boundary it actually was rather than implying commands that never existed. Bumps the Core submodule to the marker machine and emission primitive (C9, C14b). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
P10a, P10b, P11, P19 -- the integration nothing before this had to touch, and with it a Workflow now consumes records it never read itself. P10b -- the observation delta, built by the per-Run runtime handle because it is the only component that sees deliveries in the order Workflow code received them. Neither the manager (which sees buffering order) nor Core (which is annotation-blind) could reconstruct it. Emission is not conditional on records having been consumed: a first subscription to an empty stream, an activation that drained nothing, and the boundary an activation returned on all produce deltas. Control records are recorded although never yielded -- they occupy offsets inside a run, so a run's count includes them. P10a -- quiescence detected after the drain, which is a registry check rather than an event-loop change since `_run_once` already drains `_ready` to empty. Differing idle timeouts reduce by `min` in wait_id order over the blocked set and nothing else, so the result reproduces on replay. Retention is asked for only when nothing server-bound rides along. P11 -- the `_apply` branch resolves *every* waiting subscription, not only those the job's hints named: the hints are hints, and resolving only them would strand a record whose notification was coalesced away. P19 -- the runtime-only jobs are partitioned in `_handle_activation`, before the executor hand-off. Finalization is answered from in-memory state with no provider call at all, asserted against a provider that raises on every method; a park handshake slower than the deadlock timeout is still answered, and its intents are all installed before anything is rechecked, which is what closes the append/park race. Replay preparation is routed here too but belongs to P13. Three real bugs the tests caught, each of which produced a silent failure: - The annotation header was not riding the first delta. Core accumulates by byte append and never parses, so a header the runtime held but did not emit simply never reached the marker -- and the result could not be decoded at all. - The header's start cursor was captured lazily, letting a record delivered before the first emission slip in front of it; replay of that marker would never have delivered it. - `register` creates the watcher task, and it runs on the Workflow executor thread. `create_task` is not thread-safe, so the watcher was silently never scheduled -- indistinguishable from a stream that never delivers. Also: the readiness future was registered in one map and resolved from another, so nothing ever woke; and watchers outlived Worker shutdown, which matters because an idle cached Run receives no eviction activation at shutdown at all. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Replay reads the exact ranges the marker recorded and delivers them from memory in the recorded order. It never asks the backend what comes next -- the answer is already in the annotation, and consulting current stream timing would reproduce something other than what happened. The four range checks are the whole of validation, and they are sufficient because every provider guarantees a record's bytes cannot change (ADR-003). Given that, the only damage replay has to detect is a record that is no longer there, and a first, middle, or last deletion each fails a different one of the four -- which is what makes the set complete rather than merely plausible. Integrity loss is always reported, never repaired: substituting a later record for a deleted one would hand Workflow code a different history than the one its commands were derived from, and the divergence would surface much later as an unrelated nondeterminism error. Three failure kinds stay distinguishable by type. A backend that is down is a transient StreamStorageError, not integrity loss -- an operator sent to repair a backend that was merely unreachable would find nothing wrong with it. A record that validated but will not decode is a consumer-side StreamDecodeError, since the range check already proved the bytes are the bytes that were written. Delivery walks segments with one event-loop drain each rather than handing over everything at once. Collapsing them would reproduce the record order while changing how many drains occurred, and wait_condition predicates would then fire a different number of times than they did live (ADR-018). This is safe with respect to Workflow time: every segment of a marker belongs to one Workflow Task, so workflow.now() is constant across them either way. Preparation and delivery stay two different moments. The ranges are read and validated before the activation reaches the Workflow thread, so the delivering pass performs no I/O at all -- asserted against a provider that refuses every call, since an inclusive read over a whole recorded batch gets slower the more records the marker committed and would otherwise fail a healthy backend under the 2-second deadlock timeout. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The server-visible wakeup, used whenever no open Workflow Task can accept local readiness. Two of its properties are why it does not reuse the public Signal API, and both are asserted rather than assumed. Codec bypass. The envelope goes out through a raw SignalWorkflowExecution built with the protocol's own serialization rather than the user's DataConverter (ADR-025). Core is the component that must read this Signal and has no access to a user codec, so a codec that encrypted payloads would make the envelope unreadable to the only reader that matters. Nothing user-owned is in it, so bypassing the codec leaks nothing. Stable request ID. The Temporal request_id is derived deterministically from the wake's identity, so a producer retrying after an ambiguous failure sends the identical request and the server deduplicates it -- the public path's fresh UUID per attempt would turn every retry into a second wake. A parked wake is identified by its generation, since a generation is woken once. An unparked wake has no generation to identify it, so the derivation additionally includes the sender's identity and a per-sender monotonic counter, both held fixed across retries of that one attempt: without them two Workers shutting down at different times would derive the same request ID and the server would deduplicate the second wake away, turning a correct retry mechanism into silent loss. The derivation is length-prefixed so a stream named with a separator cannot collide with a different tuple, and the expected value is pinned -- a derivation that hashed anything process-local would pass a same-process equality test and fail in production. The send sequence claims each parked generation under a renewable lease before signaling. Losing the claim means another producer holds it and will send, so this one stays silent; an expired claim is taken over, because a producer that dies between claiming and signaling would otherwise strand the generation with every other producer concluding the wake was already handled. A subscription with no installed intent gets an unparked wake rather than nothing (ADR-023): cached-with-no-open-task and evicted are invisible from out here, and silence in either case loses the record until something else happens to wake the Workflow. A failed wake reports WakeNotAcknowledgedError carrying the wakes still owed, rather than reporting success. The append has already succeeded and must not be retried; the wake must be. Retrying takes the requests verbatim rather than recomputing them, since recomputing would draw a fresh wake counter for an unparked wake and defeat the deduplication that makes the retry safe. Adds parked_wait_ids to the backend contract, with a conformance case and a provider that fails it. The other five parking operations all take a wait_id, which is allocated by a per-Run counter inside subscribe() that no producer can see -- so without enumeration a producer that has just appended has no way to address a wake at all. It is the wait_id half of a key intents already carry, not a new concept. Core's own accept-and-resolve path is already covered by C11's tests. What connects the two halves is a guard asserting the four envelope constants match Core's declarations: they live in different repos, and a drift would leave both test suites passing while the wake silently stopped working. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
publish() now completes only once its wake step is acknowledged. Returning means the record is durable *and* a parked Workflow has been told about it, and that combination is what makes the "durable producer" row of the wakeup-durability boundary true. An append that lands but is never signalled leaves the Workflow parked on data already sitting in the stream, so a publish() that returned success there would be reporting a delivery that never happened. The wake happens after the append and never instead of it: only successfully appended records may trigger wakeup, since a wake for a record that did not land produces a Workflow Task that finds nothing. A failed wake raises WakeNotAcknowledgedError rather than returning, carrying both the offset and the wakes still owed. The two halves fail differently and the caller can only act on one of them -- the append succeeded and must not be retried, the wake did not and must be -- and without the offset the caller has no way to learn which is which. The un-acknowledged state is reachable only by asking for it. `wake=False` says "durable but un-signalled" in the call itself, for a caller appending a batch and waking once; the wake is idempotent either way, so it saves round trips rather than changing the outcome. finish_writing() carries the same contract. A fence is the record most likely to find the Workflow parked -- it is what a producer appends when it has nothing more to say -- so an unsignalled one strands the Workflow for its whole idle timeout at exactly the moment it was waiting to be told. The existing P6 tests now pass wake=False explicitly, which is honest about their scope: they exercise the append half against a producer built without a client. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A stream spans a whole chain, so a new Run must resume where its predecessor stopped -- on replay as well as live. A cursor is never derived from mutable backend state, on any Run: reading the current position from the backend would give replay whatever the stream holds now rather than what the Run started from, and two replays of one history could then diverge. The position travels in a reserved internal header on the Continue-As-New command, persisted in the new Run's WorkflowExecutionStarted and read before the Workflow object exists -- therefore before any subscribe() call, since a cursor restored after a subscription was established would already have been overwritten by BEGINNING (ADR-022). It fills the same annotation-header field a first execution fills with BEGINNING, so replay reads an explicit boundary in every case, including when the stream was empty for the subscription's entire life. Keyed by wait_id, so two same-stream subscriptions restore independently, and carrying each stream's name so a renumbered subscription is caught rather than resumed at another stream's offset -- which the backend would accept and no later check would notice. The header is serialized by this module rather than the user's DataConverter: one Run writes it and the next reads it, and a chain whose Runs were deployed with different converter configuration would otherwise restart at an unreadable cursor, silently, since an unreadable cursor looks exactly like no cursor at all. The encoding is order-stable because it rides on a command replay compares against History. Adds a consumption cursor, distinct from the delivery cursor. Delivery advances by whole drained batches, because that is what the annotation records and what replay must reproduce; consumption advances per record actually handed to Workflow code. A Workflow that stops iterating part-way through a batch has consumed only its prefix, and the rest dies with the Run's buffer -- so a continuation taken from the delivery cursor would step over records nothing had ever shown to Workflow code. Control records advance consumption too: they are never yielded, but they are finished with, and leaving one behind would redeliver it on every Run of the chain. Three bugs the live chain test exposed, none reachable from unit tests: The Worker looked the per-Run stream runtime up through `workflow.instance`, which under the sandbox is a proxy exposing only the WorkflowInstance protocol. It silently found nothing, so the jobs that must never reach activate() -- park, finalize, replay -- quietly went to _apply instead and failed the task as an unrecognized job. The runtime is now held on the Worker, keyed by Run, and dropped on the same RemoveFromCache path as the rest of the Run's stream state. The manager's wake sender was never wired, so a Run that could not take local readiness left the record buffered and waited for a Workflow Task nothing would create. It now sends the reserved Signal through P14's path, reading the park generation at send time rather than remembering one that may since have been abandoned, and counting each owed wake separately so two records arriving in two windows do not deduplicate into one. add_terminal returned the terminal alone and left the annotation open. A Workflow Task that observed nothing creates its accumulator right there, and the header rides that creation -- so the returned annotation began at a terminal frame and could not be decoded. It now flushes anything pending with the terminal and closes the annotation, so a later activation on the same Run opens a fresh one instead of appending a segment past the terminal. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The idle timeout is a Workflow-Task policy, not a per-subscription one. It applies to the complete set the Workflow is blocked on, so one idle stream must not park the task while another is still delivering -- parking on the idle one alone would strand records already sitting in the active one, and the Workflow would wait out a timeout for a stream that was never quiet. merge() waits on several subscriptions at once. They are already one wait set -- every subscription a Run holds is -- so it adds no coordination; what it adds is the ability to block on all of them together. Iterating them one at a time would block on the first while records piled up on the second, and the idle timer covering the set would fire against a Workflow that was not idle. It yields (subscription, value) rather than bare values, because the streams may carry different types and a merged value whose origin had to be guessed from its shape would be unusable for anything but logging. Draining goes in wait_id order on every pass. Records that arrived in one batch across two streams have no inherent order between them, so an order that depended on argument order, on dict iteration, or on which watcher happened to run first would replay differently than it ran. Every wait is marked blocked when nothing is ready, because the quiescent snapshot must name the complete set: a set missing one member would let Core park a wait the Workflow was still waiting on. The recorded delivery schedule follows from the same ordering. A run is a maximal consecutive stretch from a single wait, so alternating records across two streams encode as one run per delivery -- collapsing them into one run per stream would record an order that never happened, and replay would hand Workflow code all of one stream before any of the other. Consecutive records from one stream stay a single run, which is what keeps annotation cost per batch rather than per item. Tests cover the deliverable's four criteria: one idle stream cannot park the task while another is active; a fence on one stream alone leaves the set unparkable while all-fenced streams make it parkable, and a later record reopens a fenced stream rather than violating it; two same-stream subscriptions install distinct park intents, verified by reading both back from the backend, and each receives every record rather than competing for them; and an alternating two-stream batch encodes one run per delivery. Factors the subscription iterator into a fill/take pair so merge and single- subscription iteration share one drain, rather than having two paths that could disagree about when delivery and consumption are recorded. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Two obligations Core cannot discharge. Teardown ordering was already in place: per-Run teardown is driven by RemoveFromCache and nothing else, which is what guarantees a FinalizeExternalStreams in flight is answered before the manager's state for the Run disappears. A shutdown hook that tore Runs down itself would have no such ordering. The sweep is new. An idle cached Run gets no eviction activation at shutdown at all, and that is exactly the Run that most needs one: its records are buffered in a process about to exit, and nothing else will ever tell the Workflow they arrived. For every Run still holding subscriptions the sweep asks Core's read-only status probe -- deliberately not the readiness call, which asserts a buffered record and would manufacture a spurious Workflow Task on the way out for a Run that had nothing waiting. The four answers differ in what is owed. WftOpen: nothing, because Core owns that transition and a wake would race it into a second task for a Run already being attended to. Parked: nothing, because a producer's next append wakes it through the ordinary path. NoOpenWorkflowTask and RunNotFound: an unparked wake per subscription, awaited rather than fired, since an unacknowledged wake is the case this sweep exists to prevent. Failure is reported, never dropped. A wake that cannot be acknowledged increments external_stream_shutdown_wake_failed rather than only logging: a dropped wake is silent by nature -- the Workflow simply waits, and nothing distinguishes that from a producer having nothing to say. A wake that failed is never counted as delivered. Shutdown is never blocked past a grace period. A Worker that could not reach the server would otherwise hang on the way out, turning a recoverable delivery problem into an unrecoverable process one; a probe failure likewise logs and moves on rather than stopping teardown. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Twelve cases, all covered. Most fall out of the deliverables that produced them (P21, P15, P10b); these six did not, which is the reason a milestone gate is a separate list rather than the sum of the deliverables' own "done when"s. Readiness on one of several streams resets global quiescence: the wait that received a record re-enters the blocked state with a new generation while the others keep theirs, which is what makes Core restart the idle timer covering the set rather than let one already counting down park a Workflow that had just been given work. An alternating two-stream batch is the worst case for annotation size -- every delivery starts a new run, so the annotation grows per record rather than per batch, and it is the one shape where the byte budget is reachable in normal use. It must ask for rollover rather than grow past the limit, because an oversized marker cannot be written at all. Simultaneously ready streams are buffered together, so a batch that arrived at once costs one activation rather than one per stream. Two same-stream subscriptions stay independent through park, cancel, and Continue-As-New, verified by reading both intents back from the backend: a collapsed intent looks like success from the caller's side, since the second install simply returns. The idle-timeout reduction is a pure function of the quiescent set's configured values, so registration order cannot change it. Anything live in the inputs would produce a different WorkflowStreamQuiescent on replay and fail as nondeterminism. A wake Signal names one stream but resolves every blocked wait: by the time the Workflow Task exists other streams may have records too, and resolving only the named one would leave those buffered with the wake they needed already spent. Resolving twice is harmless, because two producers racing one generation both signal. Also fixes a watcher leak found while writing these. Re-registering a wait under a key that already had one left the previous subscription unreachable with its watcher still polling the backend for the life of the process. Nothing else would have noticed: the new subscription works and records are delivered, and the only symptom is a connection that never closes. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…ssing retry Three gaps found by auditing the Milestone 1 required-test list against what the suite actually asserts. An annotation naming a wait the Workflow did not create was raised as StreamIntegrityError, but this module's own taxonomy table puts that condition in row four -- ordinary nondeterminism, whose operator response is "fix or version the Workflow code". Integrity loss says "repair or restore the backend", so the error was sending an operator to repair a backend that is fine: the recorded ranges are exactly where they were written and it is the Workflow code that moved. This is the same mistake as calling an outage integrity loss, in the opposite direction. Both sites now raise NondeterminismError with the remedy in the message -- the replay path, and the Continue-As-New cursor restoring onto a renumbered wait. The tests asserted the wrong type and now assert the right one, including that the error is *not* reachable through the storage taxonomy. P20 documented an in-grace retry it did not have. The plan requires an unacknowledged shutdown wake to be retried within the grace period under the same request ID before being reported; the implementation reported on the first failure. A Worker shutting down is often shutting down because something is unhealthy, which makes the first attempt the one most likely to land in the middle of it. Retrying is safe and is not a second wake: the request ID is derived from the wake's identity, so the server deduplicates it against an attempt that may in fact have arrived. Bounded at three attempts, because the alternative to giving up is a Worker that never exits -- and the metric is what makes giving up visible. Two retention-loss cases were untested against real Redis. `XTRIM MAXLEN` is how a recorded range actually disappears in production -- a stream with a retention policy trims its own head while a Workflow is still parked against it -- and it must read as integrity loss rather than as a transient failure that would be retried forever against data that is not coming back. A deleted write fence is the case most easily dismissed as harmless, since a control record is never yielded to Workflow code; but it occupies an offset inside a run, so the count no longer matches and every later control position shifts by one. Also asserts what eviction is *for*: discarding prefetch state is the mechanism, but the outcome that makes it correct is that a Run which consumed records without committing a marker sees those same records again. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The required-test lists are what "the milestone is met" means, and a gate that is only a claim in a document is not a gate at all: someone renames a test, the case it covered quietly stops existing, and the list still says the milestone is met. So the lists are parsed from the plan itself -- from the copy in the vendored Core checkout, which is the branch the Core-side work is on, rather than from a duplicate that would drift -- and each case is mapped to the test that covers it. Three checks, each catching a different way the gate rots: the parsed case count still matches the count the plan's own heading declares; every case is either covered or recorded as open; and every mapped test actually exists, so a rename fails here rather than leaving the list pointing at nothing. What it deliberately does not do is assert that mapped tests pass -- that is the suite's job. Covered and open are separate maps that must partition the list. A partially covered case recorded as covered is exactly how a gate stops meaning anything, so each open case carries a sentence saying what is still missing, and the gate prints them. Milestone 2 is met: 12 of 12. Milestone 1 stands at 35 of 55, with 20 open -- three of them blocked on C15b, whose ParkReason::Shutdown is never constructed anywhere in Core, so the shutdown-with-a-task-open path cannot be produced at all yet. The remaining seventeen are gaps in the suite rather than in the implementation, and each names what it needs. The gate is not asserted green for Milestone 1. A test failing permanently because a milestone is not finished would say nothing new on any given run; what must not happen is losing track of which cases remain, and recording them is what prevents that. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Vendors the park handshake and replay marker lookahead. No Python change was needed: the park result already went out as an ExternalStreamParkResult completion command rather than through the input lane C8 removed, which is what the protocol specifies. Three required cases close. The marker's byte cost is now measured with sparse control records, which is the realistic shape and the one place a per-record cost could hide -- control positions are the only field carrying per-record information at all. Bounded in both directions: flat as record count grows by three orders of magnitude, and growing only a small constant per control record. Wake deduplication is asserted against a real server rather than by comparing request-ID strings. Two producers retrying one generation must cost one Workflow Task; that two derived strings are equal proves nothing about what the server does with them, and the server is the component that has to collapse them. Conversely two Workers' unparked wakes must both be delivered, while one sender's retry of a single attempt must still collapse -- otherwise the shutdown sweep's retry would wake the Workflow twice. Replacing the derivation with a fresh UUID per attempt fails the two deduplication tests and leaves the delivery one passing, which is the discrimination they exist for. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Everything replay-related so far drove the driver directly, which proves the
mechanism but not that a history a Worker actually wrote can be fed back through
it. It cannot, yet -- and running it through the tool users reach for found two
bugs the unit tests could not.
Lang re-derived the annotation while replaying. Every marker for a replayed
Workflow Task is already in History, so a WorkflowStreamProgress emitted there
asks Core to write a second marker for observations already recorded, and the
command is matched against the very event it was read from. Worse, what was
accumulated during a replayed activation rode the *next* completion, which is
live -- so the duplicate marker appeared even for a stream that never delivered
anything, since a registration alone is replay-visible. Lang now emits no
progress command at all while replaying and drops what it accumulated.
Delivery during replay no longer accumulates runs. Advancing the cursors is
right: the runtime has to end up in the state the live run was in. Recording
them into a new annotation is not.
Replayer gained external_stream_backends. Replay re-reads the recorded ranges
from the provider, so a history containing stream markers could not be replayed
at all -- there was no way to supply one. Asserted both ways: with the backends
the replay runs, and without them it fails naming the missing option rather than
surfacing as an attribute error.
Two of the three tests are xfail(strict) on a Core defect this exposed: C10's
replay lookahead does not claim an external stream marker that is followed by
one of lang's own commands in the same Workflow Task, so the marker event
reaches the next machine instead ("Timer machine does not handle this event").
Real histories have exactly that shape -- C14b emits the marker before lang's
commands -- and the existing Core unit histories do not. Strict, so the tests
turn red the moment Core is fixed and the markers can be removed.
The wake deduplication assertions now wait for the signal count to settle rather
than reading once. Counting immediately races persistence in both directions: too
early and a wake that was recorded is missed, too eager and a duplicate about to
appear is not seen. It flaked once in a full-suite run for that reason.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A Workflow that consumed a record, started a timer, and came back for the next one hung forever with that record sitting in its buffer. The readiness for it had already been reported and consumed. A record arriving while Workflow code is doing something else -- a timer, an activity, another stream -- is buffered and announced immediately, and by the time the Workflow asks for the next value no further notification is coming: the watcher has moved its prefetch cursor past it. The iterator filled its ready list once per batch and then blocked, so it waited on a notification that had already been spent. It now refills before every record rather than once per batch, and takes one more look after registering its wait. Filling is a buffer pop and costs nothing, and the look-after-registering is what closes the window on the other side: anything buffered earlier is found there, anything buffered later resolves the future. merge() does the same for the same reason. This is the failure mode the wake Signal exists to prevent, reached from the other direction -- the Signal arrived, Core made a Workflow Task, and lang declined to look. Only an end-to-end test could find it: every unit test either buffers before iterating or resolves the future by hand, and both paths worked. Case 22 of the Milestone 1 list closes with it: an append after a completion carrying server-bound commands wakes the subscription, asserted by finding the reserved wake Signal in the Workflow's own History rather than by observing that the record eventually arrived -- which timing alone could explain. The integration helper now continues its producer sequence across calls. Restarting it re-used a `(session_id, sequence)` idempotency key with different content, which the backend contract rejects outright -- correctly, and the test was wrong. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
… bridge Vendors the Core-owned shutdown and eviction transitions. The empty-stream replay test now passes, and its xfail comes off. The Core defect I recorded against it does not exist: `maturin develop` had been run from `temporalio/bridge/`, where it builds the bridge crate's own package into the venv, while the SDK imports `temporalio.bridge.temporal_sdk_bridge` from the package directory. That file was hours stale, so every Python run had been exercising a Core from before C8, C10, and C15b. The nondeterminism error I attributed to C10's lookahead was simply a Core that had no lookahead in it. Building from the repo root, where [tool.maturin] sets the module name, fixes it. Worth stating plainly: a stale native extension fails by being *quietly correct for the old code*, so the Rust tests stayed green and only the end-to-end tests disagreed. The lesson is to check what the process actually loaded, which is now a habit rather than a hope. One real defect remains, marked xfail(strict) with an accurate reason: a `wait_condition` evaluated alongside a subscription iteration livelocks -- the Workflow never completes and the server accumulates tens of thousands of events. That is an ordinary pattern, so it is a defect rather than a limitation, and it is kept as a marked test rather than deleted so it cannot be lost. 411 tests pass against the correctly built bridge. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Investigating the reported livelock in order to document it showed there is nothing to document. The `publish` helper in this module restarted its producer sequence on every call, so the second call re-used `producer/0` with different content and the backend rejected it -- correctly; an append is idempotent on identity, not on the key alone, and that rejection is a required conformance case. The second record was therefore never appended, and the Workflow waited for data that was never coming. What looked like a livelock between retention and quiescence was a Workflow blocked on a record that did not exist, with the exception surfacing in the test's own producer several lines away. Two test modules had grown the same helper independently and both were wrong the same way; the other was fixed earlier for the same reason. Confirmed by discrimination rather than by the fix appearing to work: the same Workflow shape without a `wait_condition`, and with an inert one, both complete normally, so the concurrent condition -- the thing the report blamed -- was never involved. With the helper fixed the original test passes in under three seconds, where it had previously exhausted a sixty-second timeout while the server accumulated tens of thousands of events. The xfail comes off. The suite is 412 passing with nothing marked, and it now runs in forty seconds rather than ninety-seven, because nothing hangs. Both hazards behind this and the earlier stale-bridge report are written up in arch_docs/streaming-poc-docs/verification-hazards.md. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A producer that never let the subscription buffer run dry made a single activation never end. The Workflow thread drained while the watcher refilled concurrently from the Worker's loop, so the iterator never blocked. Measured against the real iterator: 316,086 records consumed in two seconds, blocking zero times -- then the 2-second deadlock timeout failed the Workflow Task, and every retry did the same. A healthy producer, a healthy backend, and a Workflow that never progressed. `max_records` had been plumbed through every layer since P8 and no caller ever passed it. Delivery is now bounded by MAX_RECORDS_PER_ACTIVATION, a record count and never a duration: how many records fit in a time slice depends on machine speed, record size, and load, and because segment boundaries are recorded in the annotation, replay would divide the same records differently and fire wait_condition predicates a different number of times. The counter lives on the runtime rather than per subscription, because it is an activation budget and merge() spends one across several waits. When the budget is exhausted the subscription blocks even though records are still buffered -- that block is what ends the activation -- and readiness is re-armed for what remains. The re-arm is a correctness requirement, not liveness: readiness is announced once, when the watcher buffers a record, and that watcher has already moved its prefetch cursor past these. Without it the Workflow blocks forever on records already in front of it, and blocks *quiescently*, so Core would start the idle timer and eventually park a Workflow Task whose data had already arrived. For the same reason a wait the budget stopped no longer reports itself immediately parkable. The segment that exhausts the budget records BATCH_LIMIT. That segment_end_reason existed with no producer; the other two both assert nothing more was available, which is false when the runtime stopped with records buffered, and recording one of them would put a claim in History that the stream ran dry when it did not. Replay bypasses the budget in both directions: the recorded segments already fix how many records each activation received, and replayed records do not spend the live budget either. Verified beyond the suite: an acceptance harness sharing no helper with these tests, driving the real runtime and the real iterator against a never-empty manager, delivers exactly 256 and blocks. Independently mutation-tested the two claims that matter most -- always-BATCH_LIMIT and never-BATCH_LIMIT each fail a different test, so the reason is covered in both directions, and removing begin_activation() from activate() fails the end-to-end test, so the budget provably applies in production rather than only in unit fakes. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
None of these were visible to the 429-test suite. All four were found by writing the Milestone 1 cases end to end. `with_options(idle_timeout=...)` was silently ignored. `subscribe()` never passed it to `register`, so every subscription got the one-second default and `ExternalStreamSubscription.idle_timeout` was dead. The `min` reduction over the quiescent set worked perfectly and no user could reach it: the test that covered the reduction called `runtime.register(idle_timeout=...)` directly, bypassing the only path a Workflow has. The shutdown wake sweep never ran on a clean shutdown. It was wired into `drain_poll_queue()`, which the Worker substitutes only for a task that already raised -- so P20's entire mechanism was dead on the path users take, leaving watchers, buffers and backend connections behind. It now runs after every activation is answered, which preserves the rule that per-Run teardown is driven by RemoveFromCache, and before the bridge is finalized, because the read-only status probe still needs it. `Subscription.wait_generation` was declared and read but never assigned. Core compares it against the generation lang put in its quiescent snapshot, which advances every time a wait re-blocks, so after the first block readiness was reported against a stale zero and answered Stale -- and since the watcher's prefetch cursor was already past the record, nothing re-reported it. An append after a confirmed park never woke the Workflow at all. A replayed marker's records were delivered a second time, live. Nothing repositioned the manager's subscription after replay, so the watcher's buffer still held records the marker had already committed. The fix commits the boundary and resets the buffer, and discards a read already in flight rather than letting it undo the reposition. The terminal alone was not enough to derive that boundary: a marker written on a completion carrying server-bound commands has no terminal frame at all, so the boundary comes from the last recorded delivery per wait with the terminal overriding where present -- a terminal-only fix would have looked right in unit tests and been a no-op in production. Case 30 of the required list passes as a result and its expected-failure marker comes off. Case 29's sweep half now works and it stays marked, on a fifth defect left for its own change: an unparked wake's request ID derives from the sender's identity and a per-sender counter that restarts at one, so two Workers sharing a Client identity derive byte-identical request IDs and the server deduplicates the second away. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The wake-Signal spec says an unparked wake's request ID includes the sender's identity and a per-sender counter precisely so that two Workers shutting down at different times derive different IDs and both wakes are delivered. The implementation used the *client* identity, and two Workers in one process share a Client while each one's counter restarts at 1 -- so their first unparked wakes derived byte-identical request IDs, the server deduplicated the second, and the Run the surviving Worker had picked up never got a Workflow Task. The mechanism built to prevent silent loss was producing it. The sender identity is now drawn once per manager and held for its lifetime. Both halves matter and both are tested: two Workers must differ, or a wake is lost; and one sender's retry must not, or the shutdown sweep's in-grace retry becomes a second wake rather than a resolution of the first. The value is random, which is safe here because the wake Signal is sent from the Worker's own event loop and never from Workflow code, so it is not replay-visible. The client identity rides along as a prefix so a request ID stays traceable to a client in server-side logs. Parked wakes are untouched and must stay that way: they are identified by their generation, so every producer racing to wake one generation still derives the same ID and the server keeps one. Verified by mutation in both directions. Sharing the identity between Workers fails three tests including the two-Workers case; redrawing it per wake at the call site fails the retry case. An earlier mutation of the generator function itself changed nothing and proved nothing -- it is called once at construction, so the property under test is that the manager holds the value, not what the function returns. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
With the rollover deadline anchored at the Workflow Task's start, rollover demonstrably works: the retained task now completes at 8.0s of a 10s timeout in every run, where before it ran until the server timed it out. The Signal-latency case was a test defect. A Signal held behind a retained task is released by the rollover 16 to 65 milliseconds past the deadline, and the assertion allowed none: it compares two WorkflowTaskStarted times with a completion RPC, the server scheduling a replacement, and a poll in between, so it could never be satisfied. The tolerance is 500ms -- eight to thirty times the measured overhead, and still a quarter of the gap between the rollover deadline and the Workflow Task timeout, so an input released by the server timing the task out instead lands around 10s and still fails. That distinction is the only thing the assertion exists for. Checked for teeth by halving the rollover fraction: a three-second error is not swallowed. The other two remain expected failures, and their stated reasons were wrong -- they described the Core anchoring bug, which is fixed. Both are now blocked by a different defect, stated in present tense with the mechanism: the replacement Workflow Task activates with is_replaying=True even though it is live, so lang correctly emits nothing, Core keeps the previous wait generation, and every later readiness report is answered Stale. The watcher owes a Signal only on the three answers meaning local readiness could not be delivered; Stale is not one of them, so nothing is sent, nothing is re-announced, and the Run goes quiet with records buffered in the Worker. The watcher is right to trust the answer. The answer is wrong. A stale reason is worse than none, since the next reader takes it as the diagnosis and stops looking. Also fixes a latent flaw in the wake assertion: its baseline counted every Signal while the check counted only wake Signals, so the two could agree by accident. Both now go through one helper that names __temporal_external_stream_wake in History. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Queues the stream resolve job before the activation is built, so a live replacement Workflow Task is no longer handed to lang marked as replay. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
With the replay-flag fix in Core, the continuously fed stream runs to completion -- 30 of 30 records in about nine seconds, where it used to stall at 27. The test failed in its own teardown, on an unconditional terminate against an execution that had already finished. Two other tests in the file swallowed that same error with a bare `except Exception`, hiding it; all three now go through a helper that swallows only NOT_FOUND and re-raises anything else. Case 19 also gained the assertion it was missing: that some Workflow Task was actually held to the rollover deadline and closed inside the Workflow Task timeout. Without it the test passes on a Workflow that never rolled over at all, which is precisely the bug it exists to catch. Case 23 was looking in a window that does not exist. Instrumenting every readiness answer shows the post-rollover append answered Accepted, not Stale -- delivered live on an open Workflow Task. The rollover completion sets force_new_wft, and the replacement Workflow Task started 0.8ms later, so there is no interval between them for a client RPC to land in. No wake is owed on Accepted; the append was answered better than a Signal would have answered it. The test now reaches the state it was actually after -- a subscription active and unparked with no open Workflow Task -- by having the Workflow await a real Timer after the rollover. A completion carrying a timer is server-bound and cannot be retained, so it holds that state still for as long as the test needs. The docstring says plainly that the rollover's own window is unaddressable from outside the process and that the Timer stands in for it, so a reader is told where the substitution is rather than discovering it. It asserts position as well as presence: the wake Signal's event id must follow the TimerStarted event, because a Signal owed during an open task reaches History only when that task completes, and counting alone would accept one owed earlier. Both reasons had gone stale twice over -- first describing the rollover anchoring bug, then the replay flag -- and each time the stale text was a plausible diagnosis that would have stopped the next reader looking. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Fourteen cases move to covered as their tests landed, across Core and Python. Every mapped node id resolves to a test that exists -- which is the one thing this gate actually checks, and the reason a rename cannot quietly empty it. Case 14 gains a second test, and the reason is worth recording. It was mapped to a producer-side unit test that asserts a producer sends one Signal for a wakeable generation. That test cannot see Core's answer, which is exactly where the defect was: the manager's wait generation was never assigned, so every readiness report after the first block was answered Stale and an append after a confirmed park never woke the Workflow at all. A case can be mapped to a passing test and still be uncovered, if the test cannot observe the thing the case is about. Two remain open and are recorded as such rather than rounded up. Case 29 is under re-diagnosis. Its stated reason has gone stale three times -- first a Core rollover anchor, then a replay flag, then a wake request-ID collision -- each a plausible diagnosis that was true when written and stopped the next reader looking once it was not. The current note says what is observed and that the cause is not yet established, which is the honest state. Case 36 was never assigned to anyone. Splitting the sixteen open cases between two agents, I left it out of both lists and did not notice until reconciling the map. Its nearest test covers subscribe-and-replay but neither the park nor the eviction the case names, so it is partially covered, which counts as not covered. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The test could never reach the state the case is about. Its Workflow started a timer on the same completion it first blocked on, and a completion carrying a server-bound command asks for no quiescence -- which is the only path that registers the wait set in Core. So Core held an empty set, marked nothing ready for the wakes that did arrive, and completed three Workflow Tasks without ever activating lang. The Run was unresumable by anything at all, including the sweep it was written to exercise, and the sixty-second timeout was the client waiting on a result rather than anything about shutdown. The Workflow now blocks on the stream first, so that completion carries no command and registers the set, and only then starts its timer, so the next completion suppresses retention and leaves the Run in the no-open-task window with the set still registered. The order is load-bearing and the docstring says so. Three assertions are sharpened while the shape is being fixed. The window is verified rather than assumed -- the read-only probe must answer NoOpenWorkflowTask before shutdown. The wake assertion demands a *new* unparked Signal after shutdown, since an earlier wake in the same History proved nothing. And the record is published after the handover, because one buffered before shutdown is delivered by the watcher's own wake rather than by the sweep. It fails in about two seconds at a named mechanism now, instead of timing out after sixty at a result that was never coming. The expected-failure reason is rewritten for the fourth time, and for the fourth time the previous one was wrong: the sender identity is fine, and the two Workers do get distinct identities. What deduplicates the sweep's wake is a stale park intent. Nothing removes a confirmed park's intent when the park resolves -- remove_park_intent runs only on the abort branch -- so the backend still reports a park generation after Core has cleared it, the sweep builds a parked wake instead of the unparked one ADR-023 requires, and a parked wake's request ID ignores sender identity by design. Forcing the generation to zero moves the failure exactly where the diagnosis predicts, which is the evidence for it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Both chains share the channel surface and the stream interface, so the merge keeps one copy of each. Core moves to the unified head that carries both lineages, and the clients are regenerated from it.
The union tests against the time-skipping server too, so it carries the fix that gives concurrent result waiters one shared unlock.
The workflow info and the payload converter factory grew on main after the external chain's helpers were written.
On the shared wire zero says a producer doesn't number its records, and the Redis producer does, so its first record is one.
The shared interface reads and writes DEFAULT_TOPIC when a call names no topic, and the Redis handle stores it under that name like any other.
Without TEMPORAL_ADDRESS the cases reached the environment's own server, which has no stream service, and failed on the wire.
The workflow_streams provider answers reads from the workflow itself, so the read has to happen before the worker stops.
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.
This was referenced Oct 3, 2026
Owner
Author
|
Replaced by #88 in the v4 series: one lineage through the channel and the interface, with Max's Option 7 prototype replayed onto it. 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 native main chain and the external chain into one SDK that carries every stream provider.
What changed?
Execution,ChannelAddress, the channel surface and the interface, so the merge keeps one copy of each. Core moves to the unified head that carries both lineages, and the clients are regenerated from it. That pin is the merge's own resolution, since neither side's pin fits the merged tree.main's worker does.workflow.Infoand the payload converter factory grew onmainafter those helpers were written.TEMPORAL_ADDRESS. The constructor-publish case reads while the worker still polls.PINNED_LAYER_HOSTS_NATIVE_STREAMSis on, so the cases that need the native provider run here.labelon a field, so it runs on protobuf 6.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. The commits after the merges are the places where the two chains meet and neither one could fix alone.
How did you test it?
Link to a test plan if any -
The bridge builds,
cargo clippyandpoe lintare clean, and a regen from the pinned Core changes nothing. The streams suite, the external stream suite and the replayer and stream client modules pass on the dev server. Against a local server built from the stream server PRs, the streams suite passes withSTREAMS_LIVEset to native, redis and nexus, and so does the external suite. That includes every Nexus consumer and stream channel case. With the linked kind switched off, only the cases that need a linked channel fail, as expected. On the stock CLI dev server, the external suite and the Redis lane pass, apart from the standalone activity cases, since that server has no standalone activities.