Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ Versions follow [SemVer](https://semver.org).

## [Unreleased]

- Include GPU count and resolved GPU type in dispatched eval and baseline cache identity, including budget discounts. Legacy eval slots and baseline entries are cache misses. Upgrading: An eval dispatched by the previous kernel and still in flight at upgrade is measured again once under the new cache key (no extra budget charge). To avoid the extra run, upgrade when no evals are in flight: `touch <root>/PAUSE` (stops wakes, so no new evals are dispatched), wait until no eval jobs remain in the queue, upgrade, then `rm <root>/PAUSE` and run `outerloop start`.

- Park launch capacity refusals in capacity wait when an immediate resume is
unavailable or already refused, instead of ending the run. The next wake
delivers the refusal and a retry note, resuming the same session when supported
Expand Down
18 changes: 13 additions & 5 deletions src/outerloop/measure.py
Original file line number Diff line number Diff line change
Expand Up @@ -182,12 +182,13 @@ def read_baseline_cache(
metric: str = "",
seed_env: str = "",
gpus: int = 0,
gpu_type: str = "",
) -> dict[str, Any] | None:
"""The cached base-tree measurement for (benchmark, base sha), or None.
The entry must have been measured under the SAME determinants the
candidate will be — eval image, contract command, metric key, seed
variable, GPU count: everything the measurer's own eval identity
carries except the tree sha (the key) and the seed VALUE (fresh per
variable, GPU count and resolved GPU type: everything the measurer's
own eval identity carries except the tree sha (the key) and the seed VALUE (fresh per
attempt by design) — or it is stale (terra #178): a comparison across
determinants is not a comparison. A cache
entry is only ever written from an orchestrator-measured value (below),
Expand All @@ -207,7 +208,8 @@ def read_baseline_cache(
or data.get("command", "") != command
or data.get("metric", "") != metric
or data.get("seed_env", "") != seed_env
or int(data.get("gpus", 0) or 0) != gpus
or data.get("gpus") != gpus
or data.get("gpu_type") != gpu_type
):
return None
return data
Expand All @@ -226,6 +228,7 @@ def write_baseline_cache(
metric: str = "",
seed_env: str = "",
gpus: int = 0,
gpu_type: str = "",
) -> None:
"""Record an orchestrator-measured baseline for every later attempt on
this base, with the determinants it was measured under. Atomic (a
Expand All @@ -250,6 +253,7 @@ def write_baseline_cache(
"metric": metric,
"seed_env": seed_env,
"gpus": gpus,
"gpu_type": gpu_type,
},
fh,
)
Expand Down Expand Up @@ -302,15 +306,19 @@ def _det(self, m: Measure) -> str:
# Everything a measure's RESULT depends on and that can vary across a
# PARK/RESUME (when a fresh measurer reads this run_dir): the container
# image, the measure's logical role, the code (tree_sha), and the
# contract facts it is evaluated under (command, metric, seeded env).
# contract facts it is evaluated under (command, metric, seeded env),
# GPU count, and the resolved GPU type from the lane.
# A cache key missing any of these would return a value computed under
# DIFFERENT inputs — e.g. a resume that re-fetched the contract after
# its command changed, or ran under a rebuilt image, reading the stale
# pre-change result. (account / walltime don't change a result's value,
# only whether it completes.) NUL separators keep the parts unambiguous
# (`a`+`bc` != `ab`+`c`).
env = "".join(f"\0{k}={v}" for k, v in sorted(m.env().items()))
return f"{self.image}\0{m.name}\0{m.tree_sha}\0{m.command}\0{m.metric}{env}"
return (
f"{self.image}\0{m.name}\0{m.tree_sha}\0{m.command}\0{m.metric}"
f"\0{m.gpus}\0{self.gpu_type}{env}"
)

def _slot(self, m: Measure) -> str:
# Storage identity = the full determinant, with NOTHING truncated: the
Expand Down
3 changes: 3 additions & 0 deletions src/outerloop/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -918,6 +918,7 @@ def measure_and_decide(
metric=bench.metric,
seed_env=bench.seed_env or "",
gpus=bench.gpus,
gpu_type=str(getattr(measurer, "gpu_type", "")),
)
if bench.baseline == "cached" and cache_dir is not None
else None
Expand Down Expand Up @@ -960,6 +961,7 @@ def measure_and_decide(
metric=bench.metric,
seed_env=bench.seed_env or "",
gpus=bench.gpus,
gpu_type=str(getattr(measurer, "gpu_type", "")),
)
candidate = main["candidate"]
if not improved(baseline, candidate, bench.direction, min_relative_improvement):
Expand Down Expand Up @@ -1796,6 +1798,7 @@ def _not_run_note(request: SyscallRequest | None) -> str:
metric=bench.metric,
seed_env=bench.seed_env or "",
gpus=bench.gpus,
gpu_type=str(getattr(measurer, "gpu_type", "")),
):
main_evals = 1
capacity_refused = False
Expand Down
14 changes: 14 additions & 0 deletions tests/fixtures/README.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,19 @@
# Endpoint compatibility fixtures

