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
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/<profile>.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.

- Startup validation of `OUTERLOOP_AUTHOR_OVERRIDES` (the tick and `outerloop start`) uses the
image sessions actually run with, the default image when `OUTERLOOP_IMAGE` is unset. Before, a
codex override on a deployment without `OUTERLOOP_IMAGE` failed validation and stopped the tick.
Expand Down
37 changes: 35 additions & 2 deletions docs/endpoints.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -217,3 +218,35 @@ proxy and CLI in the test gate, set `OUTERLOOP_BRIDGE_TEST_PYTHON` to the instal
runtime's `venv/bin/python` and `OUTERLOOP_BRIDGE_TEST_CODEX` to the pinned Codex
binary before running `uv run pytest -q`. These tests never install packages or
contact a model endpoint; without these artifacts, the real-binary tests skip.

## 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/<profile>.json` journal shares the outage latch across processes
and ticks. There are no notifications or extra health probes.

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.

## 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).
19 changes: 19 additions & 0 deletions docs/install.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
11 changes: 8 additions & 3 deletions src/outerloop/attempt.py
Original file line number Diff line number Diff line change
Expand Up @@ -303,7 +303,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(
Expand Down Expand Up @@ -559,6 +561,7 @@ def _best_effort(what: str, fn: Callable[[], object], secrets: tuple[str, ...] =


STAGE_RETAINED_KEYS = (
"endpoint_wait",
LEDGER_RETRY,
"withdraw_reason",
"ledger_digits",
Expand All @@ -579,7 +582,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]
Expand Down Expand Up @@ -4783,7 +4786,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
Expand Down
17 changes: 17 additions & 0 deletions src/outerloop/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
from __future__ import annotations

import argparse
import json
import os
import shlex
import shutil
Expand Down Expand Up @@ -935,6 +936,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")
Expand All @@ -943,6 +947,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

Expand Down
132 changes: 132 additions & 0 deletions src/outerloop/endpoint_wait.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
"""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.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.
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
42 changes: 37 additions & 5 deletions src/outerloop/endpoints.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -87,19 +93,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)
Expand All @@ -124,10 +152,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()
Expand Down
Loading
Loading