From b582da8c5af7ea1e040a1b47100efaf37961f0a1 Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Thu, 6 Aug 2026 16:39:49 +0200 Subject: [PATCH] feat(bin): runner heartbeat, tiered attention windows, threaded whisper poll --- README.md | 10 +++++-- bin/train | 88 +++++++++++++++++++++++++++++++++++++++++++++---------- 2 files changed, 79 insertions(+), 19 deletions(-) diff --git a/README.md b/README.md index d105f94..a3b9fff 100644 --- a/README.md +++ b/README.md @@ -32,9 +32,13 @@ rerun after adding a tool or config. `--whispers [prefix]` flag splits the screen 50/50 with a live reeds feed (default prefix `train`, stacked instead of side-by-side on narrow windows). In tmux, `prefix+M` pops up that split view, which shows every - incomplete train as its own section (falling back to the latest train when - all are done, idling when there are none) and closes itself when the shown - trains complete on its watch; `q` or `Esc` dismisses it. A popup swallows all + attended incomplete train as its own section (falling back to the latest + train when all are done, idling when there are none) and closes itself when + the shown trains complete on its watch; `q` or `Esc` dismisses it. Attended + means: a fresh runner heartbeat (`train run` touches one every 15s), or + marks within 4h, stretched to 24h while a step waits on a human. A section + whose live step has a stale heartbeat shows "runner heartbeat lost". A + popup swallows all keys including the tmux prefix, so if one ever refuses to die, run `tmux display-popup -C` from any other terminal to close it. - `glab-merge `: waits for a GitLab MR to become mergeable, diff --git a/bin/train b/bin/train index 3520df3..33fd86e 100755 --- a/bin/train +++ b/bin/train @@ -29,6 +29,7 @@ import signal import subprocess import sys import tempfile +import threading import time import urllib.request from collections import deque @@ -59,7 +60,10 @@ class Status(StrEnum): FINISHED = {Status.DONE, Status.SKIPPED} STATUS_VALUES = {s.value for s in Status} REEDS_KIND = {Status.DONE: "done", Status.NEEDS_USER: "needs-user"} -WATCH_FRESH_SECS = 24 * 3600 +HEARTBEAT_SECS = 15 +HEARTBEAT_FRESH_SECS = 60 +WATCH_FRESH_SECS = 4 * 3600 +NEEDS_USER_FRESH_SECS = 24 * 3600 @dataclass(frozen=True) @@ -247,9 +251,19 @@ def run_step(command: str) -> tuple[int, str]: return proc.wait(), clean("".join(tail)[-500:]) +def start_heartbeat(dir_: Path) -> None: + def beat() -> None: + while True: + (dir_ / "heartbeat").touch() + time.sleep(HEARTBEAT_SECS) + + threading.Thread(target=beat, daemon=True).start() + + def cmd_run(id: str | None) -> int: train = resolve(id) signal.signal(signal.SIGTERM, lambda *_: sys.exit(143)) + start_heartbeat(train.dir) for step in train.steps(): current = train.state().get(step.key) if current is None: @@ -303,24 +317,39 @@ class WhisperFeed: self.cursor = 0 self.rows: deque[dict] = deque(maxlen=200) self.state = "init" + self.lock = threading.Lock() + + def start(self, interval: float) -> None: + def loop() -> None: + while True: + self.poll() + time.sleep(max(interval, 2)) + + threading.Thread(target=loop, daemon=True).start() + + def snapshot(self) -> tuple[list[dict], str]: + with self.lock: + return list(self.rows), self.state def poll(self) -> None: base = os.environ.get("REEDS_URL", "http://localhost:7333") try: - # two pages per tick keeps a cold catch-up from stalling the UI - for _ in range(2): + for _ in range(10): url = f"{base}/t/{self.prefix}?since={self.cursor}&limit=500" - with urllib.request.urlopen(url, timeout=1) as resp: + with urllib.request.urlopen(url, timeout=2) as resp: page = json.load(resp) - self.rows.extend(page["whispers"]) + with self.lock: + self.rows.extend(page["whispers"]) + self.state = "live" self.cursor = page["next_since"] if not page["more"]: break - self.state = "live" except OSError: - self.state = "down" + with self.lock: + self.state = "down" except (ValueError, KeyError): - self.state = "err" + with self.lock: + self.state = "err" def region(scr: "curses.window", y0: int, x0: int, height: int, width: int) -> Put: @@ -342,19 +371,43 @@ class View: train: Train steps: list[Step] marks: list[Mark] + runner_gone: bool @property def complete(self) -> bool: return bool(self.steps) and all(m.status in FINISHED for m in self.marks) +def heartbeat_age(train: Train) -> float | None: + try: + return time.time() - (train.dir / "heartbeat").stat().st_mtime + except FileNotFoundError: + return None + + def load_view(train: Train) -> View | None: try: steps = train.steps() state = train.state() except (FileNotFoundError, ParseError): return None - return View(train, steps, [state.get(s.key, Mark(Status.PENDING, "")) for s in steps]) + marks = [state.get(s.key, Mark(Status.PENDING, "")) for s in steps] + age = heartbeat_age(train) + runner_gone = ( + any(m.status is Status.LIVE for m in marks) + and age is not None + and age > HEARTBEAT_FRESH_SECS + ) + return View(train, steps, marks, runner_gone) + + +def attended(view: View, now: float) -> bool: + age = heartbeat_age(view.train) + if age is not None and age < HEARTBEAT_FRESH_SECS: + return True + idle = now - view.train.dir.stat().st_mtime + waiting = any(m.status is Status.NEEDS_USER for m in view.marks) + return idle < (NEEDS_USER_FRESH_SECS if waiting else WATCH_FRESH_SECS) def all_trains() -> list[Train]: @@ -387,6 +440,9 @@ def draw_section(put: Put, y: int, view: View, palette: dict[str, int]) -> int: put(y, 9 + len(step.label) + 2, f"({mark.detail})", curses.A_DIM) y += 1 counts = {s: sum(m.status is s for m in view.marks) for s in Status} + if view.runner_gone: + put(y, 0, "runner heartbeat lost", palette["red"]) + y += 1 if counts[Status.FAIL]: put(y, 0, f"train stalled: {counts[Status.FAIL]} step(s) broken", palette["red"] | curses.A_BOLD) @@ -406,15 +462,16 @@ WHISPER_KIND_COLOR = {"done": "green", "needs-user": "cyan", "status": "dim"} def draw_whispers(put: Put, feed: WhisperFeed, palette: dict[str, int], height: int) -> None: + rows, state = feed.snapshot() marker, color = { "live": ("[live]", "green"), "down": ("[down]", "red"), "err": ("[err] ", "yellow"), - }.get(feed.state, ("[....]", "dim")) + }.get(state, ("[....]", "dim")) put(0, 0, f"REEDS {feed.prefix}", curses.A_BOLD) put(0, 8 + len(feed.prefix), marker, palette[color]) put(1, 0, "-" * 62, 0) - visible = list(feed.rows)[-(height - 2):] + visible = rows[-(height - 2):] for i, row in enumerate(visible): stamp = time.strftime("%H:%M:%S", time.localtime(row.get("ts", 0) / 1000)) kind = str(row.get("kind", "")) @@ -450,6 +507,8 @@ def watch_loop(scr: "curses.window", fixed_id: str | None, interval: float, pass scr.timeout(int(interval * 1000)) feed = WhisperFeed(whispers) if whispers else None + if feed: + feed.start(interval) # Auto-close (exit 0) only when every shown train is complete and at least # one of them completed on our watch; already-complete trains, or no @@ -460,12 +519,9 @@ def watch_loop(scr: "curses.window", fixed_id: str | None, interval: float, views = [v for v in [load_view(Train(fixed_id))] if v] else: views = [v for t in all_trains() if (v := load_view(t))] - fresh = time.time() - WATCH_FRESH_SECS - active = [v for v in views - if not v.complete and v.train.dir.stat().st_mtime > fresh] + now = time.time() + active = [v for v in views if not v.complete and attended(v, now)] views = active or views[:1] - if feed: - feed.poll() def draw_trains(put: Put) -> None: if not views: -- 2.51.2