- `dispatched_pre_gpu_identity/`: run directories produced by
`DispatchedMeasurer._dispatch` from kernel `d4ad7529667de79aef299e2e0ad8b5fb6c9c9598`,
before GPU count/type entered the determinant. A fake Slurm submit returned
job `101`; the completed variant adds the job's exit-code and metric stdout.
Both use a synthetic repo/image, one GPU, and `SEED=7`. `identity.json`
records that kernel's slot and scheduler name. The new reader ignores and
preserves the directories (including when interrupted), so old jobs can
finish writing and
rollback can still read the old records. Legacy slots are cache misses; the
new kernel dispatches once under the GPU-aware key without charging the
resumed gate again. Newly dispatched GPU-aware slots
are not readable by the old kernel without re-measurement. Cross-run legacy
baseline entries remain misses.

- `author_route_legacy.json`: produced by `runstate.RunRecord` and
`runstate.save_record` from kernel commit `5563c46` (the parent tree before the
endpoint change), with synthetic `owner/repo` and key-file coordinates.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
cmd
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
0
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
#!/bin/sh
set -u
EV=tests/fixtures/dispatched_pre_gpu_identity/completed/eval-candidate-aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa-f69684777d523fe6502ee4ed4d09f98410a56aec
REPO=/fixture/repo
SCRATCH="${SLURM_TMPDIR:-${TMPDIR:-/tmp}}/dispatch-eval-$$"
TREE="$SCRATCH/tree"
mkdir -p "$SCRATCH/cache" "$SCRATCH/home" "$SCRATCH/work" "$TREE"
cleanup() { rm -rf "$SCRATCH"; git -C "$REPO" -c core.hooksPath=/dev/null -c core.sshCommand=false -c credential.helper= -c protocol.allow=never -c protocol.https.allow=always -c protocol.file.allow=always -c core.fsmonitor= -c core.quotePath=false worktree prune >/dev/null 2>&1 || true; }
trap 'cleanup' EXIT
trap 'echo 143 > "$EV/exit-code"; cleanup; trap - EXIT; exit 0' TERM INT HUP
[ -s "$EV/command.txt" ] || { echo 96 > "$EV/exit-code"; exit 0; }
export GIT_CONFIG_GLOBAL=/dev/null GIT_CONFIG_SYSTEM=/dev/null
export GIT_CONFIG_COUNT=1
export GIT_CONFIG_KEY_0=core.attributesFile
export GIT_CONFIG_VALUE_0=/dev/null
if git -C "$REPO" -c core.hooksPath=/dev/null -c core.sshCommand=false -c credential.helper= -c protocol.allow=never -c protocol.https.allow=always -c protocol.file.allow=always -c core.fsmonitor= -c core.quotePath=false worktree add --detach "$TREE" aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa >> "$EV/setup.log" 2>&1; then rm -f "$TREE/.git"; else echo 97 > "$EV/exit-code"; exit 0; fi
export UV_CACHE_DIR="$SCRATCH/cache" UV_LINK_MODE=copy UV_PROJECT_ENVIRONMENT="$SCRATCH/cache/venv"
export APPTAINERENV_UV_CACHE_DIR="$UV_CACHE_DIR" APPTAINERENV_UV_LINK_MODE=copy APPTAINERENV_UV_PROJECT_ENVIRONMENT="$UV_PROJECT_ENVIRONMENT"
export SEED=7 APPTAINERENV_SEED=7
apptainer exec --containall --cleanenv --nv --bind "$TREE:$TREE" --home "$SCRATCH/home:$SCRATCH/home" --bind "$SCRATCH/cache:$SCRATCH/cache" --pwd "$TREE" --workdir "$SCRATCH/work" /img.sif sh -c "$(cat "$EV/command.txt")" > "$EV/stdout" 2> "$EV/stderr"
echo $? > "$EV/exit-code"
exit 0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
{"commit": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", "author": {}}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
{"metric": "r2", "value": 0.42}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
101
4 changes: 4 additions & 0 deletions tests/fixtures/dispatched_pre_gpu_identity/identity.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"slot": "candidate-aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa-f69684777d523fe6502ee4ed4d09f98410a56aec",
"job_name": "eval-r1-candidate-b874c69af4063bd5"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
cmd
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
#!/bin/sh
set -u
EV=tests/fixtures/dispatched_pre_gpu_identity/inflight/eval-candidate-aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa-f69684777d523fe6502ee4ed4d09f98410a56aec
REPO=/fixture/repo
SCRATCH="${SLURM_TMPDIR:-${TMPDIR:-/tmp}}/dispatch-eval-$$"
TREE="$SCRATCH/tree"
mkdir -p "$SCRATCH/cache" "$SCRATCH/home" "$SCRATCH/work" "$TREE"
cleanup() { rm -rf "$SCRATCH"; git -C "$REPO" -c core.hooksPath=/dev/null -c core.sshCommand=false -c credential.helper= -c protocol.allow=never -c protocol.https.allow=always -c protocol.file.allow=always -c core.fsmonitor= -c core.quotePath=false worktree prune >/dev/null 2>&1 || true; }
trap 'cleanup' EXIT
trap 'echo 143 > "$EV/exit-code"; cleanup; trap - EXIT; exit 0' TERM INT HUP
[ -s "$EV/command.txt" ] || { echo 96 > "$EV/exit-code"; exit 0; }
export GIT_CONFIG_GLOBAL=/dev/null GIT_CONFIG_SYSTEM=/dev/null
export GIT_CONFIG_COUNT=1
export GIT_CONFIG_KEY_0=core.attributesFile
export GIT_CONFIG_VALUE_0=/dev/null
if git -C "$REPO" -c core.hooksPath=/dev/null -c core.sshCommand=false -c credential.helper= -c protocol.allow=never -c protocol.https.allow=always -c protocol.file.allow=always -c core.fsmonitor= -c core.quotePath=false worktree add --detach "$TREE" aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa >> "$EV/setup.log" 2>&1; then rm -f "$TREE/.git"; else echo 97 > "$EV/exit-code"; exit 0; fi
export UV_CACHE_DIR="$SCRATCH/cache" UV_LINK_MODE=copy UV_PROJECT_ENVIRONMENT="$SCRATCH/cache/venv"
export APPTAINERENV_UV_CACHE_DIR="$UV_CACHE_DIR" APPTAINERENV_UV_LINK_MODE=copy APPTAINERENV_UV_PROJECT_ENVIRONMENT="$UV_PROJECT_ENVIRONMENT"
export SEED=7 APPTAINERENV_SEED=7
apptainer exec --containall --cleanenv --nv --bind "$TREE:$TREE" --home "$SCRATCH/home:$SCRATCH/home" --bind "$SCRATCH/cache:$SCRATCH/cache" --pwd "$TREE" --workdir "$SCRATCH/work" /img.sif sh -c "$(cat "$EV/command.txt")" > "$EV/stdout" 2> "$EV/stderr"
echo $? > "$EV/exit-code"
exit 0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
{"commit": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", "author": {}}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
101
71 changes: 71 additions & 0 deletions tests/test_attempt.py
Original file line number Diff line number Diff line change
Expand Up @@ -7960,3 +7960,74 @@ def test_line_snapshot_restores_protected_paths_only_in_seal(tmp_path, target_re
assert not logs # Only ancestry was dropped; the final content was already restored.
else:
assert len(logs) == 1 and protected in logs[0]


def test_legacy_eval_redispatch_preserves_run_gpu_meter(tmp_path, monkeypatch):
import hashlib

from outerloop.compute import CommandResult, SlurmCompute
from outerloop.measure import DispatchSettings, plan_measures

real_measurer = DispatchSettings.measurer
state, run_id = _write_parked_candidate(tmp_path, monkeypatch, contract=CONTRACT_GPU)
monkeypatch.setattr(DispatchSettings, "measurer", real_measurer)
record = load_record(state, run_id)
record.stage.update(submitted=True, gpu_hours_used=1.0, sleeps_used=1, launches_used=0)
save_record(state, record, 1_000_000.0)
run_dir = state / "runs" / run_id
bench = load_contract(CONTRACT_GPU, "org/pilot").benchmarks[0]
measures = plan_measures(
bench.command,
bench.metric,
str(record.stage["base_sha"]),
str(record.stage["candidate_sha"]),
gpus=1,
)
# Old-format completed slots exist, but cannot satisfy the upgraded gate.
for measure in measures:
legacy_det = (
f"/img.sif\0{measure.name}\0{measure.tree_sha}\0{measure.command}\0{measure.metric}"
)
slot = run_dir / (
f"eval-{measure.name}-{measure.tree_sha}-{hashlib.sha1(legacy_det.encode()).hexdigest()}"
)
slot.mkdir()
(slot / "submitted").write_text("501")
(slot / "exit-code").write_text("0\n")
(slot / "stdout").write_text('{"metric":"mean_tour_length","value":13.0}\n')
submitted = []
live = {}

def runner(argv, timeout_s):
if argv[0] == "sbatch":
submitted.append(argv)
job = str(600 + len(submitted))
name = next(a.split("=", 1)[1] for a in argv if a.startswith("--job-name="))
live[name] = job
return CommandResult(0, job + "\n", "")
if argv[0] == "squeue" and "--name" in argv:
return CommandResult(0, live.get(argv[argv.index("--name") + 1], ""), "")
return CommandResult(0, "PENDING\n" if argv[0] == "sacct" else "", "")

dispatch = DispatchSettings(
compute=SlurmCompute(runner=runner),
image="/img.sif",
account="acct",
partition="cpu",
gpu_partition="gpu",
)
for now in (1_000_100.0, 1_000_200.0):
outcome = resume_run(
state,
run_id,
dispatch=dispatch,
github=CommentingGitHub(), # type: ignore[arg-type]
bot_auth=NoAuth(),
now=now,
)
assert outcome.outcome == "parked"
assert len(submitted) == 2
saved = load_record(state, run_id)
assert saved.stage["afterany"] == "afterany:601:602"
assert saved.stage["gpu_hours_used"] == 1.0
assert saved.stage["sleeps_used"] == 1
98 changes: 98 additions & 0 deletions tests/test_measure.py
Original file line number Diff line number Diff line change
Expand Up @@ -551,3 +551,101 @@ def runner(argv, timeout_s):
scripts = {n: Path(a[-1]).read_text() for n, a in by_job.items()}
for n, text in scripts.items():
assert ("--nv" in text) == ("-sib-speedru" in n)


@pytest.mark.parametrize(
"gpus,gpu_type,reused", [(1, "type-a", True), (2, "type-a", False), (1, "type-b", False)]
)
def test_dispatched_cache_resource_identity(tmp_path, gpus, gpu_type, reused):
from dataclasses import replace

submitted: list = []
m = _measurer(tmp_path, submitted)
m.gpu_partition = "gpu"
m.gpu_type = "type-a"
original = Measure("candidate", "a" * 40, "cmd", "r2", gpus=1)
_land(m, original, 0.42)
resumed = replace(m, gpu_type=gpu_type)
measure = replace(original, gpus=gpus)
if reused:
assert resumed.results([measure]) == {"candidate": 0.42}
assert not submitted
else:
assert resumed._job_name(measure) != m._job_name(original)
with pytest.raises(MeasurementPending):
resumed.results([measure])
assert len(submitted) == 1


@pytest.fixture
def legacy_eval_run(tmp_path):
"""Copy durable output from the previous kernel, without new-code writers."""
import shutil

source = Path(__file__).parent / "fixtures" / "dispatched_pre_gpu_identity"
identity = json.loads((source / "identity.json").read_text())

def copy(state):
run_dir = tmp_path / state
shutil.copytree(source / state, run_dir)
return run_dir, identity

return copy


@pytest.mark.parametrize("state", ["completed", "inflight"])
@pytest.mark.parametrize("interrupt", [False, True])
def test_legacy_slot_is_miss_and_redispatch_is_idempotent(
legacy_eval_run, state, interrupt, monkeypatch
):
from dataclasses import replace

run_dir, identity = legacy_eval_run(state)
legacy = run_dir / ("eval-" + identity["slot"])
before = {p.name: p.read_bytes() for p in legacy.iterdir() if p.is_file()}
assert b"#SBATCH" not in before["job.sh"]
submitted: list = []
live = {identity["job_name"]: "old-job"} if state == "inflight" else {}
m = _measurer(run_dir, submitted, live=live)
m.gpu_partition = "gpu"
m.gpu_type = "type-a"
measure = Measure("candidate", "a" * 40, "cmd", "r2", (("SEED", "7"),), 1)
assert m._slot(measure) != identity["slot"]
assert m._job_name(measure) != identity["job_name"]

if interrupt:
# Die after sbatch accepts the new job, before its marker is written.
write_text = Path.write_text

def interrupted_write(path, *args, **kwargs):
if path == m._ev(measure) / "submitted":
raise KeyboardInterrupt
return write_text(path, *args, **kwargs)

with monkeypatch.context() as patch:
patch.setattr(Path, "write_text", interrupted_write)
with pytest.raises(KeyboardInterrupt):
m.results([measure])
assert not (m._ev(measure) / "submitted").exists()
else:
with pytest.raises(MeasurementPending) as caught:
m.results([measure])
assert caught.value.job_ids == ("101",)
assert len(submitted) == 1
assert Path(submitted[0][-1]) == m._ev(measure) / "job.sh"

live[m._job_name(measure)] = "101"
resumed = _measurer(run_dir, submitted, live=live)
resumed.gpu_partition = "gpu"
resumed.gpu_type = "type-a"
for _ in range(2):
with pytest.raises(MeasurementPending) as caught:
replace(resumed).results([measure])
assert caught.value.job_ids == ("101",)
assert len(submitted) == 1
assert before == {p.name: p.read_bytes() for p in legacy.iterdir() if p.is_file()}

# Only the new slot supplies a result, even alongside a completed old slot.
_land(resumed, measure, 0.9)
assert replace(resumed).results([measure]) == {"candidate": 0.9}
assert len(submitted) == 1
15 changes: 15 additions & 0 deletions tests/test_measure_and_decide.py
Original file line number Diff line number Diff line change
Expand Up @@ -439,3 +439,18 @@ def test_cached_baseline_is_keyed_by_image_and_command(tmp_path):
)
got = read_baseline_cache(d, "main", BASE, image="/a.sif", command="run main")
assert got and got["value"] == 0.6 and not list(d.glob("*.tmp"))


@pytest.mark.parametrize("missing", [("gpus",), ("gpu_type",), ("gpus", "gpu_type")])
def test_baseline_cache_missing_resource_fields_is_a_miss(tmp_path, missing):
import json

from outerloop.measure import read_baseline_cache, write_baseline_cache

write_baseline_cache(tmp_path, "main", BASE, value=0.5, seed=0, run_tag="old")
path = tmp_path / f"main@{BASE}.json"
data = json.loads(path.read_text())
for key in missing:
del data[key]
path.write_text(json.dumps(data))
assert read_baseline_cache(tmp_path, "main", BASE) is None
Loading
Loading