From 9c3fbd9b488a6be98aea50a78793e62746aa19df Mon Sep 17 00:00:00 2001 From: Mengye Ren Date: Thu, 1 Oct 2026 10:46:03 -0400 Subject: [PATCH 1/3] A capacity refusal parks the run instead of ending it A launch refused for capacity after the in-session retry now parks the run in the capacity wait, keeping its session, candidate and meters; the next wake resumes the session with the refusal and a note that capacity may now be free. Other refusals and stale submits are unchanged. --- CHANGELOG.md | 5 ++ docs/design/lifecycle.md | 2 + src/outerloop/orchestrator.py | 39 ++++++++++++++ tests/test_orchestrator.py | 95 +++++++++++++++++++++++++++++++++++ tests/test_tick.py | 57 +++++++++++++++++++++ 5 files changed, 198 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index f987de3d..2ad2e275 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,11 @@ Versions follow [SemVer](https://semver.org). ## [Unreleased] +- Park repeated launch capacity refusals in capacity wait instead of ending the + run. The next wake resumes the same session with the refusal and a retry note; + meters are preserved and capacity waits do not exhaust stuck retries. + Existing records need no migration. + - Add operator `outerloop rebind [--root ] [--note ]` to explicitly move an existing run to its slot's current author at the next leg, including runs waiting on retired endpoints. Reports and board details show diff --git a/docs/design/lifecycle.md b/docs/design/lifecycle.md index 870854d6..9c21981b 100644 --- a/docs/design/lifecycle.md +++ b/docs/design/lifecycle.md @@ -250,6 +250,8 @@ and the contract's allowed scope (up to ten entries and 1,000 characters). Later refusals in the same run start with “Refused again:”. Nothing is sealed, launched, measured, or charged for a refused request. Refusals repeat and never end the run, including requests with malformed content or budget problems. +When a launch capacity refusal cannot be delivered in another immediate resume, +the run parks uncharged in capacity wait; its next wake receives the refusal and can retry. The session walltime and contract sleep, launch, and GPU-hour budgets bound the loop. A session that cannot resume ends as `session-error`; endpoint loss during refusal delivery ends as `session-outage` without sealing the rejected diff --git a/src/outerloop/orchestrator.py b/src/outerloop/orchestrator.py index 6f70356f..37d67828 100644 --- a/src/outerloop/orchestrator.py +++ b/src/outerloop/orchestrator.py @@ -1798,6 +1798,7 @@ def _not_run_note(request: SyscallRequest | None) -> str: gpus=bench.gpus, ): main_evals = 1 + capacity_refused = False problem = no_backend or syscall_budget_error( request, launches_used=launches_used, @@ -1913,6 +1914,7 @@ def _not_run_note(request: SyscallRequest | None) -> str: launch_afterany = launcher(sha, request) except CapacityError as exc: problem = str(exc) + capacity_refused = True else: gpu_hours_used += launches_gpu_hours(request, gpus=bench.gpus) raise RunParked( @@ -1939,6 +1941,43 @@ def _not_run_note(request: SyscallRequest | None) -> str: note=f"Stale submit refused; no gate or sibling launches ran: {problem}", run_seed=run_seed, ) + if capacity_refused: + # Bound immediate retries without turning transient capacity + # into an ending. The wake delivers this still-pending note. + append( + inbox_dir, + Message( + 0, + "note", + "kernel", + inbox_thread, + time.time(), + f"refusal:{session.session_id}:{inbox_seq}", + { + "text": "Your syscall request was REFUSED and nothing was " + "launched. The run was parked waiting for capacity. " + "Capacity may now be free; you can retry the launch.", + "quoted_text": problem, + "context_only": True, + }, + origin=inbox_dir.name, + ), + ) + raise RunParked( + phase="author-sleep", + afterany="", + base_sha=base_sha, + seed=run_seed, + suite_seed=suite_seed, + candidate_sha=sha, + session=session, + syscall=SyscallRequest(launches=()), + launches_used=launches_used, + sleeps_used=sleeps_used, + gpu_hours_used=gpu_hours_used, + judged=failed_gate, + capacity_wait=True, + ) log.warning("syscall request dropped after refusal (%s); measuring as-is", problem) break # the refusal burns no count (nothing was launched, nothing woke a diff --git a/tests/test_orchestrator.py b/tests/test_orchestrator.py index ed265023..76de320c 100644 --- a/tests/test_orchestrator.py +++ b/tests/test_orchestrator.py @@ -2364,6 +2364,101 @@ def launcher(sha, request): assert len(calls) == 1 +@pytest.mark.parametrize("can_resume", [False, True]) +def test_capacity_refusal_parks_and_next_wake_delivers_note(tmp_path, can_resume): + from outerloop.inbox import pending, wake_pending + from outerloop.operator_limits import CapacityError + from outerloop.orchestrator import RunParked + from outerloop.runstate import RunRecord + + note = "operator GPU limit for owner/repo: requested 1, running/pending 4, effective max_gpus 4" + launched = [] + meters = [] + + class Author: + supports_resume = can_resume + calls = 0 + + def run(self, brief_text, workspace, resume_session_id=None): + self.calls += 1 + if self.calls == 2: + assert resume_session_id == "s1" + assert note in brief_text + _write_syscall(workspace, {"launches": [{"name": "probe", "command": "x"}]}) + return ok_session() + + def launcher(sha, request): + launched.append(sha) + raise CapacityError(note) + + author = Author() + with pytest.raises(RunParked) as caught: + run_climb( + tmp_path, + [], + harness=author, + contract=DEEP_CONTRACT, + launcher=launcher, + launches_used=1, + sleeps_used=1, + gpu_hours_used=0.1, + on_meter=lambda *args: meters.append(args), + ) + park = caught.value + assert author.calls == len(launched) == (2 if can_resume else 1) + assert park.capacity_wait and park.phase == "author-sleep" + assert park.candidate_sha == launched[-1] + assert park.session and park.session.session_id == "s1" + assert park.syscall is not None and not park.syscall.launches + assert not park.afterany and not park.submitted + assert (park.launches_used, park.sleeps_used, park.gpu_hours_used) == (1, 1, 0.1) + assert all(meter == (1, 1, 0.1) for meter in meters) + directory = tmp_path.parent / (tmp_path.name + "-run") + refusal = pending(directory, 0)[-1] + assert refusal.source == "kernel" and refusal.kind == "note" + assert refusal.payload["quoted_text"] == note + # Deliver on wake without the note triggering a hot retry loop. + record = RunRecord("run", "owner/repo", "work", "parked", inbox_seq=refusal.seq - 1) + assert not wake_pending(directory, record) + _, harness, _ = run_climb( + tmp_path, + [], + contract=DEEP_CONTRACT, + launcher=launcher, + resume_session_id=park.session.session_id, + inbox_seq=refusal.seq - 1, + launches_used=park.launches_used, + sleeps_used=park.sleeps_used, + gpu_hours_used=park.gpu_hours_used, + ) + assert harness.calls[0][2] == "s1" + assert note in harness.calls[0][0] + assert "Capacity may now be free" in harness.calls[0][0] + + +def test_repeated_budget_refusal_keeps_existing_ending(tmp_path): + class Author: + calls = 0 + + def run(self, brief_text, workspace, resume_session_id=None): + self.calls += 1 + _write_syscall(workspace, {"launches": [{"name": "probe", "command": "x"}]}) + return ok_session() + + author = Author() + result, _, evaluator = run_climb( + tmp_path, + [], + harness=author, + contract=DEEP_CONTRACT, + launcher=lambda *args: pytest.fail("over-budget launch"), + launches_used=100, + ) + assert author.calls == 2 + assert result.outcome == "no-improvement" and result.note == "ended without a submit" + assert not evaluator.calls + + @pytest.mark.parametrize("waiting", [False, True]) def test_sibling_launch_capacity_race_keeps_evaluations_and_notifies_author(tmp_path, waiting): from outerloop.inbox import pending diff --git a/tests/test_tick.py b/tests/test_tick.py index dcdaaecb..5c9220f5 100644 --- a/tests/test_tick.py +++ b/tests/test_tick.py @@ -512,6 +512,63 @@ def test_running_experiment_is_left_alone(tmp_path: Path) -> None: assert dispatcher.dispatched == [] +@pytest.mark.parametrize("legacy", [False, True]) +def test_jobless_capacity_wait_wakes_at_deadline(tmp_path, legacy): + from outerloop.attempt import _park_run + from outerloop.harness import SessionResult + from outerloop.inbox import Message, append + from outerloop.orchestrator import RunParked + from outerloop.syscall import SyscallRequest + + record = waiting_run(tmp_path) + parked = RunParked( + phase="author-sleep", + afterany="", + base_sha="base", + seed=1, + suite_seed=1, + candidate_sha="candidate", + session=SessionResult("end_turn", False, 1.0, 1, "same-session", "", ""), + syscall=SyscallRequest(launches=()), + capacity_wait=True, + launches_used=1, + sleeps_used=2, + gpu_hours_used=0.25, + ) + append( + tmp_path / "runs" / record.run_id, + Message( + 0, + "note", + "kernel", + "", + NOW, + "capacity", + { + "text": "Capacity may now be free; retry.", + "context_only": True, + }, + ), + ) + for cycle in range(5): + now = NOW + cycle * 100 + _park_run(tmp_path, record, parked, "ref", 1, now, keep_wake_attempts=True) + saved = load_record(tmp_path, record.run_id) + assert saved.wake_attempts == 0 and saved.deadline == now + 60 + assert saved.resume_session_id == "same-session" + if legacy: + # Old records have no capacity_wait marker; the jobless deadline still works. + stage = dict(saved.stage) + del stage["capacity_wait"] + save_record(tmp_path, replace(saved, stage=stage), now) + early, _ = run_tick(tmp_path, FakeSlurm(), now=now + 59, min_tick_s=0) + assert not early.woken and not early.stuck + due, _ = run_tick(tmp_path, FakeSlurm(), now=now + 61, min_tick_s=0) + assert due.woken == ((record.run_id, "deadline"),) and not due.stuck + record = load_record(tmp_path, record.run_id) + assert record.state == PARKED and record.wake_attempts == 1 + + def test_attempts_exhausted_becomes_stuck(tmp_path: Path) -> None: waiting_run(tmp_path, wake_attempts=3) report, dispatcher = run_tick(tmp_path, FakeSlurm(states={"100": "COMPLETED"})) From e1ab3dfb1154c17e43e0ac318995339f5201a857 Mon Sep 17 00:00:00 2001 From: Mengye Ren Date: Thu, 1 Oct 2026 11:14:18 -0400 Subject: [PATCH 2/3] Capacity park only for authors that can resume --- src/outerloop/orchestrator.py | 5 +++-- tests/test_orchestrator.py | 26 +++++++++++++++++++++++++- 2 files changed, 28 insertions(+), 3 deletions(-) diff --git a/src/outerloop/orchestrator.py b/src/outerloop/orchestrator.py index 37d67828..5c9ced7c 100644 --- a/src/outerloop/orchestrator.py +++ b/src/outerloop/orchestrator.py @@ -1941,9 +1941,10 @@ def _not_run_note(request: SyscallRequest | None) -> str: note=f"Stale submit refused; no gate or sibling launches ran: {problem}", run_seed=run_seed, ) - if capacity_refused: + if capacity_refused and _can_resume(): # Bound immediate retries without turning transient capacity - # into an ending. The wake delivers this still-pending note. + # into an ending. The wake resumes this session with the note; + # an author that cannot resume keeps the old behaviour. append( inbox_dir, Message( diff --git a/tests/test_orchestrator.py b/tests/test_orchestrator.py index 76de320c..0e790950 100644 --- a/tests/test_orchestrator.py +++ b/tests/test_orchestrator.py @@ -2364,7 +2364,31 @@ def launcher(sha, request): assert len(calls) == 1 -@pytest.mark.parametrize("can_resume", [False, True]) +def test_capacity_refusal_without_resume_keeps_existing_behaviour(tmp_path): + from outerloop.operator_limits import CapacityError + + class Author: + supports_resume = False + calls = 0 + + def run(self, brief_text, workspace, resume_session_id=None): + self.calls += 1 + _write_syscall(workspace, {"launches": [{"name": "probe", "command": "x"}]}) + return ok_session() + + def launcher(sha, request): + raise CapacityError("operator GPU limit: requested 1, running/pending 4") + + author = Author() + # No park: a capacity wait would wake an author that cannot resume. + result, _, _ = run_climb( + tmp_path, [13.876, 13.876], harness=author, contract=DEEP_CONTRACT, launcher=launcher + ) + assert author.calls == 1 + assert result.outcome != "parked" # measured as-is, as before + + +@pytest.mark.parametrize("can_resume", [True]) def test_capacity_refusal_parks_and_next_wake_delivers_note(tmp_path, can_resume): from outerloop.inbox import pending, wake_pending from outerloop.operator_limits import CapacityError From 8567a2e5c871b071f64783fd6f835bb0c39d09e6 Mon Sep 17 00:00:00 2001 From: Mengye Ren Date: Thu, 1 Oct 2026 11:42:19 -0400 Subject: [PATCH 3/3] Capacity park for every author; a non-resumable author wakes into a fresh session --- CHANGELOG.md | 8 ++- docs/design/lifecycle.md | 9 ++- src/outerloop/attempt.py | 13 ++-- src/outerloop/orchestrator.py | 6 +- tests/test_attempt.py | 108 +++++++++++++++++++++++++++++++++- tests/test_orchestrator.py | 26 +------- 6 files changed, 129 insertions(+), 41 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2ad2e275..7c4ce680 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,9 +6,11 @@ Versions follow [SemVer](https://semver.org). ## [Unreleased] -- Park repeated launch capacity refusals in capacity wait instead of ending the - run. The next wake resumes the same session with the refusal and a retry note; - meters are preserved and capacity waits do not exhaust stuck retries. +- 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 + or starting a fresh session with the run context otherwise. Meters are + preserved and capacity waits do not exhaust stuck retries. Existing records need no migration. - Add operator `outerloop rebind [--root ] [--note ]` to diff --git a/docs/design/lifecycle.md b/docs/design/lifecycle.md index 9c21981b..1799cd1c 100644 --- a/docs/design/lifecycle.md +++ b/docs/design/lifecycle.md @@ -252,10 +252,13 @@ launched, measured, or charged for a refused request. Refusals repeat and never end the run, including requests with malformed content or budget problems. When a launch capacity refusal cannot be delivered in another immediate resume, the run parks uncharged in capacity wait; its next wake receives the refusal and can retry. +Authors without resume support start a fresh session with the workspace, line memory, +and inbox for orientation; resumable authors continue the same session. The session walltime and contract sleep, launch, and GPU-hour budgets bound -the loop. A session that cannot resume ends as `session-error`; endpoint loss -during refusal delivery ends as `session-outage` without sealing the rejected -tree. Abandoning a rejected tree ends without measuring, sealing, or pushing it. +the loop. Outside capacity waits, a session that cannot resume ends as +`session-error`; endpoint loss during refusal delivery ends as `session-outage` +without sealing the rejected tree. Abandoning a rejected tree ends without +measuring, sealing, or pushing it. Every refusal remains in the inbox for operators and later readers. The authoritative scope re-check in `measure_and_decide`, including wake re-entry, remains terminal. Judges still receive the run record. diff --git a/src/outerloop/attempt.py b/src/outerloop/attempt.py index 687ca263..12115239 100644 --- a/src/outerloop/attempt.py +++ b/src/outerloop/attempt.py @@ -1384,8 +1384,11 @@ def post_leg_replies(messages: tuple[dict, ...]) -> None: if record.agent_id.startswith("steward"): kwargs.setdefault("scope_validator", steward_out_of_scope) kwargs.setdefault("ruler", RULER) - if not record.resume_session_id: - # A cross-backend rebind starts fresh; workspace, line memory and inbox orient it. + resume_session_id = record.resume_session_id + if record.stage.get("capacity_wait") and not getattr(harness, "supports_resume", True): + resume_session_id = "" + if not resume_session_id: + # Fresh wakes use the workspace, line memory and inbox for orientation. kwargs.setdefault("task_hypothesis", str(record.stage.get("hypothesis") or "")) line_ref = _line_ref_for(bench, config.agent_id) kwargs.setdefault("line_ref", line_ref) @@ -1404,7 +1407,7 @@ def post_leg_replies(messages: tuple[dict, ...]) -> None: submit_preflight=lambda: _submit_preflight( ws, str(record.stage.get("base_branch") or "main"), pinned_tip ), - resume_session_id=record.resume_session_id, + resume_session_id=resume_session_id, redact_secrets=secrets, inbox_dir=directory, inbox_seq=record.inbox_seq, @@ -1567,7 +1570,7 @@ def _end(result: AttemptResult, drop_refs: list[str]) -> AttemptOutcome: drop_snapshot(ws, Snapshot(commit="", tree="", ref=ref)) return outcome - # A capacity park before the first session starts with a fresh brief. + # A capacity park starts fresh before the first session or without resume support. # Other wakes NEED the author harness and saved session. Fail as a # named ending, not a crash: the run cannot proceed and re-waking will not # help without the harness, so leaving it PARKED would just hit the stuck @@ -1580,7 +1583,7 @@ def _end(result: AttemptResult, drop_refs: list[str]) -> AttemptOutcome: and not record.author_rebind_id and not record.stage.get("capacity_wait") ) - or not getattr(harness, "supports_resume", True) + or (not getattr(harness, "supports_resume", True) and not record.stage.get("capacity_wait")) ): return _end( AttemptResult( diff --git a/src/outerloop/orchestrator.py b/src/outerloop/orchestrator.py index 5c9ced7c..54dd0f14 100644 --- a/src/outerloop/orchestrator.py +++ b/src/outerloop/orchestrator.py @@ -1941,10 +1941,10 @@ def _not_run_note(request: SyscallRequest | None) -> str: note=f"Stale submit refused; no gate or sibling launches ran: {problem}", run_seed=run_seed, ) - if capacity_refused and _can_resume(): + if capacity_refused: # Bound immediate retries without turning transient capacity - # into an ending. The wake resumes this session with the note; - # an author that cannot resume keeps the old behaviour. + # into an ending. The wake delivers the note, starting a fresh + # session when the author cannot resume. append( inbox_dir, Message( diff --git a/tests/test_attempt.py b/tests/test_attempt.py index 1474643e..cb85e108 100644 --- a/tests/test_attempt.py +++ b/tests/test_attempt.py @@ -3454,7 +3454,14 @@ def test_resume_improved_reconciles_to_an_existing_pr(tmp_path, monkeypatch) -> def _write_parked_author_sleep( - tmp_path, monkeypatch, *, raise_exc=None, run_id="tsp-9", values=None, submitted=False + tmp_path, + monkeypatch, + *, + raise_exc=None, + run_id="tsp-9", + values=None, + submitted=False, + contract_text=CONTRACT_SYSCALLS, ): """An author-sleep-parked run on disk in the REAL park state: the session's tree persisted as the author left it (uncommitted edits over base), the @@ -3470,7 +3477,7 @@ def _write_parked_author_sleep( state = tmp_path / "state" wsroot = state / "runs" / run_id / "ws" (wsroot / "src" / "pilot" / "solvers").mkdir(parents=True) - (wsroot / ".outerloop.yaml").write_text(CONTRACT_SYSCALLS) + (wsroot / ".outerloop.yaml").write_text(contract_text) (wsroot / "src" / "pilot" / "solvers" / "tsp.py").write_text("def solve(): ...\n") _git(wsroot, "init", "-q", "-b", "main") _git(wsroot, "-c", "user.name=t", "-c", "user.email=t@t", "add", "-A") @@ -3619,6 +3626,103 @@ def run(self, brief_text, workspace, resume_session_id=None): assert str(record.stage["candidate_ref"]) in refs[0] +@pytest.mark.parametrize("can_resume", [False, True]) +def test_capacity_wait_wake_delivers_refusal_with_fresh_session_when_needed( + tmp_path, monkeypatch, can_resume +) -> None: + from dataclasses import replace + + from outerloop.inbox import Message, append, pending + from outerloop.measure import MeasurementPending + from outerloop.roles import author_spec + + state, run_id, wsroot, _ = _write_parked_author_sleep( + tmp_path, + monkeypatch, + raise_exc=MeasurementPending(("701", "702")), + contract_text=CONTRACT_SYSCALLS.replace( + " direction: min\n", " direction: min\n lines: true\n" + ), + ) + (wsroot / "AGENT_MEMORY.md").write_text("Narrow sweeps were promising.\n") + record = load_record(state, run_id) + save_record( + state, + replace( + record, + agent_id="agent-01", + stage={ + **record.stage, + "capacity_wait": True, + "afterany": "", + "syscall_launches": [], + "hypothesis": "try a narrower search", + "gpu_hours_used": 0.1, + }, + ), + 1_000_001, + ) + note = "operator GPU limit: requested 1, running/pending 4" + append( + state / "runs" / run_id, + Message( + 0, + "note", + "kernel", + "", + 1_000_001, + "capacity-refusal", + { + "text": "Your syscall request was REFUSED.", + "quoted_text": note, + "context_only": True, + }, + ), + ) + calls = [] + + class RecordingHarness(ScriptedHarness): + supports_resume = can_resume + + def run(self, brief_text, workspace, resume_session_id=None): + calls.append((brief_text, workspace, resume_session_id)) + stage = load_record(state, run_id).stage + assert stage["launches_used"] == stage["sleeps_used"] == 1 + assert stage["gpu_hours_used"] == 0.1 + assert (workspace / "src/pilot/solvers/tsp.py").read_text() == ( + "def solve(): return 'wip'\n" + ) + return super().run(brief_text, workspace, resume_session_id) + + outcome = resume_run( + state, + run_id, + dispatch=_fake_dispatch(), + github=CommentingGitHub(), # type: ignore[arg-type] + bot_auth=NoAuth(), + now=1_000_100, + harness=RecordingHarness(edits={}, submit=True), + spec=author_spec(), + ) + assert outcome.outcome == "parked" + assert len(calls) == 1 + brief, workspace, session_id = calls[0] + assert workspace == wsroot + assert session_id == ("s1" if can_resume else None) + assert note in brief and "REFUSED" in brief + assert "compare against the sweep" in brief + if not can_resume: + assert "try a narrower search" in brief + assert "Narrow sweeps were promising." in brief + assert "agents/agent-01" in brief + saved = load_record(state, run_id) + assert saved.stage["phase"] == "candidate" + assert saved.stage["launches_used"] == 1 + assert saved.stage["sleeps_used"] == 2 # the fresh submit consumes a sleep + assert saved.stage["gpu_hours_used"] == 0.1 + assert not pending(state / "runs" / run_id, saved.inbox_seq) + + def test_author_sleep_wake_publishes_an_inline_improvement(tmp_path, monkeypatch) -> None: """On a synchronous backend the gate answers inline instead of parking a candidate: the woken session's improvement becomes a PR through the same diff --git a/tests/test_orchestrator.py b/tests/test_orchestrator.py index 0e790950..76de320c 100644 --- a/tests/test_orchestrator.py +++ b/tests/test_orchestrator.py @@ -2364,31 +2364,7 @@ def launcher(sha, request): assert len(calls) == 1 -def test_capacity_refusal_without_resume_keeps_existing_behaviour(tmp_path): - from outerloop.operator_limits import CapacityError - - class Author: - supports_resume = False - calls = 0 - - def run(self, brief_text, workspace, resume_session_id=None): - self.calls += 1 - _write_syscall(workspace, {"launches": [{"name": "probe", "command": "x"}]}) - return ok_session() - - def launcher(sha, request): - raise CapacityError("operator GPU limit: requested 1, running/pending 4") - - author = Author() - # No park: a capacity wait would wake an author that cannot resume. - result, _, _ = run_climb( - tmp_path, [13.876, 13.876], harness=author, contract=DEEP_CONTRACT, launcher=launcher - ) - assert author.calls == 1 - assert result.outcome != "parked" # measured as-is, as before - - -@pytest.mark.parametrize("can_resume", [True]) +@pytest.mark.parametrize("can_resume", [False, True]) def test_capacity_refusal_parks_and_next_wake_delivers_note(tmp_path, can_resume): from outerloop.inbox import pending, wake_pending from outerloop.operator_limits import CapacityError