fix(workflow): stop relaxed list reads from stranding canceled DAG steps - #26
Merged
Merged
Conversation
Cancel listed a DAG's steps once and reported success even when that list was incomplete, and a second call on an already-canceled DAG returned early without sweeping. Since the scheduler drops every event for a canceled DAG, a pending step that Cancel missed stayed pending forever, and Await on it blocked until its context expired (nightly run 34570965794, phase 06_FailCancelDelete). - An already-canceled DAG now skips the meta write but still sweeps, so re-running Cancel repairs a step an earlier call missed. Done and failed DAGs still return immediately. - Running out of CAS attempts on the meta returns ErrStaleRevision and leaves the steps alone, instead of canceling steps under a DAG whose status never became canceled. - Running out of CAS attempts on a step returns ErrStaleRevision instead of nil. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013ZSNqKsQpNDG4Z2bjKMsu5
ListSteps, ListSignals and ListDAGs enumerated the bucket with KV ListKeys and then read each key with a direct get, skipping any key whose read failed. Both halves are relaxed reads: the server places the ListKeys consumer on a random replica, and a direct get is answered by any replica, so a replica that had not applied a write yet made the record vanish from the list with a nil error. Callers that act on "all steps" (Cancel, DeleteDAG, Pause, the scheduler) then skip what they cannot see. All three now go through listLive: the key set comes from the stream leader (STREAM.INFO with a subject filter), values come from direct gets, and a listed key the direct get cannot find is re-read through the leader-only STREAM.MSG.GET, which also tells a delete tombstone apart from a lagging replica. Read errors are returned instead of dropped. The raw API requests replace two jetstream helpers that cannot serve this purpose: GetLastMsgForSubject switches to a direct get whenever the stream allows one, and Stream.Info rewrites a cached field without a lock, which races under concurrent list calls. ListSteps and ListSignals filter per DAG. ListDAGs still scans the whole bucket for keys ending in .meta, because DAG IDs may contain dots and a subject filter cannot match a suffix. Verified on a 3-node cluster with e2e/store_lists_lag_test.go: before, 70 of 1836 lists came back incomplete with a nil error; after, 0 across four runs. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013ZSNqKsQpNDG4Z2bjKMsu5
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.
What
Two fixes for the nightly E2E failure in run 34570965794 (phase
06_FailCancelDelete):dep await: got context deadline exceeded, want ErrStepCanceled.Why it failed
NatsStorelisted records with two relaxed reads: KVListKeys, whose consumer the server places on a random replica, and then a direct get per key, answered by any replica, with read errors skipped viacontinue. A replica that had not applied a write yet therefore made a record vanish from the list, with a nil error.Cancelacted on such a list once, reported success, and left a pending step uncanceled. Nothing ever repaired it: the scheduler drops every event for a canceled DAG, and a secondCancelreturned early on an already-canceled DAG.Awaiton that step then blocked until its context expired, which is what the nightly saw.The same relaxed read backed
ListSignals, which gates signal-waiting steps, andListDAGs, which drives the scheduler's sweep.Changes
workflow/cancel.go: an already-canceled DAG still gets its step sweep, so re-runningCancelrepairs a step an earlier call missed. Lost CAS races returnErrStaleRevisioninstead of nil.workflow/store_nats.go:ListSteps,ListSignalsandListDAGstake their key set from the stream leader and re-read a missed key through the leader, which also tells a delete tombstone apart from a lagging replica. Read errors are returned instead of dropped.Verification
e2e/store_lists_lag_test.go, 1836 lists per run on a 3-node clustermake testwith-race,make vet,make linte2e/store_lists_lag_test.gois new and taggede2e, so the nightly runs it. Each iteration creates a DAG through quorum-committed writes and lists it through clients pinned to each node, with background write load to widen follower lag and a follower restart at the halfway point.Notes and trade-offs
ListDAGsstill scans every key in the bucket, because DAG IDs may contain dots and a NATS subject filter cannot match a suffix. That is the same order of work as theListKeysscan it replaces, but the stream leader now answers it.meta(key<dag>.step.meta) still shows up inListDAGsas a bogus DAG with an empty ID. Confirmed on both old and new code.runStoreContractstill runs only againstMemStore. It passes 16/16 againstNatsStorewhen wired up locally.🤖 Generated with Claude Code
https://claude.ai/code/session_013ZSNqKsQpNDG4Z2bjKMsu5