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
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,13 @@ Versions follow [SemVer](https://semver.org).

## [Unreleased]

- 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 <run-id> [--root <root>] [--note <text>]` 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
Expand Down
11 changes: 8 additions & 3 deletions docs/design/lifecycle.md
Original file line number Diff line number Diff line change
Expand Up @@ -250,10 +250,15 @@ 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.
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.
Expand Down
13 changes: 8 additions & 5 deletions src/outerloop/attempt.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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"))
Comment thread
renmengye marked this conversation as resolved.
):
return _end(
AttemptResult(
Expand Down
40 changes: 40 additions & 0 deletions src/outerloop/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand All @@ -1939,6 +1941,44 @@ 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 the note, starting a fresh
# session when the author cannot resume.
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
Expand Down
108 changes: 106 additions & 2 deletions tests/test_attempt.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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")
Expand Down Expand Up @@ -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
Expand Down
95 changes: 95 additions & 0 deletions tests/test_orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading