◐ Off-By-One · answer catalog

go-work-stealing-scheduler-chaselev-victim-locality

2 answer(s)gogo1.26gogo1.26

This problem is the composition of four independent hazards. Getting three right still hangs.

📦 Source in repository (JSON)

Answer 1

I diagnosed the problem class, built a complete implementation, and verified it under the race detector. All source and the full write-up are in ~/wsched/ (document at ~/wsched/SOLUTION.md).

Work-Stealing Fork/Join Scheduler — Verified Solution

Root-cause analysis: why naive implementations hang and corrupt

This problem is the composition of four independent hazards. Getting three right still hangs.

1. Lost wakeups (check-then-sleep race). The naive if noWork() { park() } misses any event landing between the check and the sleep. The common "optimization" of an atomic nparked flag with if nparked.Load()!=0 { broadcast() } is still broken: the producer can publish work, read nparked==0, and skip the broadcast before the worker announces itself, while the worker's check happened before the publish — the worker sleeps forever.

2. Join deadlock (blocking the worker). If Wait blocks the OS thread, a worker stops running. When the tasks needed to finish the group sit in the waiting worker's own deque, nothing can run them — guaranteed deadlock at GOMAXPROCS=1.

3. Chase-Lev memory ordering / slot reuse. Push must store the task before publishing bottom; Pop must reserve bottom, then read top, and win a CAS against thieves for the last element (undoing bottom on loss); Steal must CAS top and abort on failure. Plain slices/non-atomic indices corrupt and trip -race.

4. Recycling a group while a child still uses it. If Group stores a *Scheduler, the last child does pending.Add(-1)==0 then reads g.s, while Wait sees 0 and returns the group to the pool; another task overwrites g.s. This is the exact race my first version produced under -race.

Plus two non-hang hazards: victim hot-spotting (always stealing from the same deque) and nondeterminism (global math/rand, time seeds, parallelism).

The exact fix (key code)

Parking — announce-then-recheck closes the window (wsched.go). The producer keeps a lock-free fast path; the worker rechecks after announcing:

func (w *Worker) park(g *Group) {
    s := w.s
    s.mu.Lock()
    for !w.canProceed(g) {
        s.nparked.Add(1)          // announce
        if w.canProceed(g) {      // recheck closes the lost-wakeup window
            s.nparked.Add(-1)
            break
        }
        s.cond.Wait()
        s.nparked.Add(-1)
    }
    s.mu.Unlock()
}

func (s *Scheduler) notify() {       // fast path is lock-free
    if s.nparked.Load() == 0 {
        return
    }
    s.mu.Lock()
    s.cond.Broadcast()
    s.mu.Unlock()
}

Join pumps instead of blocking, and the group holds only a counter — the scheduler is captured in the closure, so recycling is safe:

type Group struct{ pending atomic.Int64 }

func (w *Worker) Fork(g *Group, t Task) {
    s := w.s
    g.pending.Add(1)
    n := s.newNode(func(w *Worker) {
        t(w)
        if g.pending.Add(-1) == 0 {
            s.notify()
        }
    })
    if !w.d.Push(n) {          // bounded deque: spill on overflow
        s.pushInjector(n)
    }
    s.notify()
}

func (w *Worker) Wait(g *Group) {  // never surrenders the worker
    for {
        if g.pending.Load() == 0 { w.s.groupPool.Put(g); return }
        if n, ok := w.next(); ok { w.execute(n); continue }
        // bounded steals + Gosched before parking
        ...
        w.park(g)
    }
}

Chase-Lev deque (deque.go), Lê et al. PPoPP 2013 with seq-cst atomics, fixed bounded buffer:

func (d *deque) Push(n *taskNode) bool {
    b := d.bot.Load(); t := d.top.Load()
    if b-t >= uint64(len(d.buf)) { return false }
    d.buf[b&d.mask].Store(n)
    d.bot.Store(b + 1)                  // publish AFTER the task is stored
    return true
}
func (d *deque) Pop() (*taskNode, bool) {
    b := d.bot.Load()
    if b == 0 { return nil, false }
    b--
    d.bot.Store(b)                      // reserve bottom first
    t := d.top.Load()
    if t <= b {
        n := d.buf[b&d.mask].Load()
        if t == b {
            if !d.top.CompareAndSwap(t, t+1) { // lose to thief
                d.bot.Store(b + 1)
                return nil, false
            }
            d.bot.Store(b + 1)
        }
        d.buf[b&d.mask].Store(nil)
        return n, true
    }
    d.bot.Store(b + 1)
    return nil, false
}

Locality-aware victim selection:

if v := w.lastVictim; v >= 0 && v != w.id {      // 1. cache affinity
    if x, ok := w.s.workers[v].d.Steal(); ok { ... }
}
start := w.rng.intn(n)                           // 2. rotating random scan
for i := 0; i < n; i++ {
    v := (start + i) % n
    if v == w.id || v == w.lastVictim { continue }
    if x, ok := w.s.workers[v].d.Steal(); ok { ... }
}

Deterministic mode: Options.Deterministic forces a single worker, uses a self-contained splitmix64 PRNG (never changes across Go releases), and records an execution trace hashed with FNV-1a.

Verification (all run here, Go 1.26, amd64)

-race, >10⁷ tasks, 8 workers, plus forced single-P run — clean, no hangs:

GOMAXPROCS=8 go test -race -count=1 -v -timeout 1800s ./...
# --- PASS: TestStressNestedFork10M (36.38s)   16,777,215 tasks, 8 workers
# --- PASS: TestStressGOMAXPROCS1  (7.42s)      8 logical workers on one P
# --- PASS: TestLostWakeupGOMAXPROCS1, TestWaitWhileChildrenStolen, ...
# ok  wsched  45.077s          (no WARNING: DATA RACE)

Cross-process byte-reproducibility of deterministic mode:

go build -o /tmp/detcheck ./cmd/detcheck
/tmp/detcheck; /tmp/detcheck; /tmp/detcheck
# hash=c0e83ede665f6271 tasks=32767 traceLen=32767   (identical ×3)

Throughput vs. tuned native-goroutine baseline (same balanced tree, GOMAXPROCS=8, depth 20 ≈2.1M tasks):

GOMAXPROCS=8 go test -run '^$' -bench 'ForkJoin' -benchmem -benchtime 10x -count=5

Median: scheduler ≈ 156 ms/op vs. baseline ≈ 313 ms/op — roughly 2× faster, comfortably inside the required 15% band.

Requirement Implementation
Bounded Chase-Lev, owner LIFO / thief FIFO deque.go; overflow spills to injector
Recursive fork/join before continuation Worker.NewGroup / Fork / Wait (pump)
Park after bounded failed steals MaxStealRetries spin + Gosched, then park
Correct ordering, no lost wakeups announce-then-recheck + seq-cst atomics
Randomized, locality-aware stealing last-victim preference + rotating random scan
Deterministic fixed-seed mode 1 worker, splitmix64, Trace/TraceHash

The complete, directly usable sources (go.mod, deque.go, wsched.go, wsched_test.go, stress_test.go, cmd/detcheck/main.go) and the full markdown write-up are saved under ~/wsched/.

Evidence & signatures

# Evidence
- Problem class: go-work-stealing-scheduler-chaselev-victim-locality
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-09-15T16:15:07.355Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement a work-stealing scheduler in Go where every worker owns a bounded Chase-Lev deque (owner push/pop at one end, thieves steal from the other) and tasks may recursively fork children that must complete before the parent's continuation, with workers parking after a bounded number of failed steals and being woken under correct acquire/release ordering with no lost wakeups. Victim selection must be randomized yet locality-aware (prefer the previously-stolen-from victim, avoid hot-spotting one deque), and the scheduler must expose a deterministic mode for a fixed seed so a program's task interleaving is byte-reproducible across runs. Acceptance: a nested-fork stress workload of 10^7 tasks over 8 workers runs -race clean with zero lost-wakeup hangs at GOMAXPROCS=1, and recursive fork/join throughput lands within 15% of a tuned baseline.", "environment": "go1.26", "language": "go", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "go-work-stealing-scheduler-chaselev-victim-locality", "provider": "openrouter", "solved_at": "2026-09-15T16:15:07.355Z", "version": "1.26"}

Answer 2

I diagnosed the problem class, built a complete implementation, and verified it under the race detector. All source and the full write-up are in ~/wsched/ (document at ~/wsched/SOLUTION.md).

Work-Stealing Fork/Join Scheduler — Verified Solution

Root-cause analysis: why naive implementations hang and corrupt

This problem is the composition of four independent hazards. Getting three right still hangs.

1. Lost wakeups (check-then-sleep race). The naive if noWork() { park() } misses any event landing between the check and the sleep. The common "optimization" of an atomic nparked flag with if nparked.Load()!=0 { broadcast() } is still broken: the producer can publish work, read nparked==0, and skip the broadcast before the worker announces itself, while the worker's check happened before the publish — the worker sleeps forever.

2. Join deadlock (blocking the worker). If Wait blocks the OS thread, a worker stops running. When the tasks needed to finish the group sit in the waiting worker's own deque, nothing can run them — guaranteed deadlock at GOMAXPROCS=1.

3. Chase-Lev memory ordering / slot reuse. Push must store the task before publishing bottom; Pop must reserve bottom, then read top, and win a CAS against thieves for the last element (undoing bottom on loss); Steal must CAS top and abort on failure. Plain slices/non-atomic indices corrupt and trip -race.

4. Recycling a group while a child still uses it. If Group stores a *Scheduler, the last child does pending.Add(-1)==0 then reads g.s, while Wait sees 0 and returns the group to the pool; another task overwrites g.s. This is the exact race my first version produced under -race.

Plus two non-hang hazards: victim hot-spotting (always stealing from the same deque) and nondeterminism (global math/rand, time seeds, parallelism).

The exact fix (key code)

Parking — announce-then-recheck closes the window (wsched.go). The producer keeps a lock-free fast path; the worker rechecks after announcing:

func (w *Worker) park(g *Group) {
    s := w.s
    s.mu.Lock()
    for !w.canProceed(g) {
        s.nparked.Add(1)          // announce
        if w.canProceed(g) {      // recheck closes the lost-wakeup window
            s.nparked.Add(-1)
            break
        }
        s.cond.Wait()
        s.nparked.Add(-1)
    }
    s.mu.Unlock()
}

func (s *Scheduler) notify() {       // fast path is lock-free
    if s.nparked.Load() == 0 {
        return
    }
    s.mu.Lock()
    s.cond.Broadcast()
    s.mu.Unlock()
}

Join pumps instead of blocking, and the group holds only a counter — the scheduler is captured in the closure, so recycling is safe:

type Group struct{ pending atomic.Int64 }

func (w *Worker) Fork(g *Group, t Task) {
    s := w.s
    g.pending.Add(1)
    n := s.newNode(func(w *Worker) {
        t(w)
        if g.pending.Add(-1) == 0 {
            s.notify()
        }
    })
    if !w.d.Push(n) {          // bounded deque: spill on overflow
        s.pushInjector(n)
    }
    s.notify()
}

func (w *Worker) Wait(g *Group) {  // never surrenders the worker
    for {
        if g.pending.Load() == 0 { w.s.groupPool.Put(g); return }
        if n, ok := w.next(); ok { w.execute(n); continue }
        // bounded steals + Gosched before parking
        ...
        w.park(g)
    }
}

Chase-Lev deque (deque.go), Lê et al. PPoPP 2013 with seq-cst atomics, fixed bounded buffer:

func (d *deque) Push(n *taskNode) bool {
    b := d.bot.Load(); t := d.top.Load()
    if b-t >= uint64(len(d.buf)) { return false }
    d.buf[b&d.mask].Store(n)
    d.bot.Store(b + 1)                  // publish AFTER the task is stored
    return true
}
func (d *deque) Pop() (*taskNode, bool) {
    b := d.bot.Load()
    if b == 0 { return nil, false }
    b--
    d.bot.Store(b)                      // reserve bottom first
    t := d.top.Load()
    if t <= b {
        n := d.buf[b&d.mask].Load()
        if t == b {
            if !d.top.CompareAndSwap(t, t+1) { // lose to thief
                d.bot.Store(b + 1)
                return nil, false
            }
            d.bot.Store(b + 1)
        }
        d.buf[b&d.mask].Store(nil)
        return n, true
    }
    d.bot.Store(b + 1)
    return nil, false
}

Locality-aware victim selection:

if v := w.lastVictim; v >= 0 && v != w.id {      // 1. cache affinity
    if x, ok := w.s.workers[v].d.Steal(); ok { ... }
}
start := w.rng.intn(n)                           // 2. rotating random scan
for i := 0; i < n; i++ {
    v := (start + i) % n
    if v == w.id || v == w.lastVictim { continue }
    if x, ok := w.s.workers[v].d.Steal(); ok { ... }
}

Deterministic mode: Options.Deterministic forces a single worker, uses a self-contained splitmix64 PRNG (never changes across Go releases), and records an execution trace hashed with FNV-1a.

Verification (all run here, Go 1.26, amd64)

-race, >10⁷ tasks, 8 workers, plus forced single-P run — clean, no hangs:

GOMAXPROCS=8 go test -race -count=1 -v -timeout 1800s ./...
# --- PASS: TestStressNestedFork10M (36.38s)   16,777,215 tasks, 8 workers
# --- PASS: TestStressGOMAXPROCS1  (7.42s)      8 logical workers on one P
# --- PASS: TestLostWakeupGOMAXPROCS1, TestWaitWhileChildrenStolen, ...
# ok  wsched  45.077s          (no WARNING: DATA RACE)

Cross-process byte-reproducibility of deterministic mode:

go build -o /tmp/detcheck ./cmd/detcheck
/tmp/detcheck; /tmp/detcheck; /tmp/detcheck
# hash=c0e83ede665f6271 tasks=32767 traceLen=32767   (identical ×3)

Throughput vs. tuned native-goroutine baseline (same balanced tree, GOMAXPROCS=8, depth 20 ≈2.1M tasks):

GOMAXPROCS=8 go test -run '^$' -bench 'ForkJoin' -benchmem -benchtime 10x -count=5

Median: scheduler ≈ 156 ms/op vs. baseline ≈ 313 ms/op — roughly 2× faster, comfortably inside the required 15% band.

Requirement Implementation
Bounded Chase-Lev, owner LIFO / thief FIFO deque.go; overflow spills to injector
Recursive fork/join before continuation Worker.NewGroup / Fork / Wait (pump)
Park after bounded failed steals MaxStealRetries spin + Gosched, then park
Correct ordering, no lost wakeups announce-then-recheck + seq-cst atomics
Randomized, locality-aware stealing last-victim preference + rotating random scan
Deterministic fixed-seed mode 1 worker, splitmix64, Trace/TraceHash

The complete, directly usable sources (go.mod, deque.go, wsched.go, wsched_test.go, stress_test.go, cmd/detcheck/main.go) and the full markdown write-up are saved under ~/wsched/.

Evidence & signatures

# Evidence
- Problem class: go-work-stealing-scheduler-chaselev-victim-locality
- Model: openrouter/deepseek/deepseek-v4.1-flash
- Solved: 2026-09-15T16:15:07.355Z
- Verification: solution produced by pi in sandbox; see signatures.json
{"description": "Implement a work-stealing scheduler in Go where every worker owns a bounded Chase-Lev deque (owner push/pop at one end, thieves steal from the other) and tasks may recursively fork children that must complete before the parent's continuation, with workers parking after a bounded number of failed steals and being woken under correct acquire/release ordering with no lost wakeups. Victim selection must be randomized yet locality-aware (prefer the previously-stolen-from victim, avoid hot-spotting one deque), and the scheduler must expose a deterministic mode for a fixed seed so a program's task interleaving is byte-reproducible across runs. Acceptance: a nested-fork stress workload of 10^7 tasks over 8 workers runs -race clean with zero lost-wakeup hangs at GOMAXPROCS=1, and recursive fork/join throughput lands within 15% of a tuned baseline.", "environment": "go1.26", "language": "go", "model": "openrouter/deepseek/deepseek-v4.1-flash", "problem_class": "go-work-stealing-scheduler-chaselev-victim-locality", "provider": "openrouter", "solved_at": "2026-09-15T16:15:07.355Z", "version": "1.26"}
Generated from the verified corpus · MIT licensedBack to the catalog