From 63ecc32b2f08188ade99bf8e5e1b5bf25cc931cd Mon Sep 17 00:00:00 2001 From: acamilo Date: Mon, 28 Sep 2026 21:23:33 +0000 Subject: [PATCH 1/4] loop recovery: an escalation ladder that unsticks a trap on its own Restart flysim, then reset to the current rung's milestone, then to the archive below the best rung (never lower), one step per confirmed trap that outlives the previous one. Two resets a day, three-hour restarts once they are spent; the ladder starts over at a new best rung or after six quiet hours. State and history live in the unit's StateDirectory so a reboot does not forget where the ladder stood. A router model list confirms each step; a 'not stuck' answer delays it at most three probes and no answer leaves the watchdog to decide alone. Each step is announced 60 s ahead in /run/fly/wd/recovery-notice.json for the stage's recovery splash. fly-loop-reset is the one new root surface (a sudoers line); 05-deploy now converges config/fly-sudoers so a release can add it. --- docs/loop-review.md | 4 +- infra/05-deploy.sh | 12 +- infra/bin/fly-loop-recover | 287 +++++++++++++++++++++++---- infra/bin/fly-loop-reset | 36 ++++ infra/config/fly-sudoers | 4 + infra/docs/loop-recovery.md | 103 +++++++--- infra/tests/lint.sh | 2 +- infra/tests/test_loop_recover.py | 238 ++++++++++++++++------ infra/units/fly-loop-recover.service | 6 +- 9 files changed, 563 insertions(+), 129 deletions(-) create mode 100755 infra/bin/fly-loop-reset diff --git a/docs/loop-review.md b/docs/loop-review.md index 9ea2b86..de9ddda 100644 --- a/docs/loop-review.md +++ b/docs/loop-review.md @@ -9,8 +9,8 @@ the project." 1. The release watchdog (`infra/bin/fly-watchdog` check 10) flags a suspected loop: few distinct macros, a short sequence repeating, no growth in places explored. It exports `fly_loop_suspected` and writes `/run/fly/wd/loop.json`. It never acts. The separate - `fly-loop-recover.timer` can restart only flysim after two fresh suspected probes; - see `infra/docs/loop-recovery.md`. This does not replace checkpoint-based review. + `fly-loop-recover.timer` unsticks a confirmed trap on its own: flysim restart, then milestone + resets (`infra/docs/loop-recovery.md`). That buys time; it does not replace checkpoint-based review. 2. The coordinator session (Fable) checks that marker on a schedule. On a flag it pulls the live checkpoint read-only (`pct pull`, into `.local/checkpoints/`, never committed), and spawns a review agent with the trap brief: reproduce from the checkpoint with the real diff --git a/infra/05-deploy.sh b/infra/05-deploy.sh index 5d79ce4..29e6620 100644 --- a/infra/05-deploy.sh +++ b/infra/05-deploy.sh @@ -685,10 +685,20 @@ fi # --------------------------------------------------------------------------- log "05-deploy: converging bin/ helpers to /opt/fly/bin" ct_exec "$CTID" -- mkdir -p /opt/fly/bin -for name in fly-watchdog fly-loop-recover fly-recap fly-retention fly-reset-to-milestone flypush flystage-launch flycast-launch wait-for-x wait-for-stage wait-for-health; do +for name in fly-watchdog fly-loop-recover fly-loop-reset fly-recap fly-retention fly-reset-to-milestone flypush flystage-launch flycast-launch wait-for-x wait-for-stage wait-for-health; do converge_file "$CTID" "$INFRA_DIR/bin/$name" "/opt/fly/bin/$name" 0755 root:root >/dev/null done +# The NOPASSWD surface those helpers use (config/fly-sudoers). 02-base installs it at +# provision time; converging it here too means a release that adds a helper's line +# (fly-loop-reset, v0.6.4) does not need a re-provision. Validated before install. +log "05-deploy: converging config/fly-sudoers, validated before install" +TMP_SUDOERS="/tmp/fly-sudoers.$$" +ct_push_file "$CTID" "$INFRA_DIR/config/fly-sudoers" "$TMP_SUDOERS" 0440 +ct_exec "$CTID" -- visudo -c -f "$TMP_SUDOERS" || die "config/fly-sudoers failed visudo -c, refusing to install it" +ct_exec "$CTID" -- install -o root -g root -m 0440 "$TMP_SUDOERS" /etc/sudoers.d/fly-watchdog +ct_exec "$CTID" -- rm -f "$TMP_SUDOERS" + # --------------------------------------------------------------------------- # 5. reload if anything unit-shaped changed # --------------------------------------------------------------------------- diff --git a/infra/bin/fly-loop-recover b/infra/bin/fly-loop-recover index 71de1c0..3ea369f 100755 --- a/infra/bin/fly-loop-recover +++ b/infra/bin/fly-loop-recover @@ -1,67 +1,278 @@ #!/usr/bin/env python3 -"""Conservative, out-of-process loop recovery; invoked after the watchdog probe.""" +"""Out-of-process loop recovery with an escalation ladder; invoked after the watchdog probe. + +Ladder (infra/docs/loop-recovery.md): restart flysim, then reset to the current rung's +milestone, then to the rung below. It climbs one step per confirmed trap that outlives +the previous step, and starts over once the fly reaches a new best rung or has been +quiet for QUIET seconds. A router model can delay a step, never deny it. +""" import json import os from pathlib import Path +import re import subprocess import sys import time from urllib.request import Request, urlopen REPORT = Path(os.environ.get("FLY_LOOP_REPORT", "/run/fly/wd/loop.json")) -STATE = Path(os.environ.get("FLY_LOOP_RECOVERY_STATE", "/run/fly/wd/recovery.json")) +STATE = Path(os.environ.get("FLY_LOOP_RECOVERY_STATE", "/var/lib/fly-loop-recover/state.json")) +HISTORY = Path(os.environ.get("FLY_LOOP_RECOVERY_HISTORY", "/var/lib/fly-loop-recover/history.jsonl")) +NOTICE = Path(os.environ.get("FLY_RECOVERY_NOTICE", "/run/fly/wd/recovery-notice.json")) +MILESTONES = Path(os.environ.get("FLY_STATE_DIR", "/srv/fly/state")) +STATUS_URL = os.environ.get("FLY_STATUS_URL", "http://127.0.0.1:7401/status") +RESET_BIN = os.environ.get("FLY_LOOP_RESET_BIN", "/opt/fly/bin/fly-loop-reset") INTERVAL = int(os.environ.get("WD_LOOP_INTERVAL", "300")) -COOLDOWN = 3600 +COUNTDOWN = int(os.environ.get("FLY_LOOP_COUNTDOWN", "60")) +SETTLE = 1200 # a report inside this after an action may still hold the old window +QUIET = 6 * 3600 # this long without a suspected report starts the ladder over +HOLD = 3 * 3600 # restart-only spacing while the reset budget is spent +DAY = 24 * 3600 +MAX_RESETS = 2 # milestone resets per DAY +VETO_LIMIT = 3 # model vetoes of one step before it goes ahead anyway +MODEL_BUDGET = 90 # seconds across all models in one decision +VERIFY_TIMEOUT = 240 + +# docs/design/ladder.md; mirrors RANK_LADDER in flybrain-gb's pokemon_red/mod.rs (a test pins it). +LADDER = [ + "BOOT", "BEDROOM", "DOWNSTAIRS", "PALLET TOWN", "OAK'S LAB", "GOT A STARTER", "OAK'S PARCEL", + "POKEDEX", "VIRIDIAN CITY", "VIRIDIAN FOREST", "PEWTER CITY", "BOULDER BADGE", "MT. MOON", + "CERULEAN CITY", "CASCADE BADGE", "NUGGET BRIDGE", "MET BILL", "VERMILION CITY", "HM CUT", + "THUNDER BADGE", "ROCK TUNNEL", "LAVENDER TOWN", "CELADON CITY", "SILPH SCOPE", "RAINBOW BADGE", + "POKE FLUTE", "FUCHSIA CITY", "SOUL BADGE", "SILPH CO. FREED", "MARSH BADGE", "CINNABAR ISLAND", + "VOLCANO BADGE", "EARTH BADGE", "INDIGO PLATEAU", "BEAT LORELEI", "BEAT BRUNO", "BEAT AGATHA", + "CHAMPION", +] + +PROMPT = ( + "You check a Game Boy Pokemon Red run played by a simulated fly brain through macros " + "(GO OBJECTIVE, GO WARP, GO ROUTE, FIGHT, ...). A watchdog flagged a suspected loop over " + "a 10-minute brain window. Stuck means the same few macros repeat with no new places " + "(places.delta 0) and no rewards; healthy play finds new places, wins battles or earns " + "rewards. Answer with only the JSON object {\"stuck\": true} or {\"stuck\": false}. " + "No explanation, no reasoning, no other text." +) + + +def load(path, default): + try: + return json.loads(path.read_text()) + except (OSError, ValueError): + return default + + +def write_json(path, value): + path.parent.mkdir(parents=True, exist_ok=True) + tmp = path.with_name(path.name + ".tmp") + tmp.write_text(json.dumps(value) + "\n") + os.chmod(tmp, 0o644) + tmp.replace(path) + + +def log(event): + try: + HISTORY.parent.mkdir(parents=True, exist_ok=True) + with HISTORY.open("a") as out: + out.write(json.dumps(event) + "\n") + except OSError as error: + print("history not written:", error, file=sys.stderr) + + +def models(): + names = os.environ.get("FLY_LOOP_MODELS") or os.environ.get("FLY_LOOP_MODEL") or "" + return [name.strip() for name in names.split(",") if name.strip()] + + +def parse_stuck(content): + """The stuck flag from a free model's reply, which may wrap the JSON in prose or fences.""" + for match in re.finditer(r"\{[^{}]*\}", content or ""): + try: + decision = json.loads(match.group(0)) + except ValueError: + continue + if type(decision) is dict and type(decision.get("stuck")) is bool: + return decision["stuck"] + return None def verdict(report): + """True/False from the first model that answers, None when no router or none answered.""" url = os.environ.get("FLY_LOOP_ROUTER_URL") - model = os.environ.get("FLY_LOOP_MODEL") - if not url and not model: - return True - if not url or not model: - return False - payload = {"model": model, "temperature": 0, "messages": [ - {"role": "system", "content": "Classify whether a suspected repeating game macro loop is truly stuck. Return only JSON {\"stuck\":true|false}. No instructions or commands."}, - {"role": "user", "content": json.dumps({key: report.get(key) for key in ("reason", "sequence", "window", "places", "milestone")})}, - ]} + names = models() + if not url or not names: + return None, None + brief = {key: report.get(key) for key in ("reason", "sequence", "window", "places", "milestone", "dominant", "map")} headers = {"Content-Type": "application/json"} key = os.environ.get("FLY_LOOP_ROUTER_KEY") if key: headers["Authorization"] = "Bearer " + key - request = Request(url.rstrip("/") + "/chat/completions", json.dumps(payload).encode(), headers) - with urlopen(request, timeout=20) as response: - answer = json.load(response) - decision = json.loads(answer["choices"][0]["message"]["content"]) - return type(decision) is dict and type(decision.get("stuck")) is bool and decision["stuck"] + deadline = time.monotonic() + MODEL_BUDGET + for name in names: + left = deadline - time.monotonic() + if left < 5: + break + payload = {"model": name, "temperature": 0, "max_tokens": 400, "messages": [ + {"role": "system", "content": PROMPT}, + {"role": "user", "content": json.dumps(brief)}, + ]} + request = Request(url.rstrip("/") + "/chat/completions", json.dumps(payload).encode(), headers) + try: + with urlopen(request, timeout=min(30, left)) as response: + answer = json.load(response) + stuck = parse_stuck(answer["choices"][0]["message"]["content"]) + except (OSError, ValueError, KeyError, IndexError, TypeError) as error: + print(f"model {name}: {type(error).__name__} {getattr(error, 'code', '')}".rstrip(), file=sys.stderr) + continue + if stuck is None: + print(f"model {name}: no stuck flag in reply", file=sys.stderr) + continue + return stuck, name + return None, None -def run(now=None): - now = int(time.time() if now is None else now) - report = json.loads(REPORT.read_text()) - previous = json.loads(STATE.read_text()) if STATE.exists() else {} +def rungs(): + found = [] + for path in MILESTONES.glob("milestone-*.checkpoint"): + number = path.stem.removeprefix("milestone-") + if number.isdigit(): + found.append(int(number)) + return sorted(found) + + +def plan(level, rank, best, resets_left): + """(action, target rung) for a ladder level; a spent budget or no archive falls back to restart. + + Level 1 resets to the current rung's archive, level 2 and beyond to the archive below the + best rung (never lower), so repeated days of a trap cannot walk the run down the ladder. + """ + if level == 0 or not resets_left or rank is None: + return "restart", None + ceiling = rank if level == 1 else min(rank, max(best, rank) - 1) + below = [rung for rung in rungs() if rung <= ceiling] + if below: + return "reset", below[-1] + return "restart", None + + +def status(): + try: + with urlopen(STATUS_URL, timeout=5) as response: + return json.load(response) + except (OSError, ValueError): + return None + + +def verify(target, deadline): + while time.monotonic() < deadline: + now = status() + if now and now.get("status") == "running": + rank = (now.get("milestone") or {}).get("rank") + if target is None or rank == target: + return True + time.sleep(5) + return False + + +def act(action, target): + if action == "restart": + command = ["sudo", "-n", "systemctl", "restart", "flysim.service"] + else: + command = ["sudo", "-n", RESET_BIN, str(target)] + done = subprocess.run(command, check=False, timeout=VERIFY_TIMEOUT) + if done.returncode != 0: + return False + return verify(target, time.monotonic() + VERIFY_TIMEOUT) + + +def notice(base, phase, now): + try: + write_json(NOTICE, dict(base, phase=phase, updatedAt=int(now))) + except OSError as error: + print("notice not written:", error, file=sys.stderr) + + +def run(now=None, sleep=time.sleep, clock=None): + clock = clock or (time.time if now is None else (lambda: now)) + now = int(clock()) + report = load(REPORT, None) + if type(report) is not dict: + return "no watchdog report" + state = load(STATE, {}) + if type(state) is not dict: + state = {} + rank = (report.get("milestone") or {}).get("rank") + level = state.get("level", 0) + + if type(rank) is int and rank > state.get("bestRank", 0): + if level: + log({"at": now, "event": "ladder-reset", "why": "new best rung", "rank": rank}) + state.update(bestRank=rank, level=0) + level = 0 + last = state.get("lastSuspectedAt") + if level and last and now - last > QUIET: + log({"at": now, "event": "ladder-reset", "why": "quiet", "rank": rank}) + state["level"] = level = 0 + state["resets"] = [at for at in state.get("resets", []) if now - at < DAY] + at = report.get("at") if (report.get("suspected") != 1 or type(at) is not int or not 0 <= now - at <= 600 or report.get("action") != "none"): - STATE.write_text(json.dumps({key: previous[key] for key in ("restarted_at",) if key in previous}) + "\n") + state.pop("observedAt", None) + state["vetoes"] = 0 + write_json(STATE, state) return "not a fresh suspected loop" - if "restarted_at" in previous and now - previous["restarted_at"] < COOLDOWN: - return "recovery cooldown" - if previous.get("observed_at") == at: - return "waiting for next probe" - if not previous.get("observed_at") or at - previous["observed_at"] < INTERVAL - 30: - STATE.write_text(json.dumps({"observed_at": at}) + "\n") + + state["lastSuspectedAt"] = now + acted = state.get("actedAt", 0) + if at < acted + SETTLE: + write_json(STATE, state) + return f"settling after {state.get('lastAction')}" + observed = state.get("observedAt") + if not observed or observed < acted + SETTLE: + state["observedAt"] = at + write_json(STATE, state) return "waiting for second probe" - try: - approved = verdict(report) - except (OSError, ValueError, KeyError, IndexError) as error: - print("router unavailable or invalid:", type(error).__name__, file=sys.stderr) - return "router unavailable" - if not approved: - return "router did not confirm" - subprocess.run(["sudo", "-n", "systemctl", "restart", "flysim.service"], check=True) - STATE.write_text(json.dumps({"observed_at": at, "restarted_at": now}) + "\n") - return "flysim restarted; checkpoint restore keeps the rung" + if at - observed < INTERVAL - 30: + write_json(STATE, state) + return "waiting for next probe" + + resets_left = MAX_RESETS - len(state["resets"]) + action, target = plan(level, rank, state.get("bestRank", 0), resets_left) + if action == "restart" and level and now - acted < HOLD: + write_json(STATE, state) + return "holding: reset budget spent, next restart after the hold" + + stuck, model = verdict(report) + if stuck is False and state.get("vetoes", 0) < VETO_LIMIT: + state["vetoes"] = state.get("vetoes", 0) + 1 + write_json(STATE, state) + log({"at": now, "event": "veto", "model": model, "vetoes": state["vetoes"], "rank": rank}) + return f"model {model} says not stuck ({state['vetoes']}/{VETO_LIMIT})" + why = {True: f"confirmed by {model}", False: "veto limit reached", None: "no model answer, watchdog alone"}[stuck] + + base = {"v": 1, "id": f"{now}-{action}", "action": action, + "fromRung": rank, "fromLabel": (report.get("milestone") or {}).get("label"), + "reason": report.get("reason"), "loop": (report.get("sequence") or [])[:4], + "stuckSeconds": max(0, now - observed) + 600, + "announcedAt": now, "executeAt": now + COUNTDOWN} + if action == "reset": + base.update(toRung=target, toLabel=LADDER[target] if 0 <= target < len(LADDER) else None) + notice(base, "countdown", now) + sleep(COUNTDOWN) + notice(base, "acting", clock()) + ok = act(action, target) + finished = int(clock()) + notice(base, "done" if ok else "failed", finished) + + state.update(level=level + 1, actedAt=finished, lastAction=action, vetoes=0) + state.pop("observedAt", None) + if action == "reset": + state["resets"].append(finished) + write_json(STATE, state) + log({"at": finished, "event": action, "target": target, "ok": ok, "why": why, "rank": rank, + "level": level, "sequence": report.get("sequence"), "map": report.get("map")}) + outcome = "flysim restarted" if action == "restart" else f"reset to rung {target}" + return f"{outcome} ({why})" if ok else f"{action} FAILED ({why})" if __name__ == "__main__": diff --git a/infra/bin/fly-loop-reset b/infra/bin/fly-loop-reset new file mode 100755 index 0000000..de16570 --- /dev/null +++ b/infra/bin/fly-loop-reset @@ -0,0 +1,36 @@ +#!/usr/bin/env bash +# infra/bin/fly-loop-reset — the milestone step of fly-loop-recover's ladder, as root. +# +# fly-loop-recover runs as User=fly and reaches this through one NOPASSWD sudoers line +# (config/fly-sudoers). It does what infra/docs/runbook.md "Restart the run from a rung" +# does by hand for a release whose adapter already wrote the archive: stop flysim, promote +# milestone-.checkpoint with fly-reset-to-milestone, start flysim. flysim is started +# again whether or not the reset worked, so the stream never stays down on a failure; +# fly-reset-to-milestone copies both stores aside before it rewrites anything. +# +# Usage: fly-loop-reset +set -euo pipefail + +: "${FLY_STATE_DIR:=/srv/fly/state}" +: "${FLY_SERVICE:=flysim.service}" +: "${FLY_RESET_BIN:=/opt/fly/bin/fly-reset-to-milestone}" + +log() { echo "fly-loop-reset: $*" >&2; } + +RANK="${1:-}" +if [ "$#" -ne 1 ] || ! [[ "$RANK" =~ ^[0-9]{1,2}$ ]]; then + log "usage: fly-loop-reset " + exit 2 +fi +if [ ! -f "${FLY_STATE_DIR}/milestone-${RANK}.checkpoint" ]; then + log "no milestone archive for rung ${RANK}; nothing touched" + exit 1 +fi + +log "stopping ${FLY_SERVICE} to reset to rung ${RANK}" +systemctl stop "$FLY_SERVICE" +status=0 +"$FLY_RESET_BIN" "$RANK" || status=$? +[ "$status" -eq 0 ] || log "fly-reset-to-milestone ${RANK} failed (exit ${status}); starting ${FLY_SERVICE} on the state it left" +systemctl start "$FLY_SERVICE" +exit "$status" diff --git a/infra/config/fly-sudoers b/infra/config/fly-sudoers index f22752d..8963a18 100644 --- a/infra/config/fly-sudoers +++ b/infra/config/fly-sudoers @@ -16,3 +16,7 @@ fly ALL=(root) NOPASSWD: /usr/bin/systemctl restart flycast.service fly ALL=(root) NOPASSWD: /usr/bin/systemctl restart mediamtx.service fly ALL=(root) NOPASSWD: /usr/bin/systemctl restart flypush.service fly ALL=(root) NOPASSWD: /usr/sbin/reboot +# fly-loop-recover.service (also User=fly): its ladder restarts flysim (granted above) and +# resets to a milestone archive through this one root wrapper, which validates the rung, +# stops/resets/starts flysim and nothing else (infra/docs/loop-recovery.md). +fly ALL=(root) NOPASSWD: /opt/fly/bin/fly-loop-reset diff --git a/infra/docs/loop-recovery.md b/infra/docs/loop-recovery.md index 4982482..8f3df98 100644 --- a/infra/docs/loop-recovery.md +++ b/infra/docs/loop-recovery.md @@ -1,29 +1,82 @@ # Automated loop recovery -`fly-watchdog` remains report-only. The separate `fly-loop-recover.timer` checks its -`/run/fly/wd/loop.json` every five minutes. Two fresh suspected reports from -separate watchdog probes, at least ~5 minutes apart, cause **only** a -`flysim.service` restart. The ordinary checkpoint restore keeps the current -rung and learned brain state, but clears session macro ledgers. This does not -press buttons or promote a milestone archive. A successful restart imposes a -one-hour cooldown, including across intervening clear reports. The timer's -journal records every decision; inspect it with -`journalctl -u fly-loop-recover.service`. Disable the timer to stop automatic -recovery: `systemctl disable --now fly-loop-recover.timer`. +`fly-watchdog` remains report-only. The separate `fly-loop-recover.timer` reads its +`/run/fly/wd/loop.json` every five minutes and unsticks a confirmed trap on its own, climbing +a ladder one step per trap that outlives the previous step: -By default confirmation is deterministic, using watchdog's existing signal. -An optional OpenAI-compatible local router can veto a recovery: provision a -root-managed `/etc/fly/loop-recovery.env` readable by the `fly` account, -containing `FLY_LOOP_ROUTER_URL` (base URL ending in `/v1`) and -`FLY_LOOP_MODEL` (an available free-tier model). A private router can additionally -use `FLY_LOOP_ROUTER_KEY`; provision it outside this public checkout and limit -file permissions to `0640 root:fly`. Never put credentials in the unit or the -repository. When either router setting is present, both must be set; malformed or -unavailable responses prevent recovery. The model may only return `{"stuck": true|false}` and cannot -choose commands, buttons, or checkpoint paths. Confirm the model actually -exists and is reachable from the release container before configuring it. +| level | step | cost | +|---|---|---| +| 0 | `systemctl restart flysim.service` | none: the restore keeps the rung and learned state, clears session macro ledgers | +| 1 | reset to the current rung's milestone archive | progress since the rung was first reached | +| 2+ | reset to the archive below the best rung (never lower) | one visible rung | -This is an unstick mechanism, not a macro bug fix. A recurring trap still needs -the checkpoint-based loop review in `docs/loop-review.md`. The live recovery -must be recorded in the host claim log by the operator reviewing the unit -journal; the unit itself has no host claim-log access. +- **Confirmed** means two fresh suspected watchdog reports (at most 10 minutes old, `action: none`) + from separate probes about 5 minutes apart. A clear report in between starts the count again. +- **Outlives** means the reports that confirm it again are at least 20 minutes after the step, so + their 10-brain-minute window lies after it. The helper waits that long after every step. +- **Budget:** at most two milestone resets per 24 hours. With the budget spent the step is a + restart, at most every three hours, until a reset is free again. +- **Starting over:** the ladder returns to level 0 when the fly reaches a new best rung, or after + six hours with no suspected report. + +A milestone step runs `/opt/fly/bin/fly-loop-reset ` through sudo (the one line in +`config/fly-sudoers`): stop flysim, `fly-reset-to-milestone`, start flysim. flysim is started again +even when the reset fails. The reset copies both stores to `/srv/fly/state.reset-` first, as +in the runbook, and needs no deploy when the running release wrote the archive or +`FLY_ACCEPT_ADAPTERS` already names its adapter. After each step the helper waits up to four +minutes for `/status` to report `running` (and the target rank for a reset); otherwise the step is +recorded as failed and the ladder still climbs. + +## On stream + +Every step is announced 60 seconds ahead (`FLY_LOOP_COUNTDOWN`) in +`/run/fly/wd/recovery-notice.json` (`FLY_RECOVERY_NOTICE`), which the stage's recovery splash +reads. It holds the contract below, written atomically, mode 0644: + +```json +{"v": 1, "id": "1790629095-reset", "phase": "countdown|acting|done|failed", + "action": "restart|reset", "fromRung": 12, "fromLabel": "MT. MOON", + "toRung": 11, "toLabel": "BOULDER BADGE", "reason": "unrewarded", + "loop": ["GO OBJECTIVE", "GO WARP"], "stuckSeconds": 900, + "announcedAt": 1790629095, "executeAt": 1790629155, "updatedAt": 1790629160} +``` + +`toRung`/`toLabel` are present for a reset only. Consumers ignore a notice whose `updatedAt` is +more than 15 minutes old. + +## Model confirmation + +An OpenAI-compatible router can be asked before each step. The model may only answer +`{"stuck": true|false}` (the reply may wrap it in prose; the first such object counts). It cannot +choose commands, rungs or buttons. `false` delays the step by one probe, at most three times in a +row; then the step goes ahead. If no model answers (unset, rate-limited, down, malformed), the +watchdog's confirmation stands alone. A model can delay recovery by about 15 minutes, never deny +it. + +Provision a root-managed `/etc/fly/loop-recovery.env`, mode `0640 root:fly`, outside this public +checkout: + +``` +FLY_LOOP_ROUTER_URL= +FLY_LOOP_ROUTER_KEY= +FLY_LOOP_MODELS=,,... +``` + +Models are tried in order within a 90-second budget. Prefer fast free-tier chat models, and put +providers with generous free limits first: free OpenRouter models share a small daily quota and +the router's circuit breaker can close the whole provider for a while. Check each model against +a stuck and a healthy report before listing it. `FLY_LOOP_MODEL` (one model) is still read when +`FLY_LOOP_MODELS` is unset. + +## Operating it + +- Decisions: `journalctl -u fly-loop-recover.service`. Every step, veto and ladder restart is also + appended to `/var/lib/fly-loop-recover/history.jsonl`. +- Ladder state: `/var/lib/fly-loop-recover/state.json` (the unit's `StateDirectory`, so it + survives a reboot). Deleting it starts the ladder over. +- Stop automatic recovery: `systemctl disable --now fly-loop-recover.timer`. +- Record automatic steps you find in the journal in the host claim log when you next claim the + container; the unit has no access to that log. + +This is an unstick mechanism, not a macro bug fix. A recurring trap still needs the +checkpoint-based loop review in `docs/loop-review.md`. diff --git a/infra/tests/lint.sh b/infra/tests/lint.sh index 4da888a..82e9897 100755 --- a/infra/tests/lint.sh +++ b/infra/tests/lint.sh @@ -1485,7 +1485,7 @@ fi echo "--- loop recovery tests ---" if python3 -m unittest discover -s "$REPO_ROOT/infra/tests" -p 'test_loop_recover.py' >/dev/null 2>&1; then - pass "loop recovery: fresh probes, cooldown, and router refusal" + pass "loop recovery: ladder, budget, reboot-safe state, model delay and fallback, splash notice" else fail "loop recovery tests failed; run python3 -m unittest discover -s infra/tests -p test_loop_recover.py -v" fi diff --git a/infra/tests/test_loop_recover.py b/infra/tests/test_loop_recover.py index 7d0c481..881971c 100644 --- a/infra/tests/test_loop_recover.py +++ b/infra/tests/test_loop_recover.py @@ -2,89 +2,205 @@ import importlib.machinery import importlib.util import json from pathlib import Path +import re import tempfile import unittest from unittest.mock import patch -script = Path(__file__).resolve().parents[1] / "bin/fly-loop-recover" +repo = Path(__file__).resolve().parents[2] +script = repo / "infra/bin/fly-loop-recover" spec = importlib.util.spec_from_loader("recover", importlib.machinery.SourceFileLoader("recover", str(script))) recover = importlib.util.module_from_spec(spec) spec.loader.exec_module(recover) +T0 = 1_790_000_000 + class RecoveryTests(unittest.TestCase): def setUp(self): self.temp = tempfile.TemporaryDirectory() self.addCleanup(self.temp.cleanup) root = Path(self.temp.name) - recover.REPORT = root / "loop.json" - recover.STATE = root / "recovery.json" - self.report = {"suspected": 1, "at": 1000, "action": "none", "reason": "unrewarded"} - self.save() + self.root = root + for name, value in (("REPORT", "loop.json"), ("STATE", "var/state.json"), + ("HISTORY", "var/history.jsonl"), ("NOTICE", "run/notice.json")): + patcher = patch.object(recover, name, root / value) + patcher.start() + self.addCleanup(patcher.stop) + milestones = root / "state" + milestones.mkdir() + for rung in (1, 9, 10, 11, 12): + (milestones / f"milestone-{rung}.checkpoint").write_text("x") + patcher = patch.object(recover, "MILESTONES", milestones) + patcher.start() + self.addCleanup(patcher.stop) + self.env = patch.dict("os.environ", {}, clear=False) + self.env.start() + self.addCleanup(self.env.stop) + for key in ("FLY_LOOP_ROUTER_URL", "FLY_LOOP_MODELS", "FLY_LOOP_MODEL", "FLY_LOOP_ROUTER_KEY"): + recover.os.environ.pop(key, None) + self.acts = [] + self.act = patch.object(recover, "act", side_effect=lambda action, target: self.acts.append((action, target)) or True) + self.act.start() + self.addCleanup(self.act.stop) + self.report(suspected=1, at=T0) - def save(self): - recover.REPORT.write_text(json.dumps(self.report)) + def report(self, **fields): + base = {"suspected": 1, "action": "none", "reason": "unrewarded", "sequence": ["GO WARP"], + "milestone": {"rank": 12, "label": "MT. MOON"}, "map": 61} + base.update(fields) + if "rank" in fields: + base["milestone"] = {"rank": base.pop("rank"), "label": "X"} + recover.REPORT.write_text(json.dumps(base)) - def test_two_probes_and_cooldown(self): - with patch.object(recover, "verdict", return_value=True), patch.object(recover.subprocess, "run") as restart: - self.assertIn("second probe", recover.run(1001)) - self.assertIn("next probe", recover.run(1002)) - self.report["at"] = 1300 - self.save() - self.assertIn("restarted", recover.run(1301)) - restart.assert_called_once_with(["sudo", "-n", "systemctl", "restart", "flysim.service"], check=True) - self.assertIn("cooldown", recover.run(1302)) + def tick(self, at, **fields): + """One watchdog probe at `at` and the timer running right after it.""" + self.report(at=at, **fields) + return recover.run(at + 5, sleep=lambda seconds: None) - def test_stale_and_clear_never_restart(self): - with patch.object(recover.subprocess, "run") as restart: - self.assertIn("not a fresh", recover.run(1700)) - self.report.update(at=1700, suspected=0) - self.save() - self.assertIn("not a fresh", recover.run(1700)) - restart.assert_not_called() + def confirm(self, start, **fields): + """Two probes ~5 min apart after `start`; returns the second decision.""" + self.assertIn("second probe", self.tick(start, **fields)) + return self.tick(start + 300, **fields) - def test_router_failure_does_not_act(self): - recover.run(1000) - self.report["at"] = 1300 - self.save() - with patch.object(recover, "verdict", side_effect=ValueError("bad response")), patch.object(recover.subprocess, "run") as restart: - self.assertEqual(recover.run(1300), "router unavailable") - restart.assert_not_called() + def test_ladder_restart_then_current_rung_then_rung_below(self): + self.assertIn("flysim restarted", self.confirm(T0)) + self.assertEqual(self.acts, [("restart", None)]) + self.assertIn("settling", self.tick(T0 + 600)) + self.assertIn("reset to rung 12", self.confirm(T0 + 300 + recover.SETTLE + 10)) + start = T0 + 2 * (300 + recover.SETTLE + 10) + self.assertIn("reset to rung 11", self.confirm(start, rank=12)) + self.assertEqual([a for a in self.acts], [("restart", None), ("reset", 12), ("reset", 11)]) - def test_cleared_probe_preserves_cooldown(self): - with patch.object(recover, "verdict", return_value=True), patch.object(recover.subprocess, "run") as restart: - recover.run(1000) - self.report["at"] = 1300 - self.save() - recover.run(1300) - self.report.update(at=1600, suspected=0) - self.save() - recover.run(1600) - self.report.update(at=1900, suspected=1) - self.save() - recover.run(1900) - self.report["at"] = 2200 - self.save() - self.assertIn("cooldown", recover.run(2200)) - restart.assert_called_once() + def test_spent_reset_budget_holds_to_restarts_three_hours_apart(self): + at = T0 + for _ in range(3): + self.confirm(at) + at += 300 + recover.SETTLE + 10 + self.assertEqual(len(self.acts), 3) + self.assertIn("holding", self.confirm(at)) + # the trap never cleared while holding, so the first probe after the hold acts + self.assertIn("flysim restarted", self.tick(at + recover.HOLD)) + self.assertEqual(self.acts[-1], ("restart", None)) - def test_router_rejection_does_not_restart(self): - recover.run(1000) - self.report["at"] = 1300 - self.save() - with patch.object(recover, "verdict", return_value=False), patch.object(recover.subprocess, "run") as restart: - self.assertIn("did not confirm", recover.run(1300)) - restart.assert_not_called() + def test_deeper_resets_never_go_below_the_rung_under_the_best(self): + state = {"level": 5, "bestRank": 12, "resets": [], "actedAt": 0} + recover.write_json(recover.STATE, state) + self.confirm(T0, rank=11) + self.assertEqual(self.acts, [("reset", 11)]) - def test_partial_router_configuration_cannot_bypass_veto(self): - with patch.dict("os.environ", {"FLY_LOOP_ROUTER_URL": "http://router/v1"}, clear=True): - self.assertFalse(recover.verdict(self.report)) - with patch.dict("os.environ", {"FLY_LOOP_MODEL": "free"}, clear=True): - self.assertFalse(recover.verdict(self.report)) + def test_new_best_rung_starts_the_ladder_over(self): + recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0}) + self.confirm(T0, rank=13) + self.assertEqual(self.acts, [("restart", None)]) - def test_unconfigured_router_uses_deterministic_confirmation(self): - with patch.dict("os.environ", {}, clear=True): - self.assertTrue(recover.verdict(self.report)) + def test_quiet_hours_start_the_ladder_over(self): + recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "actedAt": 0, + "lastSuspectedAt": T0 - recover.QUIET - 1}) + self.confirm(T0) + self.assertEqual(self.acts, [("restart", None)]) + + def test_stale_clear_or_acted_reports_never_act(self): + self.report(at=T0 - 700) + self.assertIn("not a fresh", recover.run(T0)) + self.report(suspected=0, at=T0) + self.assertIn("not a fresh", recover.run(T0 + 1)) + self.report(at=T0, action="restart") + self.assertIn("not a fresh", recover.run(T0 + 2)) + self.assertEqual(self.acts, []) + + def test_a_clear_probe_breaks_the_streak(self): + self.tick(T0) + self.tick(T0 + 300, suspected=0) + self.assertIn("second probe", self.tick(T0 + 600)) + self.assertEqual(self.acts, []) + + def test_state_survives_a_reboot_and_ignores_a_missing_or_corrupt_file(self): + self.confirm(T0) + self.assertTrue(recover.STATE.exists()) + self.assertEqual(json.loads(recover.STATE.read_text())["level"], 1) + recover.STATE.write_text("{nope") + self.assertIn("second probe", self.tick(T0 + 5000)) + + def test_model_vetoes_delay_a_step_but_never_deny_it(self): + with patch.object(recover, "verdict", return_value=(False, "m")): + self.tick(T0) + for n in range(1, recover.VETO_LIMIT + 1): + self.assertIn(f"({n}/{recover.VETO_LIMIT})", self.tick(T0 + 300 * n)) + self.assertIn("veto limit reached", self.tick(T0 + 300 * (recover.VETO_LIMIT + 1))) + self.assertEqual(self.acts, [("restart", None)]) + + def test_no_model_answer_falls_back_to_the_watchdog(self): + with patch.object(recover, "verdict", return_value=(None, None)): + self.assertIn("watchdog alone", self.confirm(T0)) + + def test_models_are_tried_in_order_until_one_answers(self): + recover.os.environ.update(FLY_LOOP_ROUTER_URL="http://router/v1", FLY_LOOP_MODELS="a, b ,c") + replies = iter([OSError("429"), {"choices": [{"message": {"content": "thinking...\n```json\n{\"stuck\": true}\n```"}}]}]) + asked = [] + + class Reply: + def __init__(self, body): + self.body = body + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + def read(self): + return json.dumps(self.body).encode() + + def fake(request, timeout): + asked.append(json.loads(request.data)["model"]) + reply = next(replies) + if isinstance(reply, Exception): + raise reply + return Reply(reply) + + with patch.object(recover, "urlopen", side_effect=fake): + self.assertEqual(recover.verdict({"reason": "unrewarded"}), (True, "b")) + self.assertEqual(asked, ["a", "b"]) + + def test_parse_stuck_tolerates_prose_and_rejects_non_booleans(self): + self.assertIs(recover.parse_stuck('{"stuck":false}'), False) + self.assertIs(recover.parse_stuck('Answer: {"stuck": true} done'), True) + self.assertIsNone(recover.parse_stuck('{"stuck": "yes"}')) + self.assertIsNone(recover.parse_stuck("")) + self.assertIsNone(recover.parse_stuck(None)) + + def test_notice_goes_countdown_acting_done_for_the_splash(self): + phases = [] + real = recover.notice + with patch.object(recover, "notice", side_effect=lambda base, phase, now: phases.append(phase) or real(base, phase, now)): + self.confirm(T0 + 300 + recover.SETTLE) + recover.write_json(recover.STATE, dict(json.loads(recover.STATE.read_text()), actedAt=0)) + self.confirm(T0 + 2 * (300 + recover.SETTLE)) + notice = json.loads(recover.NOTICE.read_text()) + self.assertEqual(phases, ["countdown", "acting", "done"] * 2) + self.assertEqual((notice["v"], notice["action"], notice["toRung"], notice["toLabel"], notice["fromRung"]), + (1, "reset", 12, "MT. MOON", 12)) + self.assertEqual(notice["executeAt"] - notice["announcedAt"], recover.COUNTDOWN) + + def test_failed_action_is_reported_and_still_climbs(self): + self.act.stop() + with patch.object(recover, "act", return_value=False): + self.assertIn("FAILED", self.confirm(T0)) + self.act.start() + self.assertEqual(json.loads(recover.NOTICE.read_text())["phase"], "failed") + self.assertEqual(json.loads(recover.STATE.read_text())["level"], 1) + + def test_history_records_every_step(self): + self.confirm(T0) + lines = [json.loads(line) for line in recover.HISTORY.read_text().splitlines()] + self.assertEqual(lines[-1]["event"], "restart") + self.assertTrue(lines[-1]["ok"]) + + def test_ladder_labels_match_the_rust_table(self): + source = (repo / "services/flysim/crates/flybrain-gb/src/pokemon_red/mod.rs").read_text() + table = re.search(r"RANK_LADDER: \[&str; \d+\] = \[(.*?)\];", source, re.S).group(1) + self.assertEqual(recover.LADDER, re.findall(r'"([^"]*)"', table)) if __name__ == "__main__": diff --git a/infra/units/fly-loop-recover.service b/infra/units/fly-loop-recover.service index 429fb2b..a2dff4d 100644 --- a/infra/units/fly-loop-recover.service +++ b/infra/units/fly-loop-recover.service @@ -1,8 +1,12 @@ [Unit] -Description=Check confirmed macro loops for minimal flysim recovery +Description=Recover confirmed macro loops: flysim restart, then milestone resets [Service] Type=oneshot User=fly EnvironmentFile=-/etc/fly/loop-recovery.env +# The ladder's state and history outlive a reboot (/run does not). +StateDirectory=fly-loop-recover +# A step is a 60 s on-stream countdown, up to 90 s of model calls and 4 min to verify. +TimeoutStartSec=12min ExecStart=/opt/fly/bin/fly-loop-recover From 9c9cec49a92b3808da4df1ae7a3fdcf9071cc4fa Mon Sep 17 00:00:00 2001 From: acamilo Date: Mon, 28 Sep 2026 21:25:58 +0000 Subject: [PATCH 2/4] loop recovery: reset only to an archive the running build can restore fly-reset-to-milestone does not check compatibility and a flysim that refuses every checkpoint does not start. fly-loop-reset --check applies 05-deploy's rule (identical, or an adapter-only difference named in FLY_ACCEPT_ADAPTERS); the ladder picks the highest restorable rung and falls back to a restart when there is none. --- infra/bin/fly-loop-recover | 15 ++++++-- infra/bin/fly-loop-reset | 60 ++++++++++++++++++++++++++++++-- infra/docs/loop-recovery.md | 8 +++++ infra/tests/test_loop_recover.py | 16 +++++++++ 4 files changed, 93 insertions(+), 6 deletions(-) diff --git a/infra/bin/fly-loop-recover b/infra/bin/fly-loop-recover index 3ea369f..1fdc355 100755 --- a/infra/bin/fly-loop-recover +++ b/infra/bin/fly-loop-recover @@ -148,12 +148,21 @@ def plan(level, rank, best, resets_left): if level == 0 or not resets_left or rank is None: return "restart", None ceiling = rank if level == 1 else min(rank, max(best, rank) - 1) - below = [rung for rung in rungs() if rung <= ceiling] - if below: - return "reset", below[-1] + for rung in reversed([rung for rung in rungs() if rung <= ceiling]): + if restorable(rung): + return "reset", rung return "restart", None +def restorable(rung): + """Whether the running build can restore this archive (fly-loop-reset --check).""" + try: + return subprocess.run(["sudo", "-n", RESET_BIN, "--check", str(rung)], check=False, timeout=60, + stdout=subprocess.DEVNULL).returncode == 0 + except (OSError, subprocess.SubprocessError): + return False + + def status(): try: with urlopen(STATUS_URL, timeout=5) as response: diff --git a/infra/bin/fly-loop-reset b/infra/bin/fly-loop-reset index de16570..99a5e4f 100755 --- a/infra/bin/fly-loop-reset +++ b/infra/bin/fly-loop-reset @@ -8,25 +8,79 @@ # again whether or not the reset worked, so the stream never stays down on a failure; # fly-reset-to-milestone copies both stores aside before it rewrites anything. # -# Usage: fly-loop-reset +# fly-reset-to-milestone does not check that the running build can restore the archive, and a +# flysim that refuses every checkpoint refuses to start: a black stream. So an archive is only +# used when its compatibility string equals this build's, or differs only in an adapter id that +# FLY_ACCEPT_ADAPTERS in /etc/fly/fly.env names (05-deploy.sh's adapter_migration_accepted). +# +# Usage: fly-loop-reset [--check] (--check: exit 0 if restorable, 3 if not; touch nothing) set -euo pipefail : "${FLY_STATE_DIR:=/srv/fly/state}" : "${FLY_SERVICE:=flysim.service}" : "${FLY_RESET_BIN:=/opt/fly/bin/fly-reset-to-milestone}" +: "${FLY_ENV_FILE:=/etc/fly/fly.env}" +: "${FLY_BIN:=/opt/fly/current/flysim}" log() { echo "fly-loop-reset: $*" >&2; } +CHECK=0 +if [ "${1:-}" = "--check" ]; then + CHECK=1 + shift +fi RANK="${1:-}" if [ "$#" -ne 1 ] || ! [[ "$RANK" =~ ^[0-9]{1,2}$ ]]; then - log "usage: fly-loop-reset " + log "usage: fly-loop-reset [--check] " exit 2 fi -if [ ! -f "${FLY_STATE_DIR}/milestone-${RANK}.checkpoint" ]; then +ARCHIVE="${FLY_STATE_DIR}/milestone-${RANK}.checkpoint" +if [ ! -f "$ARCHIVE" ]; then log "no milestone archive for rung ${RANK}; nothing touched" exit 1 fi +# The same rule as 05-deploy.sh: exactly one '/'-separated field differs, it is the adapter +# (index 1), and the archive's adapter id is listed. +migration_accepted() { + local old="$1" new="$2" accepted="$3" i differing=0 index=-1 entry + local -a old_parts new_parts + IFS='/' read -r -a old_parts <<< "$old" + IFS='/' read -r -a new_parts <<< "$new" + [ "${#old_parts[@]}" -eq "${#new_parts[@]}" ] || return 1 + for ((i = 0; i < ${#old_parts[@]}; i++)); do + if [ "${old_parts[$i]}" != "${new_parts[$i]}" ]; then + differing=$((differing + 1)) + index=$i + fi + done + [ "$differing" -eq 1 ] && [ "$index" -eq 1 ] || return 1 + for entry in ${accepted//,/ }; do + [ "$entry" = "${old_parts[1]}" ] && return 0 + done + return 1 +} + +accepted="" +if [ -r "$FLY_ENV_FILE" ]; then + set -a + # shellcheck disable=SC1090 + . "$FLY_ENV_FILE" + set +a + accepted="${FLY_ACCEPT_ADAPTERS:-}" +fi +archive_compat="$(head -c 262144 "$ARCHIVE" 2>/dev/null | grep -a -o -m1 '"compatibility":"[^"]*"' | head -n1 | cut -d'"' -f4 || true)" +build_compat="$("$FLY_BIN" --print-compatibility 2>/dev/null | tail -n1 || true)" +if [ -z "$archive_compat" ] || [ -z "$build_compat" ]; then + log "cannot read the compatibility of rung ${RANK}'s archive or of this build; nothing touched" + exit 3 +fi +if [ "$archive_compat" != "$build_compat" ] && ! migration_accepted "$archive_compat" "$build_compat" "$accepted"; then + log "rung ${RANK}'s archive ($(echo "$archive_compat" | cut -d/ -f2)) is not restorable by this build ($(echo "$build_compat" | cut -d/ -f2)), FLY_ACCEPT_ADAPTERS='${accepted}'; nothing touched" + exit 3 +fi +[ "$CHECK" -eq 0 ] || exit 0 + log "stopping ${FLY_SERVICE} to reset to rung ${RANK}" systemctl stop "$FLY_SERVICE" status=0 diff --git a/infra/docs/loop-recovery.md b/infra/docs/loop-recovery.md index 8f3df98..cff1c24 100644 --- a/infra/docs/loop-recovery.md +++ b/infra/docs/loop-recovery.md @@ -19,6 +19,14 @@ a ladder one step per trap that outlives the previous step: - **Starting over:** the ladder returns to level 0 when the fly reaches a new best rung, or after six hours with no suspected report. +A milestone step only uses an archive the running build can restore: `fly-loop-reset --check` +compares the archive's compatibility string with `flysim --print-compatibility` and accepts an +adapter-only difference that `FLY_ACCEPT_ADAPTERS` in `/etc/fly/fly.env` names (the rule +`05-deploy.sh` applies). A flysim that refuses every checkpoint refuses to start, so an archive +from an older adapter is skipped for the next one down unless the deploy named its adapter; with +none restorable the step is a restart. Keep `FLY_ACCEPT_ADAPTERS` in the release env file so a +deploy does not drop it. + A milestone step runs `/opt/fly/bin/fly-loop-reset ` through sudo (the one line in `config/fly-sudoers`): stop flysim, `fly-reset-to-milestone`, start flysim. flysim is started again even when the reset fails. The reset copies both stores to `/srv/fly/state.reset-` first, as diff --git a/infra/tests/test_loop_recover.py b/infra/tests/test_loop_recover.py index 881971c..264a3a7 100644 --- a/infra/tests/test_loop_recover.py +++ b/infra/tests/test_loop_recover.py @@ -40,6 +40,10 @@ class RecoveryTests(unittest.TestCase): for key in ("FLY_LOOP_ROUTER_URL", "FLY_LOOP_MODELS", "FLY_LOOP_MODEL", "FLY_LOOP_ROUTER_KEY"): recover.os.environ.pop(key, None) self.acts = [] + self.unrestorable = set() + restorable = patch.object(recover, "restorable", side_effect=lambda rung: rung not in self.unrestorable) + restorable.start() + self.addCleanup(restorable.stop) self.act = patch.object(recover, "act", side_effect=lambda action, target: self.acts.append((action, target)) or True) self.act.start() self.addCleanup(self.act.stop) @@ -89,6 +93,18 @@ class RecoveryTests(unittest.TestCase): self.confirm(T0, rank=11) self.assertEqual(self.acts, [("reset", 11)]) + def test_an_archive_this_build_cannot_restore_is_skipped(self): + self.unrestorable = {11} + recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0}) + self.confirm(T0) + self.assertEqual(self.acts, [("reset", 10)]) + + def test_no_restorable_archive_means_a_restart(self): + self.unrestorable = {1, 9, 10, 11, 12} + recover.write_json(recover.STATE, {"level": 1, "bestRank": 12, "resets": [], "actedAt": 0}) + self.confirm(T0) + self.assertEqual(self.acts, [("restart", None)]) + def test_new_best_rung_starts_the_ladder_over(self): recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0}) self.confirm(T0, rank=13) From f205d95e5eed095c3ad13fbf4a67949493a6e0a4 Mon Sep 17 00:00:00 2001 From: acamilo Date: Mon, 28 Sep 2026 21:34:32 +0000 Subject: [PATCH 3/4] loop recovery: review-ladder fixes, root never runs a fly-writable binary B1: fly-loop-reset runs only /opt/fly/sbin/flysim, a root-owned copy 05-deploy installs from the release tarball after checking it against the tarball's MANIFEST; fixed paths, fly.env parsed as data, env -i. B2: the wrapper pauses fly-watchdog.timer (and waits out a running probe) for the reset, and starts flysim and the timer on every exit. H1: a step that raises or times out is a failed step; the ladder state is saved before the step runs. H2: one --list call computes the build's compatibility once; TimeoutStartSec 25 min. H3: level 2+ resets to the rung below the best or restarts, never lower. Tests run the wrapper against a fake flysim, systemctl and archives. --- infra/05-deploy.sh | 15 +++ infra/bin/fly-loop-recover | 59 ++++++---- infra/bin/fly-loop-reset | 164 +++++++++++++++++++-------- infra/docs/loop-recovery.md | 34 +++--- infra/tests/test_loop_recover.py | 117 ++++++++++++++++++- infra/units/fly-loop-recover.service | 5 +- 6 files changed, 304 insertions(+), 90 deletions(-) diff --git a/infra/05-deploy.sh b/infra/05-deploy.sh index 29e6620..061e6dd 100644 --- a/infra/05-deploy.sh +++ b/infra/05-deploy.sh @@ -272,6 +272,21 @@ if [ -n "$RELEASE_TARBALL" ]; then cpu_pin chown -R fly:fly "$release_path" fi + # A root-owned flysim for fly-loop-reset, which runs as root and so must never execute the + # fly-owned copy above (the fly account could replace it). Taken from the tarball on the host + # and checked against the tarball's own MANIFEST on every deploy, so a release installed by + # an earlier deploy is covered too. + log "05-deploy: installing a root-owned flysim at /opt/fly/sbin/flysim for fly-loop-reset" + root_flysim="$(mktemp)" + tar -xzOf "$RELEASE_TARBALL" ./flysim > "$root_flysim" + want_sha="$(tar -xzOf "$RELEASE_TARBALL" ./MANIFEST | awk '$2 == "flysim" || $2 == "./flysim" {print $1; exit}')" + [ -n "$want_sha" ] && [ "$(sha256sum "$root_flysim" | cut -d' ' -f1)" = "$want_sha" ] \ + || { rm -f "$root_flysim"; die "the tarball's flysim does not match its MANIFEST; not installing /opt/fly/sbin/flysim"; } + ct_exec "$CTID" -- install -d -o root -g root -m 0755 /opt/fly/sbin + ct_push_file "$CTID" "$root_flysim" /opt/fly/sbin/flysim.new 0755 + ct_exec "$CTID" -- sh -c 'chown root:root /opt/fly/sbin/flysim.new && mv -f /opt/fly/sbin/flysim.new /opt/fly/sbin/flysim' + rm -f "$root_flysim" + log "05-deploy: overlaying infra-owned stage/serve.{mjs,sh} (docs/design/infra.md's 'own tiny static server' clarification)" converge_file "$CTID" "$INFRA_DIR/config/serve.mjs" "${release_path}/stage/serve.mjs" 0644 fly:fly >/dev/null converge_file "$CTID" "$INFRA_DIR/config/serve.sh" "${release_path}/stage/serve.sh" 0755 fly:fly >/dev/null diff --git a/infra/bin/fly-loop-recover b/infra/bin/fly-loop-recover index 1fdc355..1d46b62 100755 --- a/infra/bin/fly-loop-recover +++ b/infra/bin/fly-loop-recover @@ -32,6 +32,9 @@ MAX_RESETS = 2 # milestone resets per DAY VETO_LIMIT = 3 # model vetoes of one step before it goes ahead anyway MODEL_BUDGET = 90 # seconds across all models in one decision VERIFY_TIMEOUT = 240 +LIST_TIMEOUT = 120 +RESTART_TIMEOUT = 180 +RESET_TIMEOUT = 480 # the wrapper waits up to 60 s for the watchdog, then stop, reset, start # docs/design/ladder.md; mirrors RANK_LADDER in flybrain-gb's pokemon_red/mod.rs (a test pins it). LADDER = [ @@ -140,27 +143,30 @@ def rungs(): def plan(level, rank, best, resets_left): - """(action, target rung) for a ladder level; a spent budget or no archive falls back to restart. + """(action, target rung) for a ladder level; a spent budget or no restorable archive is a restart. - Level 1 resets to the current rung's archive, level 2 and beyond to the archive below the - best rung (never lower), so repeated days of a trap cannot walk the run down the ladder. + Level 1 resets to the current rung's archive; level 2 and beyond to the rung below the best, + and never lower, so repeated days of a trap cannot walk the run down the ladder. """ if level == 0 or not resets_left or rank is None: return "restart", None - ceiling = rank if level == 1 else min(rank, max(best, rank) - 1) - for rung in reversed([rung for rung in rungs() if rung <= ceiling]): - if restorable(rung): - return "reset", rung - return "restart", None + target = rank if level == 1 else max(best, rank) - 1 + if target > rank or target not in restorable_rungs(): + return "restart", None + return "reset", target -def restorable(rung): - """Whether the running build can restore this archive (fly-loop-reset --check).""" +def restorable_rungs(): + """The rungs the running build can restore (fly-loop-reset --list), empty on any failure.""" try: - return subprocess.run(["sudo", "-n", RESET_BIN, "--check", str(rung)], check=False, timeout=60, - stdout=subprocess.DEVNULL).returncode == 0 - except (OSError, subprocess.SubprocessError): - return False + done = subprocess.run(["sudo", "-n", RESET_BIN, "--list"], check=False, timeout=LIST_TIMEOUT, + capture_output=True, text=True) + except (OSError, subprocess.SubprocessError) as error: + print("restorable rungs unknown:", type(error).__name__, file=sys.stderr) + return set() + if done.returncode != 0: + return set() + return {int(line) for line in done.stdout.split() if line.isdigit()} def status(): @@ -184,10 +190,14 @@ def verify(target, deadline): def act(action, target): if action == "restart": - command = ["sudo", "-n", "systemctl", "restart", "flysim.service"] + command, limit = ["sudo", "-n", "systemctl", "restart", "flysim.service"], RESTART_TIMEOUT else: - command = ["sudo", "-n", RESET_BIN, str(target)] - done = subprocess.run(command, check=False, timeout=VERIFY_TIMEOUT) + command, limit = ["sudo", "-n", RESET_BIN, str(target)], RESET_TIMEOUT + try: + done = subprocess.run(command, check=False, timeout=limit) + except (OSError, subprocess.SubprocessError) as error: + print(f"{action} did not complete: {type(error).__name__}", file=sys.stderr) + return False if done.returncode != 0: return False return verify(target, time.monotonic() + VERIFY_TIMEOUT) @@ -268,15 +278,18 @@ def run(now=None, sleep=time.sleep, clock=None): base.update(toRung=target, toLabel=LADDER[target] if 0 <= target < len(LADDER) else None) notice(base, "countdown", now) sleep(COUNTDOWN) - notice(base, "acting", clock()) + started = int(clock()) + # Recorded before acting, so a helper killed mid-step has still climbed and spent the reset. + state.update(level=level + 1, actedAt=started, lastAction=action, vetoes=0) + state.pop("observedAt", None) + if action == "reset": + state["resets"].append(started) + write_json(STATE, state) + notice(base, "acting", started) ok = act(action, target) finished = int(clock()) notice(base, "done" if ok else "failed", finished) - - state.update(level=level + 1, actedAt=finished, lastAction=action, vetoes=0) - state.pop("observedAt", None) - if action == "reset": - state["resets"].append(finished) + state["actedAt"] = finished write_json(STATE, state) log({"at": finished, "event": action, "target": target, "ok": ok, "why": why, "rank": rank, "level": level, "sequence": report.get("sequence"), "map": report.get("map")}) diff --git a/infra/bin/fly-loop-reset b/infra/bin/fly-loop-reset index 99a5e4f..00d2d9c 100755 --- a/infra/bin/fly-loop-reset +++ b/infra/bin/fly-loop-reset @@ -3,47 +3,80 @@ # # fly-loop-recover runs as User=fly and reaches this through one NOPASSWD sudoers line # (config/fly-sudoers). It does what infra/docs/runbook.md "Restart the run from a rung" -# does by hand for a release whose adapter already wrote the archive: stop flysim, promote -# milestone-.checkpoint with fly-reset-to-milestone, start flysim. flysim is started -# again whether or not the reset worked, so the stream never stays down on a failure; -# fly-reset-to-milestone copies both stores aside before it rewrites anything. +# does by hand: stop flysim, promote milestone-.checkpoint with fly-reset-to-milestone, +# start flysim. # -# fly-reset-to-milestone does not check that the running build can restore the archive, and a -# flysim that refuses every checkpoint refuses to start: a black stream. So an archive is only -# used when its compatibility string equals this build's, or differs only in an adapter id that -# FLY_ACCEPT_ADAPTERS in /etc/fly/fly.env names (05-deploy.sh's adapter_migration_accepted). +# Root never executes anything the fly account can write. The release tree under +# /opt/fly/releases is fly-owned, so the flysim run here is the root-owned copy 05-deploy +# installs at /opt/fly/sbin/flysim straight from the release tarball; without it every rung is +# refused and the ladder falls back to restarts. Paths are fixed (test overrides are honoured +# only when not root), and /etc/fly/fly.env is read as KEY=VALUE data, never sourced. # -# Usage: fly-loop-reset [--check] (--check: exit 0 if restorable, 3 if not; touch nothing) +# fly-reset-to-milestone does not check that the build can restore the archive, and a flysim +# that refuses every checkpoint refuses to start: a black stream. So a rung is only used when +# its archive's compatibility string equals this build's, or differs only in an adapter id that +# FLY_ACCEPT_ADAPTERS names (05-deploy.sh's adapter_migration_accepted). +# +# fly-watchdog restarts a flysim whose /healthz fails, which during a reset would start it on a +# half-rewritten store; its timer is stopped for the reset and started again on every exit path. +# +# Usage: fly-loop-reset reset to that rung +# fly-loop-reset --list print the rungs this build can restore, one per line set -euo pipefail -: "${FLY_STATE_DIR:=/srv/fly/state}" -: "${FLY_SERVICE:=flysim.service}" -: "${FLY_RESET_BIN:=/opt/fly/bin/fly-reset-to-milestone}" -: "${FLY_ENV_FILE:=/etc/fly/fly.env}" -: "${FLY_BIN:=/opt/fly/current/flysim}" +STATE_DIR=/srv/fly/state +ENV_FILE=/etc/fly/fly.env +FLYSIM=/opt/fly/sbin/flysim +RESET_BIN=/opt/fly/bin/fly-reset-to-milestone +if [ "$(id -u)" -ne 0 ]; then + STATE_DIR="${FLY_LOOP_RESET_TEST_STATE_DIR:-$STATE_DIR}" + ENV_FILE="${FLY_LOOP_RESET_TEST_ENV_FILE:-$ENV_FILE}" + FLYSIM="${FLY_LOOP_RESET_TEST_FLYSIM:-$FLYSIM}" + RESET_BIN="${FLY_LOOP_RESET_TEST_RESET_BIN:-$RESET_BIN}" +fi +SERVICE=flysim.service +WATCHDOG_TIMER=fly-watchdog.timer +WATCHDOG_SERVICE=fly-watchdog.service +readonly STATE_DIR ENV_FILE FLYSIM RESET_BIN SERVICE WATCHDOG_TIMER WATCHDOG_SERVICE log() { echo "fly-loop-reset: $*" >&2; } +usage() { log "usage: fly-loop-reset | --list"; exit 2; } -CHECK=0 -if [ "${1:-}" = "--check" ]; then - CHECK=1 - shift +LIST=0 +RANK="" +[ "$#" -eq 1 ] || usage +if [ "$1" = --list ]; then + LIST=1 +elif [[ "$1" =~ ^[0-9]{1,2}$ ]]; then + RANK="$1" +else + usage fi -RANK="${1:-}" -if [ "$#" -ne 1 ] || ! [[ "$RANK" =~ ^[0-9]{1,2}$ ]]; then - log "usage: fly-loop-reset [--check] " - exit 2 -fi -ARCHIVE="${FLY_STATE_DIR}/milestone-${RANK}.checkpoint" -if [ ! -f "$ARCHIVE" ]; then - log "no milestone archive for rung ${RANK}; nothing touched" - exit 1 + +# FLY_* lines of fly.env as `env` arguments. Values may be quoted the systemd way; nothing +# is evaluated. +env_args=() +accepted="" +if [ -r "$ENV_FILE" ]; then + while IFS= read -r line || [ -n "$line" ]; do + [[ "$line" =~ ^(FLY_[A-Z0-9_]*)=(.*)$ ]] || continue + key="${BASH_REMATCH[1]}" + value="${BASH_REMATCH[2]}" + if [[ "$value" =~ ^\"(.*)\"$ ]] || [[ "$value" =~ ^\'(.*)\'$ ]]; then + value="${BASH_REMATCH[1]}" + fi + env_args+=("$key=$value") + if [ "$key" = FLY_ACCEPT_ADAPTERS ]; then + accepted="$value" + fi + done < "$ENV_FILE" fi +clean_env() { env -i PATH=/usr/sbin:/usr/bin:/sbin:/bin HOME=/root "${env_args[@]}" "$@"; } # The same rule as 05-deploy.sh: exactly one '/'-separated field differs, it is the adapter # (index 1), and the archive's adapter id is listed. migration_accepted() { - local old="$1" new="$2" accepted="$3" i differing=0 index=-1 entry + local old="$1" new="$2" i differing=0 index=-1 entry local -a old_parts new_parts IFS='/' read -r -a old_parts <<< "$old" IFS='/' read -r -a new_parts <<< "$new" @@ -61,30 +94,63 @@ migration_accepted() { return 1 } -accepted="" -if [ -r "$FLY_ENV_FILE" ]; then - set -a - # shellcheck disable=SC1090 - . "$FLY_ENV_FILE" - set +a - accepted="${FLY_ACCEPT_ADAPTERS:-}" +build_compat="" +if [ -x "$FLYSIM" ] && [ "$(stat -c %u "$FLYSIM")" = "$(id -u)" ]; then + build_compat="$(clean_env timeout 90 "$FLYSIM" --print-compatibility 2>/dev/null | tail -n1 || true)" fi -archive_compat="$(head -c 262144 "$ARCHIVE" 2>/dev/null | grep -a -o -m1 '"compatibility":"[^"]*"' | head -n1 | cut -d'"' -f4 || true)" -build_compat="$("$FLY_BIN" --print-compatibility 2>/dev/null | tail -n1 || true)" -if [ -z "$archive_compat" ] || [ -z "$build_compat" ]; then - log "cannot read the compatibility of rung ${RANK}'s archive or of this build; nothing touched" +if [ -z "$build_compat" ]; then + log "no compatibility string from $FLYSIM (missing, not owned by $(id -un), or failed); no rung is restorable" + [ "$LIST" -eq 1 ] && exit 0 exit 3 fi -if [ "$archive_compat" != "$build_compat" ] && ! migration_accepted "$archive_compat" "$build_compat" "$accepted"; then - log "rung ${RANK}'s archive ($(echo "$archive_compat" | cut -d/ -f2)) is not restorable by this build ($(echo "$build_compat" | cut -d/ -f2)), FLY_ACCEPT_ADAPTERS='${accepted}'; nothing touched" - exit 3 -fi -[ "$CHECK" -eq 0 ] || exit 0 -log "stopping ${FLY_SERVICE} to reset to rung ${RANK}" -systemctl stop "$FLY_SERVICE" +restorable() { + local archive="${STATE_DIR}/milestone-$1.checkpoint" compat + [ -f "$archive" ] || return 1 + compat="$(head -c 262144 "$archive" 2>/dev/null | grep -a -o -m1 '"compatibility"[[:space:]]*:[[:space:]]*"[^"]*"' | head -n1 | cut -d'"' -f4 || true)" + [ -n "$compat" ] || return 1 + [ "$compat" = "$build_compat" ] || migration_accepted "$compat" "$build_compat" +} + +if [ "$LIST" -eq 1 ]; then + for archive in "${STATE_DIR}"/milestone-*.checkpoint; do + [ -e "$archive" ] || continue + rung="${archive##*/milestone-}" + rung="${rung%.checkpoint}" + if [[ "$rung" =~ ^[0-9]{1,2}$ ]] && restorable "$rung"; then + echo "$rung" + fi + done | sort -n + exit 0 +fi + +if ! restorable "$RANK"; then + log "rung ${RANK}'s archive is missing or not restorable by this build (FLY_ACCEPT_ADAPTERS='${accepted}'); nothing touched" + exit 3 +fi + +# shellcheck disable=SC2317 # invoked by the EXIT trap +finish() { + local status=$? + systemctl start "$SERVICE" || log "starting ${SERVICE} failed" + systemctl start "$WATCHDOG_TIMER" || log "starting ${WATCHDOG_TIMER} failed" + exit "$status" +} +trap finish EXIT +trap 'exit 143' TERM INT HUP + +log "pausing ${WATCHDOG_TIMER} and stopping ${SERVICE} to reset to rung ${RANK}" +systemctl stop "$WATCHDOG_TIMER" +for _ in $(seq 1 60); do + systemctl is-active --quiet "$WATCHDOG_SERVICE" || break + sleep 1 +done +if systemctl is-active --quiet "$WATCHDOG_SERVICE"; then + log "${WATCHDOG_SERVICE} still running after 60 s; not resetting" + exit 4 +fi +systemctl stop "$SERVICE" status=0 -"$FLY_RESET_BIN" "$RANK" || status=$? -[ "$status" -eq 0 ] || log "fly-reset-to-milestone ${RANK} failed (exit ${status}); starting ${FLY_SERVICE} on the state it left" -systemctl start "$FLY_SERVICE" +clean_env FLY_BIN="$FLYSIM" "$RESET_BIN" "$RANK" || status=$? +[ "$status" -eq 0 ] || log "fly-reset-to-milestone ${RANK} failed (exit ${status}); starting ${SERVICE} on the state it left" exit "$status" diff --git a/infra/docs/loop-recovery.md b/infra/docs/loop-recovery.md index cff1c24..953353a 100644 --- a/infra/docs/loop-recovery.md +++ b/infra/docs/loop-recovery.md @@ -19,21 +19,29 @@ a ladder one step per trap that outlives the previous step: - **Starting over:** the ladder returns to level 0 when the fly reaches a new best rung, or after six hours with no suspected report. -A milestone step only uses an archive the running build can restore: `fly-loop-reset --check` -compares the archive's compatibility string with `flysim --print-compatibility` and accepts an +A milestone step only uses an archive the running build can restore. `fly-loop-reset --list` +compares each archive's compatibility string with `flysim --print-compatibility` and accepts an adapter-only difference that `FLY_ACCEPT_ADAPTERS` in `/etc/fly/fly.env` names (the rule -`05-deploy.sh` applies). A flysim that refuses every checkpoint refuses to start, so an archive -from an older adapter is skipped for the next one down unless the deploy named its adapter; with -none restorable the step is a restart. Keep `FLY_ACCEPT_ADAPTERS` in the release env file so a -deploy does not drop it. +`05-deploy.sh` applies). A flysim that refuses every checkpoint refuses to start, so a step whose +rung is not restorable is a restart instead, never a lower rung. Keep `FLY_ACCEPT_ADAPTERS` in the +release env file so a deploy does not drop it. -A milestone step runs `/opt/fly/bin/fly-loop-reset ` through sudo (the one line in -`config/fly-sudoers`): stop flysim, `fly-reset-to-milestone`, start flysim. flysim is started again -even when the reset fails. The reset copies both stores to `/srv/fly/state.reset-` first, as -in the runbook, and needs no deploy when the running release wrote the archive or -`FLY_ACCEPT_ADAPTERS` already names its adapter. After each step the helper waits up to four -minutes for `/status` to report `running` (and the target rank for a reset); otherwise the step is -recorded as failed and the ladder still climbs. +The step runs `/opt/fly/bin/fly-loop-reset ` through sudo (the one line in +`config/fly-sudoers`). As root it: + +- runs only the root-owned `/opt/fly/sbin/flysim` that `05-deploy.sh` installs from the release + tarball after checking it against the tarball's MANIFEST, never the fly-owned release tree; + without that copy no rung is restorable; +- reads `fly.env` as `KEY=VALUE` data, never sources it, and ignores the caller's environment; +- stops `fly-watchdog.timer` and waits for a running probe to finish, so the watchdog cannot + start flysim on a half-rewritten store; stops flysim; runs `fly-reset-to-milestone`; and on every + exit path starts flysim and the watchdog timer again. + +The reset copies both stores to `/srv/fly/state.reset-` first, as in the runbook, and +clears milestone archives above the rung. After each step the helper waits up to four minutes for +`/status` to report `running` (and the target rank for a reset). Otherwise the step is recorded +as failed and the ladder still climbs. The step is recorded before it runs, so a helper killed +mid-step has still climbed and spent the reset. ## On stream diff --git a/infra/tests/test_loop_recover.py b/infra/tests/test_loop_recover.py index 264a3a7..405135c 100644 --- a/infra/tests/test_loop_recover.py +++ b/infra/tests/test_loop_recover.py @@ -12,6 +12,8 @@ script = repo / "infra/bin/fly-loop-recover" spec = importlib.util.spec_from_loader("recover", importlib.machinery.SourceFileLoader("recover", str(script))) recover = importlib.util.module_from_spec(spec) spec.loader.exec_module(recover) +RESTORABLE_RUNGS = recover.restorable_rungs +WRAPPER = repo / "infra/bin/fly-loop-reset" T0 = 1_790_000_000 @@ -41,7 +43,7 @@ class RecoveryTests(unittest.TestCase): recover.os.environ.pop(key, None) self.acts = [] self.unrestorable = set() - restorable = patch.object(recover, "restorable", side_effect=lambda rung: rung not in self.unrestorable) + restorable = patch.object(recover, "restorable_rungs", side_effect=lambda: {1, 9, 10, 11, 12} - self.unrestorable) restorable.start() self.addCleanup(restorable.stop) self.act = patch.object(recover, "act", side_effect=lambda action, target: self.acts.append((action, target)) or True) @@ -93,11 +95,16 @@ class RecoveryTests(unittest.TestCase): self.confirm(T0, rank=11) self.assertEqual(self.acts, [("reset", 11)]) - def test_an_archive_this_build_cannot_restore_is_skipped(self): + def test_an_unrestorable_rung_below_the_best_means_a_restart_never_a_lower_rung(self): self.unrestorable = {11} recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0}) self.confirm(T0) - self.assertEqual(self.acts, [("reset", 10)]) + self.assertEqual(self.acts, [("restart", None)]) + + def test_a_fly_already_below_the_rung_under_its_best_is_restarted_not_reset(self): + recover.write_json(recover.STATE, {"level": 2, "bestRank": 12, "resets": [], "actedAt": 0}) + self.confirm(T0, rank=10) + self.assertEqual(self.acts, [("restart", None)]) def test_no_restorable_archive_means_a_restart(self): self.unrestorable = {1, 9, 10, 11, 12} @@ -207,6 +214,32 @@ class RecoveryTests(unittest.TestCase): self.assertEqual(json.loads(recover.NOTICE.read_text())["phase"], "failed") self.assertEqual(json.loads(recover.STATE.read_text())["level"], 1) + def test_a_step_that_raises_or_times_out_is_a_failed_step_with_state_saved(self): + self.act.stop() + with patch.object(recover.subprocess, "run", side_effect=recover.subprocess.TimeoutExpired("sudo", 1)): + self.assertIn("FAILED", self.confirm(T0)) + self.act.start() + self.assertEqual(json.loads(recover.NOTICE.read_text())["phase"], "failed") + self.assertEqual(json.loads(recover.STATE.read_text())["level"], 1) + + def test_state_is_saved_before_the_step_runs(self): + seen = [] + self.act.stop() + with patch.object(recover, "act", side_effect=lambda a, t: seen.append(json.loads(recover.STATE.read_text())) or True): + self.confirm(T0) + self.act.start() + self.assertEqual(seen[0]["level"], 1) + self.assertNotIn("observedAt", seen[0]) + + def test_restorable_rungs_parses_the_list_and_is_empty_on_any_failure(self): + ok = recover.subprocess.CompletedProcess([], 0, stdout="11\n12\n", stderr="") + with patch.object(recover.subprocess, "run", return_value=ok): + self.assertEqual(RESTORABLE_RUNGS(), {11, 12}) + with patch.object(recover.subprocess, "run", side_effect=OSError("no sudo")): + self.assertEqual(RESTORABLE_RUNGS(), set()) + with patch.object(recover.subprocess, "run", return_value=recover.subprocess.CompletedProcess([], 1, "12", "")): + self.assertEqual(RESTORABLE_RUNGS(), set()) + def test_history_records_every_step(self): self.confirm(T0) lines = [json.loads(line) for line in recover.HISTORY.read_text().splitlines()] @@ -219,5 +252,83 @@ class RecoveryTests(unittest.TestCase): self.assertEqual(recover.LADDER, re.findall(r'"([^"]*)"', table)) +BUILD = "lif-1ms-f64-v2/pokered-unique8-v7/abc/fly-kc-mbon-rstdp-v2" + + +@unittest.skipIf(recover.os.geteuid() == 0, "the wrapper ignores test overrides as root") +class WrapperTests(unittest.TestCase): + """infra/bin/fly-loop-reset against a fake flysim, systemctl, reset tool and archives.""" + + def setUp(self): + self.temp = tempfile.TemporaryDirectory() + self.addCleanup(self.temp.cleanup) + root = Path(self.temp.name) + self.root = root + self.calls = root / "calls.log" + (root / "state").mkdir() + (root / "bin").mkdir() + self.script(root / "bin/systemctl", f'echo "systemctl $*" >> {self.calls}; [ "$1" != is-active ] || exit 3') + self.script(root / "flysim", f'[ "$1" = --print-compatibility ] && echo "{BUILD}"') + self.script(root / "reset", f'echo "reset $* FLY_BIN=$FLY_BIN" >> {self.calls}') + self.env_file = root / "fly.env" + self.env_file.write_text('FLY_MACRO_MODE=macros\nGAME_TITLE=Pokemon (Red) $(touch pwned)\n') + self.archive(12, BUILD) + self.archive(11, BUILD.replace("-v7", "-v6")) + self.archive(1, BUILD.replace("-v7", "-v5")) + + def script(self, path, body): + path.write_text("#!/bin/sh\n" + body + "\n") + path.chmod(0o755) + + def archive(self, rung, compat): + (self.root / f"state/milestone-{rung}.checkpoint").write_bytes( + b"FLYSIM01" + json.dumps({"generation": 1, "compatibility": compat}).encode() + b"\x00" * 64) + + def run_wrapper(self, *args): + env = {"PATH": f"{self.root}/bin:/usr/bin:/bin", "FLY_LOOP_RESET_TEST_STATE_DIR": str(self.root / "state"), + "FLY_LOOP_RESET_TEST_ENV_FILE": str(self.env_file), "FLY_LOOP_RESET_TEST_FLYSIM": str(self.root / "flysim"), + "FLY_LOOP_RESET_TEST_RESET_BIN": str(self.root / "reset")} + return recover.subprocess.run([str(WRAPPER), *args], env=env, capture_output=True, text=True, timeout=60) + + def log(self): + return self.calls.read_text() if self.calls.exists() else "" + + def test_list_names_only_rungs_this_build_restores(self): + self.assertEqual(self.run_wrapper("--list").stdout.split(), ["12"]) + self.env_file.write_text("FLY_ACCEPT_ADAPTERS=pokered-unique8-v6\n") + self.assertEqual(self.run_wrapper("--list").stdout.split(), ["11", "12"]) + self.assertFalse((self.root / "pwned").exists()) + + def test_reset_pauses_the_watchdog_and_always_starts_flysim_again(self): + done = self.run_wrapper("12") + self.assertEqual(done.returncode, 0, done.stderr) + self.assertEqual([line.split()[:3] for line in self.log().splitlines() if not line.startswith("systemctl is-active")], [ + ["systemctl", "stop", "fly-watchdog.timer"], ["systemctl", "stop", "flysim.service"], + ["reset", "12", f"FLY_BIN={self.root}/flysim"], + ["systemctl", "start", "flysim.service"], ["systemctl", "start", "fly-watchdog.timer"]]) + + def test_a_failed_reset_still_starts_flysim_and_the_watchdog(self): + self.script(self.root / "reset", f'echo "reset $*" >> {self.calls}; exit 7') + self.assertEqual(self.run_wrapper("12").returncode, 7) + self.assertIn("systemctl start flysim.service", self.log()) + self.assertIn("systemctl start fly-watchdog.timer", self.log()) + + def test_an_unrestorable_or_missing_rung_touches_nothing(self): + for rung in ("11", "5"): + self.assertEqual(self.run_wrapper(rung).returncode, 3) + self.assertEqual(self.log(), "") + + def test_no_build_compatibility_means_nothing_is_restorable(self): + (self.root / "flysim").unlink() + self.assertEqual(self.run_wrapper("--list").stdout, "") + self.assertEqual(self.run_wrapper("12").returncode, 3) + self.assertEqual(self.log(), "") + + def test_bad_arguments_are_refused(self): + for args in ((), ("--check", "12"), ("12", "13"), ("../12",), ("123",), ("-1",)): + self.assertEqual(self.run_wrapper(*args).returncode, 2, args) + self.assertEqual(self.log(), "") + + if __name__ == "__main__": unittest.main() diff --git a/infra/units/fly-loop-recover.service b/infra/units/fly-loop-recover.service index a2dff4d..dc93b16 100644 --- a/infra/units/fly-loop-recover.service +++ b/infra/units/fly-loop-recover.service @@ -7,6 +7,7 @@ User=fly EnvironmentFile=-/etc/fly/loop-recovery.env # The ladder's state and history outlive a reboot (/run does not). StateDirectory=fly-loop-recover -# A step is a 60 s on-stream countdown, up to 90 s of model calls and 4 min to verify. -TimeoutStartSec=12min +# A step: the rung list (<=2 min), up to 90 s of model calls, a 60 s on-stream countdown, +# the reset (<=8 min) and 4 min to verify. systemd must never kill the wrapper mid-reset. +TimeoutStartSec=25min ExecStart=/opt/fly/bin/fly-loop-recover From 26892d85303ad77f29f8f0f3cc6a326301b12485 Mon Sep 17 00:00:00 2001 From: acamilo Date: Mon, 28 Sep 2026 21:37:40 +0000 Subject: [PATCH 4/4] loop recovery: wait on the watchdog's ActiveState, and only a root flysim that is the running build review-ladder r2: a oneshot probe mid-run is 'activating', which is-active does not count, so the wait never waited. The root copy must also be byte-identical to /opt/fly/current/flysim, so a manual rollback cannot approve an archive the running build refuses. Docs: the hold after an unrestorable rung, the post-deploy --list check, and never stopping the unit mid-step. --- infra/bin/fly-loop-reset | 24 ++++++++++++++++++++---- infra/docs/loop-recovery.md | 16 +++++++++++----- infra/tests/test_loop_recover.py | 26 ++++++++++++++++++++++++-- 3 files changed, 55 insertions(+), 11 deletions(-) diff --git a/infra/bin/fly-loop-reset b/infra/bin/fly-loop-reset index 00d2d9c..95ae93a 100755 --- a/infra/bin/fly-loop-reset +++ b/infra/bin/fly-loop-reset @@ -94,12 +94,21 @@ migration_accepted() { return 1 } +# The root copy must be the flysim that runs: a manual rollback of /opt/fly/current without a +# deploy would otherwise approve archives the running build refuses. +LIVE_FLYSIM=/opt/fly/current/flysim +[ "$(id -u)" -eq 0 ] || LIVE_FLYSIM="${FLY_LOOP_RESET_TEST_LIVE_FLYSIM:-$FLYSIM}" +same_build() { + [ -f "$LIVE_FLYSIM" ] \ + && [ "$(sha256sum < "$FLYSIM" | cut -d' ' -f1)" = "$(sha256sum < "$LIVE_FLYSIM" | cut -d' ' -f1)" ] +} + build_compat="" -if [ -x "$FLYSIM" ] && [ "$(stat -c %u "$FLYSIM")" = "$(id -u)" ]; then +if [ -x "$FLYSIM" ] && [ "$(stat -c %u "$FLYSIM")" = "$(id -u)" ] && same_build; then build_compat="$(clean_env timeout 90 "$FLYSIM" --print-compatibility 2>/dev/null | tail -n1 || true)" fi if [ -z "$build_compat" ]; then - log "no compatibility string from $FLYSIM (missing, not owned by $(id -un), or failed); no rung is restorable" + log "no compatibility string from $FLYSIM (missing, not owned by $(id -un), not the running build, or failed); no rung is restorable" [ "$LIST" -eq 1 ] && exit 0 exit 3 fi @@ -141,11 +150,18 @@ trap 'exit 143' TERM INT HUP log "pausing ${WATCHDOG_TIMER} and stopping ${SERVICE} to reset to rung ${RANK}" systemctl stop "$WATCHDOG_TIMER" +# A oneshot probe mid-run is "activating", which `is-active` does not count; wait on ActiveState. +watchdog_idle() { + case "$(systemctl show -p ActiveState --value "$WATCHDOG_SERVICE")" in + inactive|failed) return 0 ;; + *) return 1 ;; + esac +} for _ in $(seq 1 60); do - systemctl is-active --quiet "$WATCHDOG_SERVICE" || break + watchdog_idle && break sleep 1 done -if systemctl is-active --quiet "$WATCHDOG_SERVICE"; then +if ! watchdog_idle; then log "${WATCHDOG_SERVICE} still running after 60 s; not resetting" exit 4 fi diff --git a/infra/docs/loop-recovery.md b/infra/docs/loop-recovery.md index 953353a..f06d878 100644 --- a/infra/docs/loop-recovery.md +++ b/infra/docs/loop-recovery.md @@ -15,7 +15,8 @@ a ladder one step per trap that outlives the previous step: - **Outlives** means the reports that confirm it again are at least 20 minutes after the step, so their 10-brain-minute window lies after it. The helper waits that long after every step. - **Budget:** at most two milestone resets per 24 hours. With the budget spent the step is a - restart, at most every three hours, until a reset is free again. + restart, at most every three hours, until a reset is free again. The same three-hour spacing applies + when a reset level finds no restorable rung and restarts instead. - **Starting over:** the ladder returns to level 0 when the fly reaches a new best rung, or after six hours with no suspected report. @@ -29,9 +30,12 @@ release env file so a deploy does not drop it. The step runs `/opt/fly/bin/fly-loop-reset ` through sudo (the one line in `config/fly-sudoers`). As root it: -- runs only the root-owned `/opt/fly/sbin/flysim` that `05-deploy.sh` installs from the release - tarball after checking it against the tarball's MANIFEST, never the fly-owned release tree; - without that copy no rung is restorable; +- runs only the root-owned `/opt/fly/sbin/flysim`, and only while it is byte-identical to the + flysim `/opt/fly/current` points to, that `05-deploy.sh` installs from the release + tarball after checking it against the tarball's MANIFEST (never the fly-owned release tree). + Without that copy (an infra-only deploy, or a manual rollback) no rung is restorable and the + ladder only restarts; after a release deploy, check as `fly` that + `sudo -n /opt/fly/bin/fly-loop-reset --list` names the current rung; - reads `fly.env` as `KEY=VALUE` data, never sources it, and ignores the caller's environment; - stops `fly-watchdog.timer` and waits for a running probe to finish, so the watchdog cannot start flysim on a half-rewritten store; stops flysim; runs `fly-reset-to-milestone`; and on every @@ -41,7 +45,9 @@ The reset copies both stores to `/srv/fly/state.reset-` first, as in the ru clears milestone archives above the rung. After each step the helper waits up to four minutes for `/status` to report `running` (and the target rank for a reset). Otherwise the step is recorded as failed and the ladder still climbs. The step is recorded before it runs, so a helper killed -mid-step has still climbed and spent the reset. +mid-step has still climbed and spent the reset. Do not stop `fly-loop-recover.service` during a +step: that kills the reset too, and the wrapper's exit trap starts flysim on whatever state is +left. ## On stream diff --git a/infra/tests/test_loop_recover.py b/infra/tests/test_loop_recover.py index 405135c..c3b15ff 100644 --- a/infra/tests/test_loop_recover.py +++ b/infra/tests/test_loop_recover.py @@ -267,7 +267,12 @@ class WrapperTests(unittest.TestCase): self.calls = root / "calls.log" (root / "state").mkdir() (root / "bin").mkdir() - self.script(root / "bin/systemctl", f'echo "systemctl $*" >> {self.calls}; [ "$1" != is-active ] || exit 3') + # The watchdog probe is mid-run ("activating") for the first two looks. + self.script(root / "bin/systemctl", f'''echo "systemctl $*" >> {self.calls} +if [ "$1" = show ]; then + n=$(grep -c "^systemctl show" {self.calls}) + if [ "$n" -le 2 ]; then echo activating; else echo inactive; fi +fi''') self.script(root / "flysim", f'[ "$1" = --print-compatibility ] && echo "{BUILD}"') self.script(root / "reset", f'echo "reset $* FLY_BIN=$FLY_BIN" >> {self.calls}') self.env_file = root / "fly.env" @@ -288,6 +293,7 @@ class WrapperTests(unittest.TestCase): env = {"PATH": f"{self.root}/bin:/usr/bin:/bin", "FLY_LOOP_RESET_TEST_STATE_DIR": str(self.root / "state"), "FLY_LOOP_RESET_TEST_ENV_FILE": str(self.env_file), "FLY_LOOP_RESET_TEST_FLYSIM": str(self.root / "flysim"), "FLY_LOOP_RESET_TEST_RESET_BIN": str(self.root / "reset")} + env.update(getattr(self, "extra_env", {})) return recover.subprocess.run([str(WRAPPER), *args], env=env, capture_output=True, text=True, timeout=60) def log(self): @@ -303,7 +309,9 @@ class WrapperTests(unittest.TestCase): done = self.run_wrapper("12") self.assertEqual(done.returncode, 0, done.stderr) self.assertEqual([line.split()[:3] for line in self.log().splitlines() if not line.startswith("systemctl is-active")], [ - ["systemctl", "stop", "fly-watchdog.timer"], ["systemctl", "stop", "flysim.service"], + ["systemctl", "stop", "fly-watchdog.timer"], + ["systemctl", "show", "-p"], ["systemctl", "show", "-p"], ["systemctl", "show", "-p"], ["systemctl", "show", "-p"], + ["systemctl", "stop", "flysim.service"], ["reset", "12", f"FLY_BIN={self.root}/flysim"], ["systemctl", "start", "flysim.service"], ["systemctl", "start", "fly-watchdog.timer"]]) @@ -324,6 +332,20 @@ class WrapperTests(unittest.TestCase): self.assertEqual(self.run_wrapper("12").returncode, 3) self.assertEqual(self.log(), "") + def test_a_root_copy_that_is_not_the_running_build_restores_nothing(self): + other = self.root / "live-flysim" + self.script(other, "echo other build") + self.extra_env = {"FLY_LOOP_RESET_TEST_LIVE_FLYSIM": str(other)} + self.assertEqual(self.run_wrapper("--list").stdout, "") + self.assertEqual(self.run_wrapper("12").returncode, 3) + + def test_a_watchdog_probe_that_never_finishes_blocks_the_reset(self): + self.script(self.root / "bin/systemctl", f'echo "systemctl $*" >> {self.calls}; [ "$1" != show ] || echo activating') + self.script(self.root / "bin/sleep", "exit 0") + self.assertEqual(self.run_wrapper("12").returncode, 4) + self.assertNotIn("systemctl stop flysim.service", self.log()) + self.assertIn("systemctl start fly-watchdog.timer", self.log()) + def test_bad_arguments_are_refused(self): for args in ((), ("--check", "12"), ("12", "13"), ("../12",), ("123",), ("-1",)): self.assertEqual(self.run_wrapper(*args).returncode, 2, args)