diff --git a/CHANGELOG.md b/CHANGELOG.md index 3a2fc625..2250de48 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,13 @@ Versions follow [SemVer](https://semver.org). ## [Unreleased] +- Add `outerloop end [--root ] [--note ]` to request an + operator ending at the next tick, including runs waiting for an unavailable + endpoint. The tick uses the existing ending cleanup and refuses later publish. +- Upgrading: existing run records need no backfill; an absent `end-request.json` + means no request. Older kernels ignore requests and do not recognize the new + `operator` ending when writing records; keep the updated kernel for these runs. + - Author overrides accept operator-only `session_minutes` (10–240) and `session_max_turns` (10–300). Limits bind with the author selection and survive settings changes across wakes and review replies. Contracts can still lower diff --git a/docs/design/agent-protocols.md b/docs/design/agent-protocols.md index 549e6ba4..36faa4d9 100644 --- a/docs/design/agent-protocols.md +++ b/docs/design/agent-protocols.md @@ -94,7 +94,7 @@ park, terminal at the ending. Legs are turns inside it. An ending is a completion with an outcome, not a task state. A negative result is successful work; a rejected PR is a human's decision, not the agent rejecting the task. Only kernel-side failure (stuck, aborted) maps to -failed or canceled. The six endings travel as data on the final message. +failed or canceled. The endings travel as data on the final message. | Outerloop | A2A | Note | | --- | --- | --- | diff --git a/docs/design/lifecycle.md b/docs/design/lifecycle.md index 02f9a9d6..12a614d0 100644 --- a/docs/design/lifecycle.md +++ b/docs/design/lifecycle.md @@ -44,7 +44,7 @@ through `attempt.run_author_leg`, including the steward's review work. `followup.py` and its job are gone. GitHub collection positions live beside the inbox, while the run record carries the delivered sequence. Replies are posted by the author leg; only a submit invokes the gate and one publish. -The six endings and reports remain. Old state names and collection positions +The endings and reports remain. Old state names and collection positions migrate on read. Each tick logs `legacy follow-up records: N` from raw live records; operators must confirm zero across every fleet before deployment. @@ -55,14 +55,14 @@ records; operators must confirm zero across every fleet before deployment. | State | Meaning | Leaves it | | --- | --- | --- | | `running` | a session is live in a job | the session ends: it slept (park), or it stopped (end) | -| `parked` | no session; the run waits for jobs it launched, for messages, or both | a wake (the same session resumes), or a human ends the PR | +| `parked` | no session; the run waits for jobs it launched, for messages, or both | a wake (the same session resumes), or a PR or operator ending | | `ended` | terminal, with a report | never | A PR being open is a fact about a run, recorded in `pr_url`, not a state. A parked run with a PR is what `in-review` was. `implementing` is `running`; -`waiting` and `in-review` are `parked`; `concluding` is deleted. The six -endings stay as they are: merged, rejected, negative result, budget exhausted, -aborted, stuck. They are how a human reads the board, and every one still +`waiting` and `in-review` are `parked`; `concluding` is deleted. The endings +are merged, rejected, negative result, budget exhausted, aborted, stuck, and +operator. They are how a human reads the board, and every one still produces a report. ### One engine: park and wake @@ -213,14 +213,22 @@ nothing. A run that has spent everything can still reply and end. ### Endings A run ends when the author ends it, when its PR is merged or closed, -when the meter runs out, or when the kernel cannot continue: a crash, a -tampered workspace, or a kernel action that made no progress -`MAX_WAKE_ATTEMPTS` times, a failed publish retry included. A merge or close +when an operator requests it with `outerloop end `, when the meter +runs out, or when the kernel cannot continue: a crash, a tampered workspace, +or a kernel action that made no progress +`MAX_WAKE_ATTEMPTS` times, a failed publish retry included. The operator command +atomically records the request time and optional note in the run directory; it +does not end the run itself. The tick records the ending as `operator` and keeps +the note as the ending note. A merge or close ends the run at the next tick whatever it is doing: pending jobs are cancelled, a session in flight finishes its leg and its publish is refused. +An operator request waits for the run's lease instead: a session in flight +finishes its leg, and its publish is refused unless it had already started, a queued wake exits without a leg, +and the next tick that holds the lease ends the run. If that leg ends the run +itself (a negative result, budget exhausted, stuck), its own ending stands. Every ending writes the report, seals the line notebook, releases the issue -claim when no PR exists, and cancels the run's live launches. No other path -ends a run. A gate verdict never ends a run by itself, and a reviewer's +claim when no PR exists, and cancels the run's live launches. These are the only +paths that end a run. A gate verdict never ends a run by itself, and a reviewer's comment never does. ## What stays rigid diff --git a/docs/install.md b/docs/install.md index e8043389..014532eb 100644 --- a/docs/install.md +++ b/docs/install.md @@ -699,6 +699,21 @@ message delivery, GitHub polling, self-merge sweep, board, and ending records continue; existing runs keep spending, including their panels, author sessions, and the authors' own `launch` submissions. +To end one run, use `outerloop end [--root ] [--note ]`. +The root defaults to `OUTERLOOP_ROOT` in the environment or operator settings, +then `~/.outerloop`. The command atomically writes `end-request.json` in the +run directory with the requested time and note. Unknown and already ended runs +are refused; repeating a pending request preserves its original time and note. +The run ends as `operator`, with the same cleanup as a PR merge or close, at the +first active tick when no session holds it, even without a PR or an available +model endpoint. A session in flight finishes its leg and cannot start a publish (one +already under way completes, and the run ends right after the leg); a queued +wake exits without starting one. Pending jobs are cancelled. A run with an issue +waits until GitHub is reachable, so the issue is told. If a session in flight +ends the run itself first, its own ending stands. +The slot becomes free, and the next claim uses the current settings, including +`OUTERLOOP_AUTHOR_OVERRIDES`. A paused loop must resume to process the request. + **Live operator ceilings.** Create `/limits.toml` to limit this fleet while leaving target contracts under their normal review process: diff --git a/src/outerloop/attempt.py b/src/outerloop/attempt.py index 3154d6da..2cae6ed4 100644 --- a/src/outerloop/attempt.py +++ b/src/outerloop/attempt.py @@ -3753,7 +3753,7 @@ def publish( ) -> AttemptOutcome: """Publish a credited sealed tree: open a PR or fast-forward its head.""" latest = load_record(run_root, run_id) - if latest.state == ENDED: + if latest.state == ENDED or end_requested(run_root, run_id): return AttemptOutcome(run_id=run_id, outcome="publish-refused", pr_url=latest.pr_url) meter = { k: v @@ -5267,6 +5267,11 @@ def _run_id(value: str) -> str: # a wake must never crash on an unreadable/odd record — fall back to # the claude author (resume_author), same fail-safe as the sweep _wake_record = None + if end_requested(Path(args.run_root), args.resume): + # an operator ending is pending: no new leg; the tick ends the run + # once this wake's lease is free + _release_own_lease(args.run_root, args.resume) + return 0 wake_limits = bound_limits(getattr(_wake_record, "author_limits", None)) if wake_limits is not None: if args.job_minutes: @@ -5366,6 +5371,10 @@ def _run_id(value: str) -> str: else None, ) wake_secrets = tuple(k for k in (bot_auth.token(), *wake_panel_secrets, wake_api_key) if k) + if end_requested(Path(args.run_root), args.resume): + # requested during setup: still no new leg + _release_own_lease(args.run_root, args.resume) + return 0 try: resumed = resume_run( args.run_root, @@ -5575,6 +5584,39 @@ def withdraw_pr( return "" +def end_requested(run_root: Path, run_id: str) -> bool: + from outerloop.runstate import END_REQUEST_NAME + + return (run_dir_of(run_root, run_id) / END_REQUEST_NAME).is_file() + + +def end_on_request( + run_root: Path, record: RunRecord, github: GitHubClient | None, now: float +) -> str: + """End a run an operator asked to end. The caller holds the run's lease, so + no session leg (and no publish) is in flight.""" + from outerloop.runstate import OPERATOR + + record = load_record(run_root, record.run_id) + if record.state == ENDED or not end_requested(run_root, record.run_id): + return "" + if record.issue_number and github is None: + return "" # the issue comment needs GitHub; a later tick ends it + finish_run(run_root, record, OPERATOR, requested_note(run_root, record.run_id), now, github) + return OPERATOR + + +def requested_note(run_root: Path, run_id: str) -> str: + from outerloop.runstate import END_REQUEST_NAME + + try: + request = run_dir_of(run_root, run_id) / END_REQUEST_NAME + note = json.loads(request.read_text()).get("note", "") + return note if isinstance(note, str) else "" + except (OSError, ValueError, AttributeError): + return "" # the request still stands; a damaged note never blocks the ending + + def close_if_done(run_root: Path, record: RunRecord, github: GitHubClient, now: float) -> str: """Finish a withdrawal intent, then route the PR ending through the run terminal.""" from outerloop.github import GitHubError @@ -5733,6 +5775,10 @@ def _ending_comment(record: RunRecord, ending: str) -> str: """ from outerloop.steward import MAX_STEWARD_ATTEMPTS, RELEASE_MARKER + if ending == "operator": + # As for any ending without a PR, nothing else will resolve the issue. + release = "" if record.pr_url else f"{RELEASE_MARKER}\n" + return f"{release}Run `{record.run_id}` was ended by an operator." if ending == "merged": return ( f"Pull request {record.pr_url} was merged; run `{record.run_id}` is " diff --git a/src/outerloop/cli.py b/src/outerloop/cli.py index f0528317..8822e718 100644 --- a/src/outerloop/cli.py +++ b/src/outerloop/cli.py @@ -860,6 +860,32 @@ def upgrade(args: argparse.Namespace) -> int: return 0 +def end(args: argparse.Namespace) -> int: + import time + + from outerloop.runstate import request_end + + try: + values = env_file_values(keys=("OUTERLOOP_ROOT",)) + root = Path( + args.root + or os.environ.get("OUTERLOOP_ROOT") + or values.get("OUTERLOOP_ROOT") + or DEFAULT_LOCAL_ROOT + ).expanduser() + created = request_end(root, args.run_id, args.note, time.time()) + except (ValueError, OSError) as exc: + print(f"outerloop end: {exc}", file=sys.stderr) + return 2 + status = "requested" if created else "already requested" + print( + f"Run {args.run_id}: ending {status}. At the next tick, the run will end as " + "operator and its pending jobs will be cancelled; a session in flight will " + "finish its leg and its publish will be refused." + ) + return 0 + + def main(argv: list[str] | None = None) -> int: from outerloop import __version__ @@ -934,6 +960,10 @@ def main(argv: list[str] | None = None) -> int: from outerloop import init return init.main(argv[1:]) + p = sub.add_parser("end", help="request an operator ending at the next tick") + p.add_argument("run_id") + p.add_argument("--root", help="state root (defaults to OUTERLOOP_ROOT or ~/.outerloop)") + p.add_argument("--note", default="", help="reason for ending the run") p = sub.add_parser("limits", help="show live operator ceilings and fleet GPU usage") p.add_argument("--root", help="state root (defaults to OUTERLOOP_ROOT or ~/.outerloop)") p = sub.add_parser("status", help="show local runs and endpoint outages (read-only)") @@ -981,6 +1011,8 @@ def main(argv: list[str] | None = None) -> int: print(str(exc), file=sys.stderr) return 1 return 0 + if args.command == "end": + return end(args) if args.command == "permissions": return permissions(args) if args.command == "upgrade": diff --git a/src/outerloop/runstate.py b/src/outerloop/runstate.py index 3e6ee75e..7388d026 100644 --- a/src/outerloop/runstate.py +++ b/src/outerloop/runstate.py @@ -34,16 +34,18 @@ STATES = (RUNNING, PARKED, ENDED) -# The six endings ("The life of a run" — every one produces a report). +# The endings ("The life of a run" — every one produces a report). MERGED = "merged" REJECTED = "rejected" NEGATIVE_RESULT = "negative-result" BUDGET_EXHAUSTED = "budget-exhausted" ABORTED = "aborted" STUCK = "stuck" +OPERATOR = "operator" -ENDINGS = (MERGED, REJECTED, NEGATIVE_RESULT, BUDGET_EXHAUSTED, ABORTED, STUCK) +ENDINGS = (MERGED, REJECTED, NEGATIVE_RESULT, BUDGET_EXHAUSTED, ABORTED, STUCK, OPERATOR) +END_REQUEST_NAME = "end-request.json" RECORD_NAME = "state.json" LEASE_NAME = "lease.json" @@ -190,6 +192,30 @@ def run_dir(root: Path, run_id: str) -> Path: return root / "runs" / run_id +def request_end(root: Path, run_id: str, note: str, now: float) -> bool: + """Atomically retain the first operator request without changing the run.""" + if not run_id or run_id in (".", "..") or Path(run_id).name != run_id: + raise ValueError(f"unknown run id: {run_id}") + directory = run_dir(root, run_id) + if not (directory / RECORD_NAME).is_file(): + raise ValueError(f"unknown run id: {run_id}") + with (directory / ".record-lock").open("a") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + record = load_record(root, run_id) + if record.ended(): + raise ValueError(f"run {run_id} already ended ({record.ending})") + path = directory / END_REQUEST_NAME + if path.exists(): + return False + tmp = directory / f".{END_REQUEST_NAME}.{os.getpid()}.tmp" + try: + tmp.write_text(json.dumps({"requested_at": now, "note": note}) + "\n") + os.replace(tmp, path) + finally: + tmp.unlink(missing_ok=True) + return True + + def save_record(root: Path, record: RunRecord, now: float) -> None: """Serialize record writers; a persisted ending cannot be replaced.""" directory = run_dir(root, record.run_id) diff --git a/src/outerloop/tick.py b/src/outerloop/tick.py index 325f6e3e..171f3dff 100644 --- a/src/outerloop/tick.py +++ b/src/outerloop/tick.py @@ -1141,6 +1141,25 @@ def wake(record: RunRecord, reason: str, tag: str) -> None: migrate_inbox(root, record.run_id, now) finally: release_lease(root, record.run_id) + from outerloop.attempt import end_on_request, end_requested + + # An operator ending waits for the lease: a live session finishes its + # leg (its publish is refused) and a queued wake exits without one. + # A running record is a session in flight (a fresh climb holds no lease): + # it finishes its leg; _sweep_running ends it if its job died. + if ( + not dry_run + and record.state != RUNNING + and end_requested(root, record.run_id) + and acquire_lease(root, record.run_id, holder, "", now) + ): + try: + ending = end_on_request(root, record, github, now) + finally: + release_lease(root, record.run_id) + if ending: + ended.append((record.run_id, ending)) + continue merged = False blessed_before = record.auto_blessed_head try: @@ -1570,14 +1589,26 @@ def _sweep_running( # after its run is declared dead must never survive one with contextlib.suppress(Exception): compute.cancel(jid) - from outerloop.attempt import finish_run - + from outerloop.attempt import end_requested, finish_run, requested_note + from outerloop.runstate import OPERATOR + + # A pending operator request names the ending: its session died first. + requested = end_requested(root, fresh.run_id) + if requested and fresh.issue_number and github is None: + continue # the operator ending must tell its issue; a later tick ends it + if requested: + ending, note = OPERATOR, requested_note(root, fresh.run_id) + else: + ending = ABORTED + note += " — ended by the sweep (a killed climb leaves no exception to contain)" finish_run( root, fresh, - ABORTED, - f"{note} — ended by the sweep (a killed climb leaves no exception to contain)", + ending, + note, now, + # an operator ending tells the requesting issue, like every other path + github if requested else None, auth=getattr(github, "auth", None), bot_login=bot_login, ) @@ -1588,7 +1619,7 @@ def _sweep_running( try: report_path.write_text( f"# Run report — {record.target} / {record.benchmark}\n" - f"Outcome: **aborted** (climb job killed)\n" + f"Outcome: **{ending}** (climb job killed)\n" f"Note: {note}\n" ) except OSError as exc: @@ -1669,6 +1700,11 @@ def _sweep_one( return # a concurrent tick reaped it first; it owns redelivery reaped.append(record.run_id) + from outerloop.attempt import end_requested + + if end_requested(root, record.run_id): + return # an operator ending is pending: never wake the run again + from outerloop.inbox import wake_pending job_ids = _poll_targets(record) diff --git a/tests/test_attempt.py b/tests/test_attempt.py index 1a69a392..4b31d410 100644 --- a/tests/test_attempt.py +++ b/tests/test_attempt.py @@ -7423,7 +7423,7 @@ def successful_resume(self, brief, workspace, resume_session_id=None): @pytest.mark.parametrize("wake", [False, True]) -@pytest.mark.parametrize("failure", ["missing", "dead", "term", "interrupt"]) +@pytest.mark.parametrize("failure", ["missing", "dead", "term", "interrupt", "operator"]) def test_endpoint_unavailable_parks_and_recovers( tmp_path, target_repo_syscalls, monkeypatch, wake, failure ): @@ -7453,7 +7453,7 @@ def test_endpoint_unavailable_parks_and_recovers( key.write_text("secret") key.chmod(0o600) endpoint = EndpointProfile("local", "", key, "model", ("anthropic",), address) - if failure != "missing": + if failure not in ("missing", "operator"): address.write_text("http://localhost:8000/v1") connection = Mock() connection.request.side_effect = { @@ -7461,6 +7461,7 @@ def test_endpoint_unavailable_parks_and_recovers( "term": climb_mod.Terminated(), "interrupt": KeyboardInterrupt(), "missing": None, + "operator": None, }[failure] monkeypatch.setattr("outerloop.endpoints.HTTPConnection", Mock(return_value=connection)) original = ScriptedHarness.run @@ -7516,6 +7517,18 @@ def resume(): assert load_record(root, "tsp-1").wake_attempts == waiting.wake_attempts assert load_record(root, "tsp-1").stage["endpoint_wait"] == wait + if failure == "operator": + from fakes import RecordingDispatcher + from outerloop.compute import LocalCompute + from outerloop.runstate import request_end + from outerloop.tick import sweep + + request_end(root, "tsp-1", "endpoint retired", 1_000_102) + report = sweep(root, LocalCompute(root), RecordingDispatcher(), 1_000_103) + assert report.review_ended == (("tsp-1", "operator"),) + assert load_record(root, "tsp-1").ending_note == "endpoint retired" + return + address.write_text("http://localhost:8000/v1") connection.request.side_effect = None connection.getresponse.return_value.status = 200 diff --git a/tests/test_attempt_review.py b/tests/test_attempt_review.py index 2bd7dd51..1975002a 100644 --- a/tests/test_attempt_review.py +++ b/tests/test_attempt_review.py @@ -1461,7 +1461,7 @@ def run(self, *args, **kwargs): assert outage_active(root, NOW) -@pytest.mark.parametrize("ending", ["merged", "rejected"]) +@pytest.mark.parametrize("ending", ["merged", "rejected", "operator"]) @pytest.mark.parametrize("during", ["session", "post"]) def test_reply_return_preserves_concurrent_ending(review_run, ending, during): from outerloop.attempt import finish_run diff --git a/tests/test_operator_end.py b/tests/test_operator_end.py new file mode 100644 index 00000000..cbf3c232 --- /dev/null +++ b/tests/test_operator_end.py @@ -0,0 +1,291 @@ +"""Operator requests use the tick's terminal path, including live sessions.""" + +import json +from dataclasses import replace +from unittest.mock import Mock + +import pytest + +from fakes import RecordingDispatcher +from outerloop.attempt import end_on_request, publish +from outerloop.cli import main +from outerloop.endpoints import EndpointUnavailable +from outerloop.runstate import ( + END_REQUEST_NAME, + ENDED, + PARKED, + RUNNING, + RunRecord, + acquire_lease, + load_record, + release_lease, + request_end, + run_dir, + save_record, +) +from outerloop.tick import sweep + + +@pytest.fixture +def run(tmp_path, monkeypatch): + monkeypatch.setattr("outerloop.cli.env_file_values", lambda **kwargs: {}) + record = RunRecord(run_id="r1", target="owner/repo", task_title="Try a change", state=PARKED) + save_record(tmp_path, record, 1) + return record + + +def test_command_request_is_idempotent_and_atomic(tmp_path, run, monkeypatch, capsys): + from outerloop import runstate + + path = run_dir(tmp_path, run.run_id) / END_REQUEST_NAME + original = runstate.os.replace + writes = [] + + def replace_file(source, destination): + assert not path.exists() + assert source.parent == path.parent + payload = json.loads(source.read_text()) + assert payload["note"] == "endpoint retired" + assert payload["requested_at"] > 0 + writes.append(destination) + original(source, destination) + + monkeypatch.setattr(runstate.os, "replace", replace_file) + args = ["end", "r1", "--root", str(tmp_path), "--note", "endpoint retired"] + assert main(args) == 0 + before = path.read_bytes() + assert main([*args[:-1], "another reason"]) == 0 + assert path.read_bytes() == before + assert writes == [path] + assert load_record(tmp_path, "r1").state == PARKED + assert "already requested" in capsys.readouterr().out + + +@pytest.mark.parametrize("run_id", ["missing", "../r1", ".", ".."]) +def test_command_refuses_unknown_ids(tmp_path, run, run_id, capsys): + assert main(["end", run_id, "--root", str(tmp_path)]) == 2 + assert "unknown run id" in capsys.readouterr().err + assert not (run_dir(tmp_path, "r1") / END_REQUEST_NAME).exists() + + +def test_command_refuses_ended_run(tmp_path, run, capsys): + save_record(tmp_path, replace(run, state=ENDED, ending="operator"), 2) + assert main(["end", "r1", "--root", str(tmp_path)]) == 2 + assert "already ended" in capsys.readouterr().err + + +def test_failed_atomic_write_leaves_no_request(tmp_path, run, monkeypatch): + def fail(*args): + raise OSError("write failed") + + monkeypatch.setattr("outerloop.runstate.os.replace", fail) + assert main(["end", "r1", "--root", str(tmp_path)]) == 2 + assert not list(run_dir(tmp_path, "r1").glob("*end-request*")) + + +@pytest.mark.parametrize("state", [PARKED, RUNNING]) +@pytest.mark.parametrize("issue", [0, 12]) +def test_tick_operator_ending(tmp_path, run, monkeypatch, state, issue): + record = replace(run, state=state, issue_number=issue, stage={"launch_afterany": "afterany:7"}) + save_record(tmp_path, record, 2) + # An unavailable endpoint must never be consulted to end a run. + dispatcher = RecordingDispatcher() + dispatch = Mock(side_effect=EndpointUnavailable("unavailable")) + monkeypatch.setattr(dispatcher, "dispatch", dispatch) + compute = Mock() + compute.status.return_value = "PENDING" + monkeypatch.setattr("outerloop.compute.compute_from_env", lambda: compute) + github = Mock() if issue else None + if state == RUNNING: + assert acquire_lease(tmp_path, "r1", "session", "99", 2) + assert request_end(tmp_path, "r1", "endpoint retired", 3) + dry = sweep(tmp_path, compute, dispatcher, 4, github=github, dry_run=True) + assert not dry.review_ended + assert load_record(tmp_path, "r1").state == state + if state == RUNNING: + # The live session keeps its lease: the tick waits, and the session's + # publish is refused while the request is pending. + assert not sweep(tmp_path, compute, dispatcher, 5, github=github).review_ended + assert load_record(tmp_path, "r1").state == RUNNING + outcome = publish( + result=Mock(), + ws=Mock(), + workspace=tmp_path, + run_root=tmp_path, + run_dir=run_dir(tmp_path, "r1"), + run_id="r1", + record=record, + config=Mock(), + contract=Mock(), + github=Mock(), + now=8, + secrets=(), + base_branch="main", + base_sha="base", + issue_number=issue, + line_ref="", + date="", + ) + assert outcome.outcome == "publish-refused" + release_lease(tmp_path, "r1") + # The leg finishes and parks; only then does the tick end the run. + save_record(tmp_path, replace(load_record(tmp_path, "r1"), state=PARKED), 4) + report = sweep(tmp_path, compute, dispatcher, 5, github=github) + assert report.review_ended == (("r1", "operator"),) + final = load_record(tmp_path, "r1") + assert (final.state, final.ending, final.ending_note) == (ENDED, "operator", "endpoint retired") + assert "operator: endpoint retired" in (run_dir(tmp_path, "r1") / "report.md").read_text() + from outerloop.climbboard import collect_rows + + row = collect_rows(tmp_path, "owner/repo")["benchmark"][0] + assert (row.outcome, row.note) == ("operator", "endpoint retired") + dispatch.assert_not_called() + compute.cancel.assert_called_once_with("7") + if github: + from outerloop.intake import RELEASE_MARKER + + assert RELEASE_MARKER in github.comment.call_args.args[2] + github.get_pull_request.assert_not_called() + # The request stays as an audit record, but neither a later tick nor a stale + # session record can act on it twice. + assert (run_dir(tmp_path, "r1") / END_REQUEST_NAME).exists() + assert not sweep(tmp_path, compute, dispatcher, 6, github=github).review_ended + assert end_on_request(tmp_path, record, github, 7) == "" + compute.cancel.assert_called_once() + if github: + github.comment.assert_called_once() + if state == RUNNING: + save_record(tmp_path, record, 9) + assert load_record(tmp_path, "r1") == final + + +def test_legacy_record_without_request_is_unchanged(tmp_path, run): + assert end_on_request(tmp_path, run, None, 3) == "" + assert load_record(tmp_path, "r1").state == PARKED + + +@pytest.mark.parametrize("note", [None, 3, ["x"]]) +def test_non_text_note_is_recorded_empty(tmp_path, run, note): + (run_dir(tmp_path, "r1") / END_REQUEST_NAME).write_text(json.dumps({"note": note})) + assert end_on_request(tmp_path, load_record(tmp_path, "r1"), None, 10.0) == "operator" + assert load_record(tmp_path, "r1").ending_note == "" + + +def test_damaged_request_still_ends_the_run_once(tmp_path, run): + (run_dir(tmp_path, "r1") / END_REQUEST_NAME).write_text("{not json") + assert end_on_request(tmp_path, load_record(tmp_path, "r1"), None, 10.0) == "operator" + final = load_record(tmp_path, "r1") + assert (final.state, final.ending, final.ending_note) == (ENDED, "operator", "") + assert end_on_request(tmp_path, final, None, 11.0) == "" + + +def test_issue_run_waits_for_github_before_ending(tmp_path, run): + save_record(tmp_path, replace(run, issue_number=12), 2) + assert request_end(tmp_path, "r1", "", 3) + assert end_on_request(tmp_path, load_record(tmp_path, "r1"), None, 4) == "" + assert load_record(tmp_path, "r1").state == PARKED + github = Mock() + assert end_on_request(tmp_path, load_record(tmp_path, "r1"), github, 5) == "operator" + github.comment.assert_called_once() + + +def test_queued_wake_exits_without_a_leg(tmp_path, run, monkeypatch): + from outerloop import attempt + + assert request_end(tmp_path, "r1", "", 3) + image = tmp_path / "image.sif" + image.touch() + argv = ["climb", "--resume", "r1", "--run-root", str(tmp_path), "--image", str(image)] + monkeypatch.setattr("sys.argv", argv) + monkeypatch.setattr(attempt, "_lease_held_by_another_job", lambda *args: "") + released = [] + monkeypatch.setattr(attempt, "_release_own_lease", lambda *args: released.append(args)) + monkeypatch.setattr(attempt, "resume_run", Mock(side_effect=AssertionError("leg started"))) + monkeypatch.setattr(attempt, "build_harness", Mock(side_effect=AssertionError("harness built"))) + assert attempt.main() == 0 + assert released == [(str(tmp_path), "r1")] or released == [(tmp_path, "r1")] + assert load_record(tmp_path, "r1").state == PARKED + + +def test_issue_run_without_github_is_never_woken_while_it_waits(tmp_path, run, monkeypatch): + record = replace(run, issue_number=12, experiment_job_id="7", deadline=10.0) + save_record(tmp_path, record, 2) + assert request_end(tmp_path, "r1", "retired", 3) + compute = Mock() + compute.status.return_value = "COMPLETED" + monkeypatch.setattr("outerloop.compute.compute_from_env", lambda: compute) + dispatcher = RecordingDispatcher() + dispatch = Mock() + monkeypatch.setattr(dispatcher, "dispatch", dispatch) + for now in (100, 200, 300, 400): + assert not sweep(tmp_path, compute, dispatcher, now, grace_s=0).review_ended + dispatch.assert_not_called() + waiting = load_record(tmp_path, "r1") + assert (waiting.state, waiting.wake_attempts) == (PARKED, 0) + github = Mock() + report = sweep(tmp_path, compute, dispatcher, 500, grace_s=0, github=github) + assert report.review_ended == (("r1", "operator"),) + assert load_record(tmp_path, "r1").ending_note == "retired" + + +def test_dead_session_of_a_requested_run_ends_as_operator(tmp_path, run, monkeypatch): + save_record(tmp_path, replace(run, state=RUNNING, run_job_id="55", issue_number=12), 2) + assert acquire_lease(tmp_path, "r1", "session", "55", 2) + assert request_end(tmp_path, "r1", "retired", 3) + compute = Mock() + compute.status.return_value = "FAILED" + monkeypatch.setattr("outerloop.compute.compute_from_env", lambda: compute) + dispatcher = RecordingDispatcher() + sweep(tmp_path, compute, dispatcher, 10, grace_s=1) # stamps the kill + # Without GitHub the issue cannot be told, so the dead run waits. + assert "r1" not in sweep(tmp_path, compute, dispatcher, 20, grace_s=1).running_ended + assert load_record(tmp_path, "r1").state == RUNNING + github = Mock() + report = sweep(tmp_path, compute, dispatcher, 30, grace_s=1, github=github) + assert "r1" in report.running_ended + final = load_record(tmp_path, "r1") + assert (final.state, final.ending, final.ending_note) == (ENDED, "operator", "retired") + from outerloop.intake import RELEASE_MARKER + + assert RELEASE_MARKER in github.comment.call_args.args[2] + assert "operator: retired" in (run_dir(tmp_path, "r1") / "report.md").read_text() + + +def test_request_during_wake_setup_still_prevents_the_leg(tmp_path, run, monkeypatch): + from types import SimpleNamespace + + from outerloop import attempt + + image = tmp_path / "image.sif" + image.touch() + argv = ["climb", "--resume", "r1", "--run-root", str(tmp_path), "--image", str(image)] + argv += ["--panel-skip", "test"] + monkeypatch.setattr("sys.argv", argv) + monkeypatch.setenv("OUTERLOOP_AUTHOR_BACKEND", "claude") + monkeypatch.setenv("OUTERLOOP_CLAUDE_MODEL", "claude-native") + monkeypatch.setattr(attempt, "_lease_held_by_another_job", lambda *args: "") + monkeypatch.setattr(attempt, "resolve_bot_auth", lambda *a: SimpleNamespace(token=lambda: "t")) + monkeypatch.setattr(attempt, "model_key", lambda *args: "key") + + def limits(*args): + request_end(tmp_path, "r1", "late", 4) # lands during setup, after the first check + return None + + monkeypatch.setattr(attempt, "bound_limits", limits) + monkeypatch.setattr(attempt, "build_harness", lambda *a, **k: object()) + released = [] + monkeypatch.setattr(attempt, "_release_own_lease", lambda *args: released.append(args)) + monkeypatch.setattr(attempt, "resume_run", Mock(side_effect=AssertionError("leg started"))) + assert attempt.main() == 0 + assert released + + +def test_fresh_climb_without_a_lease_finishes_its_leg(tmp_path, run, monkeypatch): + save_record(tmp_path, replace(run, state=RUNNING, run_job_id="42"), 2) + assert request_end(tmp_path, "r1", "", 3) + compute = Mock() + compute.status.return_value = "RUNNING" # the climb job is alive + monkeypatch.setattr("outerloop.compute.compute_from_env", lambda: compute) + report = sweep(tmp_path, compute, RecordingDispatcher(), 10, grace_s=1) + assert not report.review_ended and "r1" not in report.running_ended + assert load_record(tmp_path, "r1").state == RUNNING