Skip to content

fix(workflow): stop relaxed list reads from stranding canceled DAG steps - #26

Merged
F1bonacc1 merged 2 commits into
mainfrom
fix/leader-authoritative-store-lists
Sep 11, 2026
Merged

F1bonacc1 merged 2 commits into
mainfrom
fix/leader-authoritative-store-lists

Conversation

@F1bonacc1

Copy link
Copy Markdown
Owner

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

NatsStore listed records with two relaxed reads: KV ListKeys, whose consumer the server places on a random replica, and then a direct get per key, answered by any replica, with read errors skipped via continue. A replica that had not applied a write yet therefore made a record vanish from the list, with a nil error.

Cancel acted 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 second Cancel returned early on an already-canceled DAG. Await on 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, and ListDAGs, which drives the scheduler's sweep.

Changes

  • workflow/cancel.go: an already-canceled DAG still gets its step sweep, so re-running Cancel repairs a step an earlier call missed. Lost CAS races return ErrStaleRevision instead of nil.
  • workflow/store_nats.go: ListSteps, ListSignals and ListDAGs take 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

Check Before After
Unit tests for a missed read (steps, signals, DAGs) fail: record dropped, nil error pass
e2e/store_lists_lag_test.go, 1836 lists per run on a 3-node cluster 70 incomplete with a nil error (56 step lists, 14 signal lists) 0 incomplete across 4 runs
make test with -race, make vet, make lint n/a clean

e2e/store_lists_lag_test.go is new and tagged e2e, 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

  • ListDAGs still 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 the ListKeys scan it replaces, but the stream leader now answers it.
  • These lists now return errors in windows where they used to return a partial result. Event handling redelivers on error, and the sweep runs again on its next interval.
  • Pre-existing, not addressed here: a step named meta (key <dag>.step.meta) still shows up in ListDAGs as a bogus DAG with an empty ID. Confirmed on both old and new code.
  • runStoreContract still runs only against MemStore. It passes 16/16 against NatsStore when wired up locally.

🤖 Generated with Claude Code

https://claude.ai/code/session_013ZSNqKsQpNDG4Z2bjKMsu5

F1bonacc1 and others added 2 commits September 12, 2026 01:40
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
@F1bonacc1
F1bonacc1 merged commit edd239f into main Sep 11, 2026
8 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant