Skip to content

Server panics in the scheduler's gap computation (workerload.rs:160 index out of bounds) with prioritised multi-variant requests and workers with different resource sets #1135

Description

@chansigit

Title: Server panics in the scheduler's gap computation (workerload.rs:160 index out of bounds) when prioritised tasks with multi-variant resource requests meet workers with different resource sets

Summary

With task priorities in use, the server dies within ~20 s of a mixed job stream:

thread 'main' panicked at crates/tako/src/internal/server/workerload.rs:160:37:
index out of bounds: the len is 3 but the index is 4

Frames (nightly 2026-09-25-21f2d2e8, RUST_BACKTRACE):
WorkerResources::remove (workerload.rs:160) ← GapCache::get_gap (scheduler/gap.rs) ← run_scheduling_inner ← server main loop.

We first hit this on v0.26.2 in production (2026-09-23, same line, same message with index 3), two minutes after
starting to submit with --priority; without priorities the same workload has run for days.

Setup that reproduces it (single node, all processes local)

Three workers on one host, started with --detect-resources none:

  • two "CPU" workers: --cpus 8 --resource mem=sum(60000) --resource runtime/aaa=sum(8)
  • one "GPU" worker: the same plus --resource "gpuSlot/0=[u0#0]" --resource "gpuMemoryMB/0=sum(72000)" --resource "gpuHostMB=sum(30000)"

A stream submitted every 2 s (see driver.py):

  • 6 × hq submit --priority 1 --cpus 1 --resource mem=2000 --resource runtime/aaa=1 -- sleep 8
  • 1 × --priority 2 --cpus 2 ... sleep 12, every 10 s 1 × --priority 4 --cpus 4 ... sleep 20
  • 1 × --priority 0 --cpus 0.125 --resource runtime/aaa=0.125 ... sleep 20
  • every 6 s a job file whose single task has priority = 4 and [[task.request]] variants
    {cpus=4, runtime/aaa=1, gpuSlot/<s>=1, gpuMemoryMB/<s>=20000, gpuHostMB=12000} for s in 0..N (N = 64 or 2)
    plus a CPU-only variant {cpus=4, runtime/aaa=1, mem=12000}
  • every 20 s the same without the CPU-only variant, priority = 8

The panic is racy but frequent: in our trials the server died within 16–18 s in most runs that included the
multi-variant, GPU-preferred jobs (with N = 2 as well as 64); with those jobs removed (only single-variant
prioritised tasks) it never died in 45 s trials. v0.26.2 survived 3 min of the same stream once but died in
production with the same message.

What the code says

WorkerResources::n_resources is sized by the worker's own highest declared resource id
(from_description), so a CPU-only worker registered after the GPU resources exist has a vector of length 3
while gpuSlot/0, gpuMemoryMB/0, gpuHostMB are ids 3–5. get() is bounds-safe (.get().unwrap_or(ZERO)),
but remove() / remove_multiple() index directly. In GapCache::get_gap, free.remove(rq) over the
assigned tasks' chosen variants (and remove_multiple(h_rq, count) for a trivial blocker, even with
count = 0) reaches a resource id the worker never declared → panic. main still has the direct indexing
(workerload.rs remove, remove_multiple).

Proposed fix

WorkerResources::remove / remove_multiple / remove_multiple_masked should treat a resource id past the
worker's vector the way get does (the worker holds none of it, nothing to subtract). A pull request with that
change and a regression test in gap.rs (test_compute_gap_resource_the_worker_lacks, which panics on
current main) follows.

Reproduction files

run.sh (starts server + 3 workers in a scratch server dir, runs the driver, prints the panic) and
driver.py (the stream, knobs via env: SLOTS, AGENT, RUNTIME, WIDE, TWO, PREF, ONLY).
RUST_LOG=hq=debug,tako=debug logs can be provided on request.

HyperQueue versions: v0.26.2 (production, 2026-09-23) and nightly-2026-09-25-21f2d2e8bdc3e5c20b1aae528bf4125b593aa5a7 (this reproduction).

run.sh
#!/bin/bash
# HQ priority-panic reproduction on a throwaway server: production-shaped workers and job stream.
# usage: run.sh <hq binary> <label> [seconds]   (all HQ processes run inside the control container)
set -u
HQBIN=$1; LABEL=$2; SECS=${3:-180}
SIF=/scratch/users/chensj16/containers/rsi-control-20260915-1.sif
DIR=/tmp/hq-repro-$LABEL; OUT=/scratch/users/chensj16/eca-runs/gen2-acceptance-20260917/ops/hq-repro/$LABEL
mkdir -p $DIR/jobs $OUT
HQ="apptainer exec --cleanenv --bind /scratch,/tmp $SIF $HQBIN --server-dir $DIR/server"
$HQ server start > $DIR/server.log 2>&1 &
SERVER=$!
sleep 3
COMMON="worker start --manager none --detect-resources none --resource mem=sum(60000)"; [ "${RUNTIME:-1}" = 1 ] && COMMON="$COMMON --resource runtime/aaa=sum(8)"
$HQ $COMMON --cpus 8 > $DIR/worker-cpu1.log 2>&1 & W1=$!
$HQ $COMMON --cpus 8 > $DIR/worker-cpu2.log 2>&1 & W2=$!
$HQ $COMMON --cpus 8 --resource "gpuSlot/0=[u0#0]" --resource "gpuMemoryMB/0=sum(72000)" --resource "gpuHostMB=sum(30000)" > $DIR/worker-gpu.log 2>&1 & W3=$!
sleep 3
$HQ worker list > $DIR/workers.txt 2>&1
python3 $OUT/../driver.py "$HQ" $DIR $SERVER $SECS > $DIR/driver.log 2>&1
echo "driver exit $?" >> $DIR/driver.log
$HQ job list --all > $DIR/jobs.txt 2>&1
kill $W1 $W2 $W3 2>/dev/null; sleep 1; kill $SERVER 2>/dev/null
cp $DIR/*.log $DIR/*.txt $OUT/ 2>/dev/null
tail -3 $DIR/driver.log; grep -m2 -iE "panick|index out of bounds" $DIR/server.log || echo "no panic in server.log"
driver.py
"""Submit a production-shaped stream of prioritised tasks and watch whether the HQ server survives."""
import os, subprocess, sys, time
hq, root, server_pid, seconds = sys.argv[1].split(), sys.argv[2], int(sys.argv[3]), int(sys.argv[4])
SLOTS = int(os.environ.get("SLOTS", 64))
ON = lambda k: os.environ.get(k, "1") == "1"   # AGENT, RUNTIME, WIDE, TWO, PREF, ONLY: knobs of the reduction series
RT = (["--resource", "runtime/aaa=1"] if ON("RUNTIME") else [])

def submit(*args):
    return subprocess.run(hq + list(args), capture_output=True, text=True, timeout=60)

def jobfile(name, priority, sleep, cpus, mem, cpu_variant):
    lines = [f'name = "{name}"', "[[task]]", f'command = ["sleep", "{sleep}"]', f'cwd = "{root}"',
             'pin = "taskset"', 'crash_limit = "never-restart"', f"priority = {priority}", 'stdout = "none"', 'stderr = "none"']
    base = {"cpus": cpus, **({"runtime/aaa": 1} if ON("RUNTIME") else {})}
    variants = [{**base, f"gpuSlot/{s}": 1, f"gpuMemoryMB/{s}": 20000, "gpuHostMB": mem} for s in range(SLOTS)]
    if cpu_variant:
        variants.append({**base, "mem": mem})
    for v in variants:
        lines += ["", "[[task.request]]", 'time_request = "600s"',
                  "resources = { " + ", ".join(f'"{k}" = {val}' for k, val in v.items()) + " }"]
    path = os.path.join(root, "jobs", name + ".toml")
    with open(path, "w") as f:
        f.write("\n".join(lines) + "\n")
    return path

def alive():
    try:
        os.kill(server_pid, 0)
        return True
    except OSError:
        return False

start, tick, n = time.time(), 0, 0
while time.time() - start < seconds:
    if not alive():
        print(f"SERVER DIED after {time.time()-start:.0f}s, {n} submissions"); sys.exit(2)
    batch = []
    for _ in range(6):
        batch.append(["submit", "--name", f"deg-{n}", "--priority", "1", "--cpus", "1", "--resource", "mem=2000", *RT, "--", "sleep", "8"]); n += 1
    if ON("TWO"):
        batch.append(["submit", "--name", f"two-{n}", "--priority", "2", "--cpus", "2", "--resource", "mem=6000", *RT, "--", "sleep", "12"]); n += 1
    if ON("AGENT"):
        batch.append(["submit", "--name", f"agent-{n}", "--priority", "0", "--cpus", "0.125", "--resource", "mem=384", *(["--resource", "runtime/aaa=0.125"] if ON("RUNTIME") else []), "--", "sleep", "20"]); n += 1
    if tick % 5 == 0 and ON("WIDE"):
        batch.append(["submit", "--name", f"wide-{n}", "--priority", "4", "--cpus", "4", "--resource", "mem=12000", *RT, "--", "sleep", "20"]); n += 1
    if tick % 3 == 0 and ON("PREF"):
        batch.append(["job", "submit-file", jobfile(f"gpu-pref-{n}", 4, 15, 4, 12000, True)]); n += 1
    if tick % 10 == 0 and ON("ONLY"):
        batch.append(["job", "submit-file", jobfile(f"gpu-only-{n}", 8, 10, 4, 12000, False)]); n += 1
    for args in batch:
        r = submit(*args)
        if r.returncode != 0:
            print(f"submit failed rc={r.returncode}: {' '.join(args[:4])}: {r.stderr.strip()[:200]}")
            if not alive():
                print(f"SERVER DIED after {time.time()-start:.0f}s, {n} submissions"); sys.exit(2)
    if tick % 10 == 0:
        print(f"{time.time()-start:.0f}s tick {tick} submissions {n} server alive {alive()}", flush=True)
    tick += 1
    time.sleep(2)
print(f"survived {seconds}s, {n} submissions, server alive {alive()}")

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions