From d355e116567e8566198d281d45fc2d070d3d4e5b Mon Sep 17 00:00:00 2001 From: Mengye Ren Date: Wed, 30 Sep 2026 21:49:41 -0400 Subject: [PATCH 1/3] Endpoint outages: one log line per outage, a wait marker on each run, and a read-only outerloop status A run whose on-prem server is unavailable now records when it started waiting, the fleet logs one line when an endpoint goes down and one when it recovers, and outerloop status (text or --json) lists active runs and outages without touching state. A server record may name its model and an expiry; a mismatch or an expired record counts as unavailable. --- CHANGELOG.md | 10 ++ docs/endpoints.md | 46 +++++- docs/install.md | 19 +++ src/outerloop/attempt.py | 11 +- src/outerloop/cli.py | 17 +++ src/outerloop/endpoint_wait.py | 130 +++++++++++++++++ src/outerloop/endpoints.py | 42 +++++- src/outerloop/harness.py | 8 +- src/outerloop/orchestrator.py | 4 + src/outerloop/status.py | 106 ++++++++++++++ src/outerloop/tick.py | 16 ++- tests/test_attempt.py | 17 ++- tests/test_climbboard.py | 56 ++++++++ tests/test_endpoint_wait.py | 256 +++++++++++++++++++++++++++++++++ tests/test_endpoints.py | 6 +- tests/test_status.py | 123 ++++++++++++++++ 16 files changed, 848 insertions(+), 19 deletions(-) create mode 100644 src/outerloop/endpoint_wait.py create mode 100644 src/outerloop/status.py create mode 100644 tests/test_endpoint_wait.py create mode 100644 tests/test_status.py diff --git a/CHANGELOG.md b/CHANGELOG.md index c3a3441e..f50c933a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,16 @@ Versions follow [SemVer](https://semver.org). ## [Unreleased] +- Add read-only `outerloop status` (text/`--json`) for local runs and endpoint + outages. Endpoint waits stay out of the published board/status strip and never + trigger research-log commits; log one shared outage start and recovery with + duration and run IDs. + Validate optional served model and expiry in bounded endpoint address records. +- Upgrading: no action needed; the first endpoint deferral adds + `stage.endpoint_wait` and an `endpoint-waits/.json` log latch. Missing + keys/journals are tolerated; ended runs and in-flight PRs are unchanged. + Rollback to the preceding kernel safely ignores the additive state. + - Support separate instances on one cluster account: process-only absolute `OUTERLOOP_ENV_FILE`, with the existing ownership/write-permission checks, and stable settings-path suffixes for resident and per-cadence scheduler jobs. diff --git a/docs/endpoints.md b/docs/endpoints.md index 0410f026..c6ec2974 100644 --- a/docs/endpoints.md +++ b/docs/endpoints.md @@ -24,8 +24,9 @@ OUTERLOOP_IMAGE=/opt/agent.sif The files must be readable, nonempty, and private (`chmod 600`). Paths must be absolute (`~` is expanded). Profile names start with a letter and contain only letters, digits, and underscores; references are case-insensitive and their env -keys are uppercase. Only `_URL`, `_KEY_FILE`, `_MODEL`, and `_API` are forwarded by the -profile allowlist. URLs cannot contain credentials, a query, or a fragment. +keys are uppercase. The profile allowlist forwards `_URL`, `_URL_FILE`, +`_KEY_FILE`, `_MODEL`, and `_API`. URLs cannot contain credentials, a query, or a +fragment. A profile owns its model. `OUTERLOOP_AUTHOR_MODEL` may be omitted; if set, it must match the profile's served model. Panel syntax is @@ -186,3 +187,44 @@ request checks the server with a three-second timeout and the profile key in an Authorization header. Missing files, dead servers, and interrupted checks park fresh runs and defer wakes without consuming a wake retry. There is no polling. Malformed contents are configuration errors. Profile names, served model, API capabilities and credential paths retain their existing semantics. + +## Endpoint wait visibility + +When a run defers, `stage.endpoint_wait` records the canonical (lowercase) +profile name as `endpoint` and its first unavailable time as `since` (Unix +seconds). `outerloop status` shows the endpoint and start time locally. +Endpoint infrastructure state is excluded from the GitHub board and status strip +and never triggers a research-log commit. A successful session probe clears that +run's wait. Other waiting runs retain their own start times until they resume. + +The kernel logs one outage-start line for the first waiting run and one recovery +line with the duration and all run IDs that waited. A locked, atomically written +`endpoint-waits/.json` journal shares the outage latch across processes +and ticks. There are no notifications or extra health probes. As with +ordinary logging, a process crash between journal persistence and log emission +can lose a line; the journal prevents repetition on subsequent ticks. + +JSON address records may also contain `model` and `expires_at` (finite Unix +seconds, at most the end of year 9999). When present, the model must match the +profile and expiry must be in the future. Mismatched models, expired/invalid +expiry, and records larger than 64 KiB are unavailable. Bare URLs and JSON without these optional fields retain +their existing behavior. Reads are bounded to 64 KiB plus one byte; the +existing health request retains its three-second timeout. + +Compatibility: legacy run records without `stage.endpoint_wait` and state roots +without an outage journal mean no recorded wait; no backfill is needed. Ended +runs are left untouched. The legacy author-route fixture exercises repeated +reads and retry after an interrupted write. Unknown stage and journal keys are +preserved. The preceding kernel can ignore these additive keys on rollback; +wait visibility is lost, but run and PR lifecycle semantics are unchanged. + +## Local operator status + +Run `outerloop status` or `outerloop status --root /path/to/state --json` to +read active runs and current endpoint outages without GitHub, scheduler calls, +health probes, or state writes. Each run includes its author route and override +flag, phase, recorded GPU-hours used/budget, and any endpoint wait. Text times +are UTC; JSON times are Unix seconds. Outage waiting-run IDs come from the shared +journal; after recovery, individual runs can still retain waits until they resume. +GPU budgets use the local workspace contract (including review top-ups); missing +or unreadable contracts show `unknown` (`null` in JSON). diff --git a/docs/install.md b/docs/install.md index 13b99d08..de7bf3f8 100644 --- a/docs/install.md +++ b/docs/install.md @@ -479,6 +479,25 @@ The login loop does **not auto-update**, even with `OUTERLOOP_AUTO_UPDATE=main`. To restart or upgrade: stop the process, run `git pull`, run `uv sync`, then run `outerloop start` again. Settings from `.env` are exported only at launch. +### Local status + +```bash +outerloop status +outerloop status --root /path/to/state --json +``` + +This read-only command lists non-ended runs with target, agent, state/phase, +author backend/model and override flag, recorded GPU-hours used/budget, and +endpoint waits, followed by current endpoint outages and their waiting run IDs. +It reads local files only, with no GitHub, scheduler, or health-probe calls, so it +is safe on a login node. GPU budgets come from each local workspace contract, +including review top-ups; unavailable contracts show unknown/null. +Root precedence is `--root`, process `OUTERLOOP_ROOT`, the selected operator +settings file's `OUTERLOOP_ROOT`, then `~/.outerloop`—the resolution used by +`start` when launching the tick (`tick` itself requires `--root`). +An empty or nonexistent root reports no runs/outages without creating files. +See [endpoint waits](endpoints.md#local-operator-status) for outage semantics. + ### Upgrading 1. Run `outerloop upgrade` (add `--pre` for pre-releases). diff --git a/src/outerloop/attempt.py b/src/outerloop/attempt.py index 0ffd167d..b68948a6 100644 --- a/src/outerloop/attempt.py +++ b/src/outerloop/attempt.py @@ -297,7 +297,9 @@ def _defer_endpoint(root: Path, record: RunRecord, exc: EndpointUnavailable) -> dc_replace(record, state=PARKED, wake_attempts=max(0, record.wake_attempts - 1)), time.time(), ) - log.warning("run %s: %s", record.run_id, exc) + from outerloop.endpoint_wait import unavailable + + unavailable(root, record.run_id, exc, time.time()) def resume_author( @@ -553,6 +555,7 @@ def _best_effort(what: str, fn: Callable[[], object], secrets: tuple[str, ...] = STAGE_RETAINED_KEYS = ( + "endpoint_wait", LEDGER_RETRY, "withdraw_reason", "ledger_digits", @@ -573,7 +576,7 @@ def _message_stage(run_root: Path, record: RunRecord) -> dict[str, object]: """Delivery owns these keys; captured leg records are never authoritative.""" current = load_record(run_root, record.run_id) stage = dict(record.stage) - for key in ("message_counter", "message_delivery"): + for key in ("message_counter", "message_delivery", "endpoint_wait"): stage.pop(key, None) if key in current.stage: stage[key] = current.stage[key] @@ -4704,7 +4707,9 @@ def acknowledge(seq: int) -> None: base_branch=base_branch, ) parked = p - log.warning("run %s: %s", run_id, exc) + from outerloop.endpoint_wait import unavailable + + unavailable(run_root, run_id, exc, time.time()) return AttemptOutcome(run_id=run_id, outcome="parked") except RunParked as p: # The climb dispatched its measures and hibernated. Persist the diff --git a/src/outerloop/cli.py b/src/outerloop/cli.py index ef86b452..63b7c61d 100644 --- a/src/outerloop/cli.py +++ b/src/outerloop/cli.py @@ -13,6 +13,7 @@ from __future__ import annotations import argparse +import json import os import shlex import shutil @@ -920,6 +921,9 @@ def main(argv: list[str] | None = None) -> int: return init.main(argv[1:]) 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)") + p.add_argument("--root", help="state root (defaults to OUTERLOOP_ROOT or ~/.outerloop)") + p.add_argument("--json", action="store_true", help="print structured JSON") p = sub.add_parser("permissions", help="check and update the App's required permissions") p.add_argument("--open", action="store_true", help="open the next permission settings page") p = sub.add_parser("migrate-ledger", help="seed research-log from a pinned main ledger") @@ -928,6 +932,19 @@ def main(argv: list[str] | None = None) -> int: p.add_argument("--force", action="store_true", help="replace an existing branch ledger") p.add_argument("--dry-run", action="store_true", help="print the table without writing") args = parser.parse_args(argv) + if args.command == "status": + from outerloop.status import collect_status, render_text + + values = env_file_values(keys=("OUTERLOOP_ROOT",)) if not args.root else {} + root = Path( + args.root + or os.environ.get("OUTERLOOP_ROOT") + or values.get("OUTERLOOP_ROOT") + or DEFAULT_LOCAL_ROOT + ).expanduser() + status = collect_status(root) + print(json.dumps(status, indent=2) if args.json else render_text(status)) + return 0 if args.command == "migrate-ledger": from outerloop.ledger_migrate import migrate diff --git a/src/outerloop/endpoint_wait.py b/src/outerloop/endpoint_wait.py new file mode 100644 index 00000000..a50f6868 --- /dev/null +++ b/src/outerloop/endpoint_wait.py @@ -0,0 +1,130 @@ +"""Durable per-run waits and shared endpoint outage log latches.""" + +from __future__ import annotations + +import contextlib +import fcntl +import json +import logging +import os +import time +from collections.abc import Iterator +from dataclasses import replace +from pathlib import Path +from typing import Any + +from outerloop.endpoints import EndpointProfile, EndpointUnavailable +from outerloop.runstate import _save_record, load_record, run_dir + +log = logging.getLogger(__name__) + + +class EndpointWaitReason(str): + """Transient intake deferral, not a configuration error to log every tick.""" + + +@contextlib.contextmanager +def journal(root: Path, name: str) -> Iterator[dict[str, Any]]: + directory = root / "endpoint-waits" + directory.mkdir(parents=True, exist_ok=True) + path = directory / f"{name}.json" + with (directory / f"{name}.lock").open("a") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + try: + state = json.loads(path.read_text()) + except FileNotFoundError: + state = {} + yield state + temporary = path.with_suffix(".tmp") + temporary.write_text(json.dumps(state)) + os.replace(temporary, path) + + +def waiting(root: Path, run_id: str) -> dict[str, Any]: + value = load_record(root, run_id).stage.get("endpoint_wait") + return dict(value) if isinstance(value, dict) else {} + + +def _set_wait(root: Path, run_id: str, name: str, now: float, *, clear: bool = False) -> None: + """Merge only our stage key under the existing record writer lock.""" + with (run_dir(root, run_id) / ".record-lock").open("a") as lock: + fcntl.flock(lock, fcntl.LOCK_EX) + record = load_record(root, run_id) + if record.ended(): + return + stage = dict(record.stage) + previous = stage.get("endpoint_wait") + if clear: + if not isinstance(previous, dict) or previous.get("endpoint") != name: + return + stage.pop("endpoint_wait") + elif isinstance(previous, dict) and previous.get("endpoint") == name: + return + else: + stage["endpoint_wait"] = {"endpoint": name, "since": now} + _save_record(root, replace(record, stage=stage), now) + + +def unavailable(root: Path | None, run_id: str, exc: EndpointUnavailable, now: float) -> None: + """Every deferral path uses this helper; intake/dry runs have no record to mark.""" + if root is None or not run_id or not exc.endpoint: + return + name = exc.endpoint + with journal(root, name) as state: + record = load_record(root, run_id) + if record.ended(): + return + previous = record.stage.get("endpoint_wait") + started = ( + previous["since"] + if isinstance(previous, dict) and previous.get("endpoint") == name + else now + ) + first = "since" not in state + state.setdefault("since", now) + state.setdefault("runs", {}).setdefault(run_id, started) + _set_wait(root, run_id, name, state["runs"][run_id]) + if first: + log.warning("endpoint %s unavailable; waiting runs: %s", name, run_id) + + +def recovered(root: Path, run_id: str, name: str, now: float) -> None: + """A successful session probe ends the shared outage; only this run resumes.""" + name = name.lower() + path = root / "endpoint-waits" / f"{name}.json" + if not path.exists(): + # A crash can persist the run stage before the shared journal rename. + if waiting(root, run_id).get("endpoint") == name: + _set_wait(root, run_id, name, now, clear=True) + return # Healthy/legacy runs need no journal or record writes. + with journal(root, name) as state: + since = state.pop("since", None) + runs = state.pop("runs", {}) + _set_wait(root, run_id, name, now, clear=True) + if since is not None: + log.warning( + "endpoint %s recovered after %.1fs; waited runs: %s", + name, + max(0, now - since), + ", ".join(sorted(runs)), + ) + + +def run_context(workspace: Path) -> tuple[Path | None, str]: + directory = next( + ( + p + for p in (workspace, *workspace.parents) + if p.parent.name == "runs" and (p / "state.json").is_file() + ), + None, + ) + return (directory.parent.parent, directory.name) if directory is not None else (None, "") + + +def session_url(profile: EndpointProfile, workspace: Path) -> str: + url = profile.session_url() + root, run_id = run_context(workspace) + if root is not None: + recovered(root, run_id, profile.name, time.time()) + return url diff --git a/src/outerloop/endpoints.py b/src/outerloop/endpoints.py index 834c7caf..f6c394f7 100644 --- a/src/outerloop/endpoints.py +++ b/src/outerloop/endpoints.py @@ -3,8 +3,10 @@ from __future__ import annotations import json +import math import os import re +import time from collections.abc import Mapping from dataclasses import dataclass from http.client import HTTPConnection, HTTPException, HTTPSConnection @@ -34,6 +36,10 @@ def split_endpoint(model: str) -> tuple[str, str]: class EndpointUnavailable(Exception): """A configured server address is temporarily unavailable.""" + def __init__(self, message: str, endpoint: str = "") -> None: + super().__init__(message) + self.endpoint = endpoint.lower() + def validate_url(value: str, name: str) -> str: url = urlsplit(value) @@ -72,19 +78,41 @@ def url(self) -> str: if self.url_file is None: return self.fixed_url try: - value = self.url_file.read_text().strip() + with self.url_file.open("rb") as source: + raw = source.read(65537) + if len(raw) > 65536: + raise EndpointUnavailable(f"endpoint {self.name!r}: record too large", self.name) + value = raw.decode("utf-8").strip() except OSError as exc: raise EndpointUnavailable( f"endpoint {self.name!r}: URL_FILE {self.url_file} unavailable; " - "waiting for server address" + "waiting for server address", + self.name, ) from exc if value.startswith("{"): try: - value = json.loads(value)["url"] + record = json.loads(value) + value = record["url"] except (ValueError, KeyError, TypeError) as exc: raise ValueError( f"endpoint {self.name!r}: URL_FILE must contain a URL or JSON with a url key" ) from exc + if "model" in record and record["model"] != self.model: + raise EndpointUnavailable( + f"endpoint {self.name!r}: served model mismatch", self.name + ) + if "expires_at" in record: + expiry = record["expires_at"] + if ( + isinstance(expiry, bool) + or not isinstance(expiry, (int, float)) + or not 0 < expiry <= 253402300799 + or not math.isfinite(expiry) + or expiry <= time.time() + ): + raise EndpointUnavailable( + f"endpoint {self.name!r}: expired or invalid expiry", self.name + ) if not isinstance(value, str): raise ValueError(f"endpoint {self.name!r}: URL_FILE url must be a string") return validate_url(value, self.name) @@ -109,10 +137,14 @@ def session_url(self) -> str: ) response = connection.getresponse() if response.status != 200: - raise EndpointUnavailable(f"endpoint {self.name!r}: models request failed") + raise EndpointUnavailable( + f"endpoint {self.name!r}: models request failed", self.name + ) return value except (OSError, HTTPException, Terminated, KeyboardInterrupt) as exc: - raise EndpointUnavailable(f"endpoint {self.name!r}: server unavailable") from exc + raise EndpointUnavailable( + f"endpoint {self.name!r}: server unavailable", self.name + ) from exc finally: if connection is not None: connection.close() diff --git a/src/outerloop/harness.py b/src/outerloop/harness.py index 8e732943..c9c7afc2 100644 --- a/src/outerloop/harness.py +++ b/src/outerloop/harness.py @@ -27,6 +27,7 @@ from pathlib import Path from typing import Any, Protocol +from outerloop.endpoint_wait import session_url as endpoint_session_url from outerloop.endpoints import EndpointProfile from outerloop.hermes_install import hermes_ready, hermes_runtime from outerloop.image import apptainer_from_env @@ -731,7 +732,7 @@ def run( { # Claude Code appends /v1/messages itself; profiles use the # OpenAI-style base, so drop a trailing /v1 - "ANTHROPIC_BASE_URL": self.endpoint.session_url() + "ANTHROPIC_BASE_URL": endpoint_session_url(self.endpoint, workspace) .rstrip("/") .removesuffix("/v1"), "CLAUDE_CODE_USE_VERTEX": "0", @@ -1195,7 +1196,7 @@ def run( 'model_provider = "outerloop_endpoint"\n' "[model_providers.outerloop_endpoint]\n" 'name = "Outerloop endpoint"\n' - f"base_url = {json.dumps(self.endpoint.session_url())}\n" + f"base_url = {json.dumps(endpoint_session_url(self.endpoint, workspace))}\n" 'env_key = "OUTERLOOP_SESSION_KEY"\n' 'wire_api = "responses"\n' "requires_openai_auth = false\n" @@ -1508,10 +1509,11 @@ def run( config_lines.insert(1, f" default: {json.dumps(self.model)}\n") if self.endpoint: config_lines.append(" reasoning_echo: true\n") + endpoint_url = endpoint_session_url(self.endpoint, workspace) config_lines.append( "custom_providers:\n" f" - name: {json.dumps(self.provider)}\n" - f" base_url: {json.dumps(self.endpoint.session_url())}\n" + f" base_url: {json.dumps(endpoint_url)}\n" f" key_env: {json.dumps(self.key_env)}\n" " api_mode: chat_completions\n" ) diff --git a/src/outerloop/orchestrator.py b/src/outerloop/orchestrator.py index af15a24d..03432ee8 100644 --- a/src/outerloop/orchestrator.py +++ b/src/outerloop/orchestrator.py @@ -1533,6 +1533,10 @@ def _resume(message: Message, *, allow_checkpoint: bool = True) -> AttemptResult spec, harness, prompt, workspace, resume_session_id=session.session_id ) except EndpointUnavailable as exc: + from outerloop.endpoint_wait import run_context, unavailable + + wait_root, wait_run = run_context(workspace) + unavailable(wait_root, wait_run, exc, time.time()) if not allow_checkpoint or scope_validator(list(changed_paths()), contract): # A rejected tree cannot be sealed even to preserve an outage park. return AttemptResult( diff --git a/src/outerloop/status.py b/src/outerloop/status.py new file mode 100644 index 00000000..d8b29dd9 --- /dev/null +++ b/src/outerloop/status.py @@ -0,0 +1,106 @@ +"""Read-only operator status from local records and the shared outage journal.""" + +from __future__ import annotations + +import json +import logging +from datetime import UTC, datetime +from pathlib import Path +from typing import Any + +from outerloop.contract import CONTRACT_NAME, MAX_CONTRACT_BYTES, load_contract +from outerloop.runstate import RunRecord, list_runs, run_dir + +log = logging.getLogger(__name__) + + +def _gpu_budget(root: Path, record: RunRecord) -> float | None: + """Best available local contract; an absent checkout means unknown.""" + path = run_dir(root, record.run_id) / "ws" / CONTRACT_NAME + try: + with path.open() as stream: + contract = load_contract(stream.read(MAX_CONTRACT_BYTES + 1), record.target) + budget = contract.budgets.gpu_hours_per_run + if record.stage.get("review_topup"): + budget += contract.budgets.review_topup.gpu_hours + return budget + except (OSError, ValueError): + return None + + +def collect_status(root: Path) -> dict[str, Any]: + """No probes, scheduler queries, GitHub calls, locks, or state writes.""" + runs = [] + for record in list_runs(root): + if record.ended(): + continue + stage = record.stage or {} + wait = stage.get("endpoint_wait") + runs.append( + { + "run_id": record.run_id, + "target": record.target, + "agent": record.agent_id, + "state": record.state, + "phase": stage.get("phase", ""), + "author_backend": record.author_backend or "claude", + "author_model": record.author_model, + "author_overridden": record.author_overridden, + "gpu_hours_used": stage.get("gpu_hours_used", 0.0), + "gpu_hours_budget": _gpu_budget(root, record), + "endpoint_wait": dict(wait) if isinstance(wait, dict) else None, + } + ) + outages = [] + # Writers use atomic rename, so reading requires no lock or journal creation. + for path in sorted((root / "endpoint-waits").glob("*.json")): + try: + state = json.loads(path.read_text()) + if "since" not in state: + continue + outages.append( + { + "endpoint": path.stem, + "since": state["since"], + "waiting_runs": sorted(state.get("runs", {})), + } + ) + except (OSError, ValueError, TypeError, KeyError) as exc: + log.warning("unreadable endpoint journal %s: %s", path, exc) + return {"runs": runs, "outages": outages} + + +def _time(value: Any) -> str: + try: + return datetime.fromtimestamp(float(value), tz=UTC).isoformat() + except (ValueError, TypeError, OverflowError, OSError): + return str(value) + + +def render_text(status: dict[str, Any]) -> str: + lines = [] + for run in status["runs"]: + budget = run["gpu_hours_budget"] + line = ( + f"{run['run_id']} target={run['target']} agent={run['agent']} " + f"state={run['state']} phase={run['phase'] or '-'} " + f"author={run['author_backend']}/{run['author_model'] or '(default)'} " + f"overridden={'yes' if run['author_overridden'] else 'no'} " + f"GPU-hours={run['gpu_hours_used']}/{budget if budget is not None else 'unknown'}" + ) + if wait := run["endpoint_wait"]: + line += ( + f" waiting for endpoint {wait.get('endpoint', '?')} " + f"since {_time(wait.get('since'))}" + ) + lines.append(line) + if not lines: + lines.append("No active runs.") + for outage in status["outages"]: + lines.append( + f"Endpoint outage {outage['endpoint']} since {_time(outage['since'])}; " + f"waiting runs: {', '.join(outage['waiting_runs']) or 'none'}" + ) + if not status["outages"]: + lines.append("No endpoint outages.") + return "\n".join(lines) diff --git a/src/outerloop/tick.py b/src/outerloop/tick.py index 7aa2e305..babea703 100644 --- a/src/outerloop/tick.py +++ b/src/outerloop/tick.py @@ -41,6 +41,7 @@ quote_command, ) from outerloop.disk import DEFAULT_MIN_FREE_BYTES, check_disk +from outerloop.endpoint_wait import EndpointWaitReason from outerloop.gpu_lanes import GpuLane, gpu_lanes_from_env from outerloop.harness import DEFAULT_MAX_TURNS, ClaudeModelUnset, default_claude_model, redact from outerloop.housekeeping import shed_ended_workspaces @@ -1693,7 +1694,9 @@ def _sweep_one( if profile: _ = profile.url except EndpointUnavailable as exc: - log.warning("run %s: %s", record.run_id, exc) + from outerloop.endpoint_wait import unavailable + + unavailable(None if dry_run else root, record.run_id, exc, now) deferred.append(record.run_id) return except ValueError: @@ -2476,7 +2479,12 @@ def _author_config_error(spec: ServiceSpec, agent_id: str = "agent-01") -> str: if profile: _ = profile.url # Readiness before claim: a missing address never spends an attempt. return author_config_error(backend, model, spec.image) - except (ClaudeModelUnset, ValueError, EndpointUnavailable) as exc: + except EndpointUnavailable as exc: + from outerloop.endpoint_wait import unavailable + + unavailable(None, "", exc, 0) + return EndpointWaitReason(str(exc)) + except (ClaudeModelUnset, ValueError) as exc: return str(exc) @@ -2838,6 +2846,8 @@ def service_self_initiated( return None author_error = _author_config_error(spec, slot_agent) if author_error: + if isinstance(author_error, EndpointWaitReason): + return None log.error( "climb on %s not launched: author misconfigured — %s " "(fix OUTERLOOP_AUTHOR_BACKEND/_MODEL)", @@ -3116,6 +3126,8 @@ def service_intake( return None author_error = _author_config_error(spec) if author_error: + if isinstance(author_error, EndpointWaitReason): + return None log.error( "issue #%d not claimed: author misconfigured — %s " "(fix OUTERLOOP_AUTHOR_BACKEND/_MODEL)", diff --git a/tests/test_attempt.py b/tests/test_attempt.py index 89ff676e..1023c60d 100644 --- a/tests/test_attempt.py +++ b/tests/test_attempt.py @@ -7369,6 +7369,7 @@ def test_endpoint_unavailable_parks_and_recovers( from dataclasses import replace from unittest.mock import Mock + from outerloop.endpoint_wait import session_url from outerloop.endpoints import EndpointProfile from outerloop.roles import author_spec @@ -7437,8 +7438,11 @@ def resume(): waiting = load_record(root, "tsp-1") assert waiting.state == "parked" and not waiting.ending assert waiting.wake_attempts == (2 if wake else 0) + wait = waiting.stage["endpoint_wait"] + assert isinstance(wait, dict) and wait["endpoint"] == "local" + assert isinstance(wait["since"], float) if wake: - assert waiting.stage == before.stage + assert {k: v for k, v in waiting.stage.items() if k != "endpoint_wait"} == before.stage assert waiting.resume_session_id == before.resume_session_id else: assert waiting.stage["capacity_wait"] and not waiting.resume_session_id @@ -7450,15 +7454,22 @@ def resume(): assert resume().outcome == "parked" assert load_record(root, "tsp-1").wake_attempts == waiting.wake_attempts - def recovered(self, brief, *args, **kwargs): + assert load_record(root, "tsp-1").stage["endpoint_wait"] == wait + address.write_text("http://localhost:8000/v1") + connection.request.side_effect = None + connection.getresponse.return_value.status = 200 + + def recovered(self, brief, workspace, *args, **kwargs): + session_url(endpoint, workspace) seen_briefs.append(brief) - return original(self, brief, *args, **kwargs) + return original(self, brief, workspace, *args, **kwargs) monkeypatch.setattr(ScriptedHarness, "run", recovered) assert resume().outcome != "parked" if not wake: assert "try a new move" in seen_briefs[0] assert load_record(root, "tsp-1").ending != "aborted" + assert "endpoint_wait" not in load_record(root, "tsp-1").stage @pytest.mark.parametrize("entry", ["fresh", "wake"]) diff --git a/tests/test_climbboard.py b/tests/test_climbboard.py index b3c7e6e2..c656f622 100644 --- a/tests/test_climbboard.py +++ b/tests/test_climbboard.py @@ -1436,3 +1436,59 @@ def test_hypothesis_needs_a_real_label_and_status_falls_back(tmp_path: Path) -> save_record(tmp_path, record, 2.0) (r,) = collect_status(tmp_path, "org/repo", 3.0)["runs"] assert r["hypothesis"] == "EMA helps." and r["direction"] + + +def test_endpoint_wait_is_invisible_to_published_board_and_strip(tmp_path): + from outerloop.climbboard import collect_status + from outerloop.endpoint_wait import recovered, unavailable + from outerloop.endpoints import EndpointUnavailable + + _terminal_run(tmp_path, "finished") + live = RunRecord( + "live", + "org/repo", + "Research", + "parked", + stage={"phase": "candidate", "hypothesis": "Improve the result"}, + ) + save_record(tmp_path, live, 10) + gh = _BoardGitHub() + service_climb_board(tmp_path, gh, "org/repo") + assert service_status(tmp_path, gh, "org/repo", 20) + published = dict(gh.files) + gh.puts.clear() + for now in (30, 40): + unavailable(tmp_path, "live", EndpointUnavailable("down", "local"), now) + assert service_climb_board(tmp_path, gh, "org/repo") == 0 + assert not service_status(tmp_path, gh, "org/repo", now) + assert gh.files == published + assert gh.puts == [] + recovered(tmp_path, "live", "local", 50) + assert not service_status(tmp_path, gh, "org/repo", 50) + assert service_climb_board(tmp_path, gh, "org/repo") == 0 + assert gh.files == published + assert gh.puts == [] + + # Identical records apart from the additive key produce identical wire bytes, + # including a pre-existing configuration block and an ended run's board row. + for stage in (live.stage, {**live.stage, "hermes_resume_required_chars": 10**9}): + plain = dc_replace(live, stage=stage) + waiting = dc_replace( + plain, stage={**stage, "endpoint_wait": {"endpoint": "local", "since": 1}} + ) + before = collect_status(tmp_path, "org/repo", 60, records=[plain]) + after = collect_status(tmp_path, "org/repo", 60, records=[waiting]) + assert json.dumps(before, indent=1) == json.dumps(after, indent=1) + assert "endpoint_wait" not in after["runs"][0] + + from outerloop.runstate import load_record + + ended = load_record(tmp_path, "finished") + waited = dc_replace( + ended, stage={**ended.stage, "endpoint_wait": {"endpoint": "local", "since": 1}} + ) + fresh = _BoardGitHub() + with_wait = _BoardGitHub() + service_climb_board(tmp_path, fresh, "org/repo", records=[ended]) + service_climb_board(tmp_path, with_wait, "org/repo", records=[waited]) + assert fresh.files == with_wait.files diff --git a/tests/test_endpoint_wait.py b/tests/test_endpoint_wait.py new file mode 100644 index 00000000..44522e26 --- /dev/null +++ b/tests/test_endpoint_wait.py @@ -0,0 +1,256 @@ +"""Endpoint visibility persists across ticks without changing scheduling.""" + +import json +import logging +from dataclasses import replace +from pathlib import Path +from unittest.mock import Mock + +import pytest + +from outerloop.endpoint_wait import recovered, session_url, unavailable, waiting +from outerloop.endpoints import EndpointProfile, EndpointUnavailable +from outerloop.runstate import RunRecord, load_record, save_record + + +@pytest.fixture +def profile(tmp_path): + key = tmp_path / "key" + key.write_text("endpoint-secret") + key.chmod(0o600) + return EndpointProfile("local", "", key, "open-model", ("chat",), tmp_path / "address") + + +def record(root, run_id="one"): + item = RunRecord(run_id, "owner/repo", "Check endpoint", "parked", stage={"custom": 42}) + save_record(root, item, 1) + workspace = root / "runs" / run_id / "ws" + workspace.mkdir(exist_ok=True) + return workspace + + +def down(name="local"): + return EndpointUnavailable("server unavailable", name) + + +def test_waits_status_and_shared_outage_logs(tmp_path, caplog): + caplog.set_level(logging.WARNING) + record(tmp_path) + record(tmp_path, "two") + for tick in range(100, 120): + unavailable(tmp_path, "one", down("LOCAL"), tick) + unavailable(tmp_path, "two", down(), tick + 0.5) + assert waiting(tmp_path, "one") == {"endpoint": "local", "since": 100} + assert waiting(tmp_path, "two")["since"] == 100.5 + assert load_record(tmp_path, "one").stage["custom"] == 42 + assert len(caplog.records) == 1 + recovered(tmp_path, "one", "local", 130) + assert not waiting(tmp_path, "one") + assert waiting(tmp_path, "two") # This run has not resumed yet. + recovered(tmp_path, "two", "LOCAL", 131) + recovered(tmp_path, "two", "local", 132) + assert not waiting(tmp_path, "two") + assert len(caplog.records) == 2 + assert "30.0s" in caplog.records[1].message + assert "one, two" in caplog.records[1].message + unavailable(tmp_path, "one", down(), 200) + assert len(caplog.records) == 3 + + +def test_healthy_session_clears_wait_without_extra_probe(tmp_path, profile, monkeypatch): + workspace = record(tmp_path) + profile.url_file.write_text("http://localhost:8000/v1") + connection = Mock() + connection.getresponse.return_value.status = 200 + factory = Mock(return_value=connection) + monkeypatch.setattr("outerloop.endpoints.HTTPConnection", factory) + original = (tmp_path / "runs/one/state.json").read_bytes() + assert session_url(profile, workspace) == "http://localhost:8000/v1" + assert (tmp_path / "runs/one/state.json").read_bytes() == original + assert not (tmp_path / "endpoint-waits").exists() + factory.assert_called_once_with("localhost", 8000, timeout=3) + unavailable(tmp_path, "one", down(), 100) + session_url(profile, workspace) + assert not waiting(tmp_path, "one") + assert connection.request.call_count == 2 + + +@pytest.mark.parametrize("ended", [False, True]) +def test_legacy_fixture_repeat_interruption_retry(tmp_path, monkeypatch, ended): + directory = tmp_path / "runs/legacy-author" + directory.mkdir(parents=True) + directory.joinpath("state.json").write_text( + Path("tests/fixtures/author_route_legacy.json").read_text() + ) + old = load_record(tmp_path, "legacy-author") + if ended: + save_record(tmp_path, replace(old, state="ended", ending="aborted"), 90) + for _ in range(2): + assert waiting(tmp_path, old.run_id) == {} + recovered(tmp_path, old.run_id, "local", 100) + original = __import__("os").replace + + def interrupted(src, dst): + if Path(dst).parent.name == "endpoint-waits": + raise OSError("interrupted before atomic rename") + original(src, dst) + + with monkeypatch.context() as patch: + patch.setattr("outerloop.endpoint_wait.os.replace", interrupted) + with pytest.raises(OSError): + unavailable(tmp_path, old.run_id, down(), 110) + for now in (120, 130): + unavailable(tmp_path, old.run_id, down(), now) + if ended: + assert not waiting(tmp_path, old.run_id) + else: + assert waiting(tmp_path, old.run_id)["since"] == 110 + recovered(tmp_path, old.run_id, "local", 140) + assert not waiting(tmp_path, old.run_id) + assert load_record(tmp_path, old.run_id).author_model == old.author_model + + +@pytest.mark.parametrize( + "metadata", + [ + {"model": "wrong-model"}, + {"model": None}, + {"expires_at": 0}, + {"expires_at": 100}, + {"expires_at": True}, + {"expires_at": "tomorrow"}, + {"expires_at": float("nan")}, + {"expires_at": float("inf")}, + {"expires_at": 10**400}, + ], +) +def test_invalid_metadata_is_unavailable(tmp_path, profile, monkeypatch, metadata): + profile.url_file.write_text(json.dumps({"url": "http://localhost/v1", **metadata})) + monkeypatch.setattr("outerloop.endpoints.time.time", lambda: 100) + with pytest.raises(EndpointUnavailable) as error: + profile.session_url() + assert error.value.endpoint == "local" + + +def test_valid_metadata_and_size_bound(tmp_path, profile, monkeypatch): + monkeypatch.setattr("outerloop.endpoints.time.time", lambda: 100) + profile.url_file.write_text( + json.dumps( + { + "url": "http://localhost/v1", + "model": "open-model", + "expires_at": 101, + } + ) + ) + assert profile.url == "http://localhost/v1" + profile.url_file.write_text(" " * 65537) + with pytest.raises(EndpointUnavailable, match="record too large"): + _ = profile.url + + +def test_wake_deferral_records_wait_and_refunds_retry(tmp_path, monkeypatch, profile): + from outerloop.attempt import defer_endpoint_wake + + record(tmp_path) + current = replace(load_record(tmp_path, "one"), wake_attempts=2) + save_record(tmp_path, current, 50) + monkeypatch.setattr("outerloop.attempt.resolve_endpoint", lambda *_: (profile.model, profile)) + monkeypatch.setattr("outerloop.attempt.time.time", lambda: 100) + assert defer_endpoint_wake(tmp_path, current) + assert load_record(tmp_path, "one").wake_attempts == 1 + assert waiting(tmp_path, "one") == {"endpoint": "local", "since": 100} + + +def test_dry_run_and_intake_do_not_write(tmp_path): + unavailable(None, "one", down(), 100) + unavailable(tmp_path, "", down(), 100) + assert list(tmp_path.iterdir()) == [] + + +def test_parked_sweep_uses_shared_wait_across_ticks(tmp_path, profile, monkeypatch, caplog): + from outerloop.tick import _sweep_one + + monkeypatch.setattr("outerloop.endpoints.resolve_endpoint", lambda *_: (profile.model, profile)) + for run_id in ("one", "two"): + record(tmp_path, run_id) + save_record(tmp_path, replace(load_record(tmp_path, run_id), deadline=10), 1) + wake = Mock() + for now in range(100, 105): + deferred: list[str] = [] + for run_id in ("one", "two"): + _sweep_one( + tmp_path, + Mock(), + Mock(), + now, + 60, + 60, + False, + load_record(tmp_path, run_id), + "holder", + wake, + deferred, + [], + [], + ) + assert deferred == ["one", "two"] + assert len(caplog.records) == 1 + assert waiting(tmp_path, "one")["since"] == 100 + assert waiting(tmp_path, "two")["since"] == 100 + wake.assert_not_called() + + +def test_checkpoint_preserves_wait_and_does_not_resurrect_after_recovery(tmp_path): + from outerloop.attempt import _clear_stage, _park_run + from outerloop.orchestrator import RunParked + + record(tmp_path) + stale = load_record(tmp_path, "one") + unavailable(tmp_path, "one", down(), 100) + parked = RunParked( + phase="author-sleep", + afterany="", + base_sha="base", + seed=0, + suite_seed=0, + candidate_sha="candidate", + capacity_wait=True, + ) + _park_run(tmp_path, stale, parked, "keep", 1, 101, ()) + assert waiting(tmp_path, "one")["since"] == 100 + stale = load_record(tmp_path, "one") + recovered(tmp_path, "one", "local", 110) + cleared = _clear_stage(stale, tmp_path) + save_record(tmp_path, cleared, 111) + assert not waiting(tmp_path, "one") + + +def test_intake_preflight_is_transient_without_claiming_run(tmp_path, profile, monkeypatch): + from outerloop.endpoint_wait import EndpointWaitReason + from outerloop.tick import _author_config_error + + monkeypatch.setattr("outerloop.tick._selected_author", lambda *_: ("hermes", profile.model)) + monkeypatch.setattr("outerloop.endpoints.resolve_endpoint", lambda *_: (profile.model, profile)) + for _ in range(3): + assert isinstance(_author_config_error(Mock(image="/opt/agent.sif")), EndpointWaitReason) + assert not (tmp_path / "runs").exists() + assert not (tmp_path / "endpoint-waits").exists() + + +def test_resume_after_interrupted_first_outage_write(tmp_path, monkeypatch): + record(tmp_path) + original = __import__("os").replace + + def interrupted(src, dst): + if Path(dst).parent.name == "endpoint-waits": + raise OSError("interrupted") + original(src, dst) + + with monkeypatch.context() as patch: + patch.setattr("outerloop.endpoint_wait.os.replace", interrupted) + with pytest.raises(OSError): + unavailable(tmp_path, "one", down(), 100) + assert waiting(tmp_path, "one") + recovered(tmp_path, "one", "local", 110) + assert not waiting(tmp_path, "one") diff --git a/tests/test_endpoints.py b/tests/test_endpoints.py index 269c9370..769995d6 100644 --- a/tests/test_endpoints.py +++ b/tests/test_endpoints.py @@ -668,7 +668,11 @@ def test_url_file_wake_keeps_park_and_refunds_retry(profile, tmp_path, monkeypat assert defer_endpoint_wake(tmp_path, record) waiting = load_record(tmp_path, record.run_id) assert waiting.state == "parked" and waiting.wake_attempts == 0 - assert waiting.resume_session_id == "existing" and waiting.stage == record.stage + assert waiting.resume_session_id == "existing" + wait = waiting.stage["endpoint_wait"] + assert isinstance(wait, dict) and wait["endpoint"] == "local" + assert isinstance(wait["since"], float) + assert {k: v for k, v in waiting.stage.items() if k != "endpoint_wait"} == record.stage path.write_text("http://localhost:8000/v1") assert not defer_endpoint_wake(tmp_path, waiting) diff --git a/tests/test_status.py b/tests/test_status.py new file mode 100644 index 00000000..475bb724 --- /dev/null +++ b/tests/test_status.py @@ -0,0 +1,123 @@ +import json +from dataclasses import replace +from pathlib import Path + +import pytest + +from outerloop.cli import main +from outerloop.endpoint_wait import recovered, unavailable +from outerloop.endpoints import EndpointUnavailable +from outerloop.runstate import RunRecord, save_record + + +def test_status_wait_and_outage_read_only(tmp_path, capsys, monkeypatch): + run = RunRecord( + "one", + "owner/repo", + "Research", + "parked", + author_backend="hermes", + author_model="open-model[endpoint=local]", + author_overridden=True, + stage={"phase": "candidate", "gpu_hours_used": 2.5}, + ) + save_record(tmp_path, run, 1) + save_record(tmp_path, replace(run, run_id="ended", state="ended", ending="aborted"), 1) + workspace = tmp_path / "runs/one/ws" + workspace.mkdir() + (workspace / ".outerloop.yaml").write_text( + """ +benchmarks: + - name: benchmark + command: python evaluate.py + metric: score + direction: max +budgets: + gpu_hours_per_run: 8 + runs_per_week: 10 +scope: + allowed: [src/] +roadmap: README.md +""" + ) + unavailable(tmp_path, "one", EndpointUnavailable("down", "local"), 100) + before = {p: p.read_bytes() for p in tmp_path.rglob("*") if p.is_file()} + + def forbidden(*args, **kwargs): + pytest.fail("status must not call subprocesses or the network") + + monkeypatch.setattr("subprocess.Popen", forbidden) + monkeypatch.setattr("socket.socket", forbidden) + assert main(["status", "--root", str(tmp_path)]) == 0 + text = capsys.readouterr().out + assert len(text.splitlines()) == 2 + assert "one target=owner/repo agent=agent-01 state=parked phase=candidate" in text + assert "author=hermes/open-model[endpoint=local] overridden=yes GPU-hours=2.5/8.0" in text + assert "waiting for endpoint local since 1970-01-01T00:01:40+00:00" in text + assert "Endpoint outage local since 1970-01-01T00:01:40+00:00; waiting runs: one" in text + assert main(["status", "--root", str(tmp_path), "--json"]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload == { + "runs": [ + { + "run_id": "one", + "target": "owner/repo", + "agent": "agent-01", + "state": "parked", + "phase": "candidate", + "author_backend": "hermes", + "author_model": "open-model[endpoint=local]", + "author_overridden": True, + "gpu_hours_used": 2.5, + "gpu_hours_budget": 8.0, + "endpoint_wait": {"endpoint": "local", "since": 100}, + } + ], + "outages": [{"endpoint": "local", "since": 100, "waiting_runs": ["one"]}], + } + assert before == {p: p.read_bytes() for p in tmp_path.rglob("*") if p.is_file()} + recovered(tmp_path, "one", "local", 110) + assert main(["status", "--root", str(tmp_path), "--json"]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["outages"] == [] + assert payload["runs"][0]["endpoint_wait"] is None + + +@pytest.mark.parametrize("exists", [True, False]) +def test_status_empty_root(tmp_path, capsys, exists): + root = tmp_path / "state" + if exists: + root.mkdir() + assert main(["status", "--root", str(root)]) == 0 + assert capsys.readouterr().out == "No active runs.\nNo endpoint outages.\n" + assert main(["status", "--root", str(root), "--json"]) == 0 + assert json.loads(capsys.readouterr().out) == {"runs": [], "outages": []} + assert root.exists() == exists + assert list(root.glob("*")) == [] + + +def test_status_legacy_and_root_precedence(tmp_path, monkeypatch, capsys): + root = tmp_path / "configured" + directory = root / "runs/legacy-author" + directory.mkdir(parents=True) + directory.joinpath("state.json").write_bytes( + Path("tests/fixtures/author_route_legacy.json").read_bytes() + ) + monkeypatch.delenv("OUTERLOOP_ROOT", raising=False) + monkeypatch.setattr("outerloop.cli.env_file_values", lambda **_: {"OUTERLOOP_ROOT": str(root)}) + assert main(["status", "--json"]) == 0 + run = json.loads(capsys.readouterr().out)["runs"][0] + assert run["endpoint_wait"] is None + assert run["gpu_hours_budget"] is None + assert run["author_backend"] == "codex" + assert run["author_overridden"] is False + monkeypatch.setenv("OUTERLOOP_ROOT", str(tmp_path / "empty")) + assert main(["status", "--json"]) == 0 + assert json.loads(capsys.readouterr().out)["runs"] == [] + assert main(["status", "--root", str(root), "--json"]) == 0 + assert len(json.loads(capsys.readouterr().out)["runs"]) == 1 + monkeypatch.delenv("OUTERLOOP_ROOT") + monkeypatch.setattr("outerloop.cli.env_file_values", lambda **_: {}) + monkeypatch.setattr("outerloop.cli.DEFAULT_LOCAL_ROOT", root) + assert main(["status", "--json"]) == 0 + assert len(json.loads(capsys.readouterr().out)["runs"]) == 1 From 1b695d7a295e76f1c98119b2fc1762c47b9fa51e Mon Sep 17 00:00:00 2001 From: Mengye Ren Date: Wed, 30 Sep 2026 21:58:14 -0400 Subject: [PATCH 2/3] Codex harness: keep the bridge-local base URL within the line limit --- src/outerloop/harness.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/src/outerloop/harness.py b/src/outerloop/harness.py index 1c490cdf..be18f164 100644 --- a/src/outerloop/harness.py +++ b/src/outerloop/harness.py @@ -1206,9 +1206,10 @@ def run( codex_dir.mkdir(mode=0o700, exist_ok=True) except OSError: return _error_result("workspace-error", detail="could not create codex config dir") - base_url = ( - "http://127.0.0.1:1/v1" if bridge else endpoint_session_url(self.endpoint, workspace) - ) + if bridge: + base_url = "http://127.0.0.1:1/v1" + else: + base_url = endpoint_session_url(self.endpoint, workspace) config = ( 'model_provider = "outerloop_endpoint"\n' "[model_providers.outerloop_endpoint]\n" From 3e6838474f7ca2de81c911b41d7f9e7282fd8821 Mon Sep 17 00:00:00 2001 From: Mengye Ren Date: Wed, 30 Sep 2026 22:31:13 -0400 Subject: [PATCH 3/3] Endpoint recovery takes the journal lock once an outage folder exists --- src/outerloop/endpoint_wait.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/src/outerloop/endpoint_wait.py b/src/outerloop/endpoint_wait.py index a50f6868..79df62ae 100644 --- a/src/outerloop/endpoint_wait.py +++ b/src/outerloop/endpoint_wait.py @@ -92,8 +92,10 @@ def recovered(root: Path, run_id: str, name: str, now: float) -> None: """A successful session probe ends the shared outage; only this run resumes.""" name = name.lower() path = root / "endpoint-waits" / f"{name}.json" - if not path.exists(): - # A crash can persist the run stage before the shared journal rename. + if not path.parent.exists(): + # No outage was ever journaled here. A crash can still persist the run + # stage before the journal rename; once the folder exists, recovery + # always takes the journal lock so a concurrent outage write is not lost. if waiting(root, run_id).get("endpoint") == name: _set_wait(root, run_id, name, now, clear=True) return # Healthy/legacy runs need no journal or record writes.