#634 devpulse: watchdog stall detector — kill false-positives on long tool calls + surface stall live.

Two rough edges on the JSONL stall detector, both hardened in one pass on my own module (apps/handlers/watchdog/agent.py).

PART 1 (false-positive): _has_jsonl_activity inferred liveness purely from JSONL file-size growth over the 120s window. An agent doing ONE genuinely long operation (big Read, long Bash, heavy compute) writes no new JSONL lines for that span -> read as idle -> STALLED fires WHILE the agent is actively working. Fix: watch_agent now also treats an in-flight tool_use as activity. While a tool runs, the assistant's tool_use is the last transcript entry; new _last_entry_is_inflight_tool() tail-reads the newest .jsonl and detects it (fully defensive -> False on any parse/shape drift, degrading to size-based). LIVE-PROVEN against real Claude Code transcripts: sampled my own session across a 10s in-flight bash -> tool_use line is written at tool START and persists the whole call (the sub-second flush lag is irrelevant at the 120s horizon).

PART 2 (invisible stall): the stall only hit _stderr()+logger. The Monitor tool that arms the watchdog turns each STDOUT line into a live event but only captures stderr to a file (never surfaced) -> devpulse never saw the stall until the 600s timeout. Fix: new _stdout_event() emits the stall (+ a long-running-tool advisory for a possibly-hung tool, + a resumed signal) to stdout so Monitor relays it live; the verbose trail stays on stderr+logger.

Stall logic extracted into a StallTracker class (kills deep-nesting). +9 tests (unit + full-loop stdout proofs + real-transcript schema check); 142 watchdog tests green, seedgo audit 100%, no type errors.

Rides PR#659 (issue-clearing campaign, no main-merge).
This commit is contained in:
AIOSAI
2026-07-10 03:57:19 -07:00
parent cded2993f5
commit c970c5761f
3 changed files with 388 additions and 26 deletions
+14
View File
@@ -33,6 +33,20 @@ PyPI version — not the changelog header.
### Fixed
- **Watchdog stall detector no longer false-fires on a long single tool call, and
a real stall now reaches devpulse live (issue #634).** Liveness was inferred
purely from JSONL file-size growth, so an agent doing one genuinely long
operation (big read, long-running Bash, heavy compute) wrote no new lines for
the span and was misread as `STALLED` while actively working. `watch_agent` now
also treats an in-flight `tool_use` (the assistant's last transcript entry while
a tool runs) as activity — verified live against real Claude Code transcripts:
the `tool_use` line is written at tool *start* and persists for the whole call.
Part 2: the stall (and a new long-running-tool advisory, plus a resumed signal)
is emitted to **stdout** — which the Monitor-tool wrapper surfaces as a live
event — instead of only `stderr`+logger, which Monitor captures but never
relays. Stall logic extracted into a `StallTracker` for clarity; +9 tests
(142 green), devpulse audit 100%. (devpulse)
- **`aipass install` no longer hard-fails (exit 2, silently) when it can't create
global symlinks (issue #660 follow-up).** `setup.sh` runs under
`set -euo pipefail`; the #660 `safe_symlink` refactor returns `2` on `ln`
@@ -1,9 +1,9 @@
# =================== AIPass ====================
# Name: agent.py
# Description: Watchdog Agent Handler — block until dispatched agent exits
# Version: 1.0.0
# Version: 1.1.0
# Created: 2026-04-14
# Modified: 2026-04-14
# Modified: 2026-07-10
# =============================================
# Signal choice: ai_mail dispatch lock file polling.
@@ -44,6 +44,19 @@ def _stderr(msg: str) -> None:
sys.stderr.flush()
def _stdout_event(msg: str) -> None:
"""Write one flushed event line to stdout so a Monitor-tool wrapper turns it
into a live notification to devpulse.
The Monitor tool treats each stdout line as an event but only captures stderr
to a file (never surfaced). So mid-watch signals a caller must ACT on — a stall
or a possibly-hung tool — go here, while the verbose debug trail stays on
_stderr + logger. (#634 part 2.)
"""
sys.stdout.write(msg + "\n")
sys.stdout.flush()
def _find_repo_root(start: Path | None = None) -> Path | None:
"""Walk upward looking for AIPASS_REGISTRY.json. Returns None if not found."""
cur = (start or Path.cwd()).resolve()
@@ -231,6 +244,87 @@ def _has_jsonl_activity(projects_dir: Path, baseline: dict) -> bool:
return False
def _newest_jsonl(projects_dir: Path) -> Path | None:
"""Return the most-recently-modified .jsonl in projects_dir, or None."""
if not projects_dir.exists():
return None
try:
files = list(projects_dir.glob("*.jsonl"))
except OSError as exc:
logger.info("[watchdog.agent] newest jsonl glob failed: %s", exc)
return None
newest: Path | None = None
newest_mtime = -1.0
for f in files:
try:
mtime = f.stat().st_mtime
except OSError as exc:
logger.info("[watchdog.agent] newest jsonl stat failed for %s: %s", f.name, exc)
continue
if mtime > newest_mtime:
newest_mtime = mtime
newest = f
return newest
def _tail_last_line(path: Path, max_bytes: int = 1_000_000) -> str | None:
"""Read the last non-blank newline-delimited line of a file via a bounded tail
read. Reads at most ``max_bytes`` from the end so a multi-MB transcript stays
cheap; a single line longer than that decodes partially and simply fails to
parse downstream (→ treated as no in-flight tool)."""
try:
size = path.stat().st_size
with path.open("rb") as fh:
if size > max_bytes:
fh.seek(size - max_bytes)
chunk = fh.read()
except OSError as exc:
logger.info("[watchdog.agent] tail read failed for %s: %s", path.name, exc)
return None
if not chunk:
return None
text = chunk.decode("utf-8", errors="replace")
lines = [ln for ln in text.splitlines() if ln.strip()]
return lines[-1] if lines else None
def _last_entry_is_inflight_tool(projects_dir: Path) -> bool:
"""True if the newest JSONL's last event is an assistant message dispatching a
tool call — an in-flight ``tool_use`` awaiting its result.
While a tool runs (a big Read, a long Bash, heavy compute) the agent writes NO
new JSONL lines, so size-growth alone misreads that span as idle and false-fires
STALLED (#634 part 1). The last line being an assistant ``tool_use`` is the
precise signal that the agent is actively working, not stuck.
Best-effort: any read/parse/shape surprise returns False, degrading to the
size-based liveness check so the stall detector never crashes on a format drift.
"""
newest = _newest_jsonl(projects_dir)
if newest is None:
return False
last_line = _tail_last_line(newest)
if not last_line:
return False
try:
entry = json.loads(last_line)
except (json.JSONDecodeError, ValueError) as exc:
logger.info("[watchdog.agent] last jsonl entry unparseable in %s: %s", newest.name, exc)
return False
if not isinstance(entry, dict):
return False
# Schemas vary: role/content may sit under "message" or at the top level.
message = entry.get("message")
if not isinstance(message, dict):
message = entry
if message.get("role") != "assistant":
return False
content = message.get("content")
if not isinstance(content, list):
return False
return any(isinstance(block, dict) and block.get("type") == "tool_use" for block in content)
def _read_lock(lock_file: Path) -> dict | None:
"""Read lock file, return dict or None on miss/error."""
if not lock_file.exists():
@@ -316,6 +410,98 @@ def _classify_exit(
return ("completed_silent", reason, 0)
class StallTracker:
"""Per-watch JSONL liveness/stall state — one ``observe()`` call per poll tick.
Liveness signal = new JSONL lines (size growth) OR an in-flight tool call. A
long single tool call (big Read, long Bash, heavy compute) writes no new JSONL
lines while it runs but leaves an assistant ``tool_use`` as the last entry, so
it counts as activity instead of false-firing STALLED (#634 part 1).
Actionable signals — stall, long-running tool, resumed — go to stdout via
``_stdout_event`` so a Monitor-tool wrapper surfaces them to devpulse live
(#634 part 2). The verbose trail stays on ``_stderr`` + logger.
"""
STALL_THRESHOLD = 120.0
# A single tool call held in-flight this long is surfaced as a soft advisory
# (heavy op or a hung tool). Below the 600s default timeout so long watches
# get a mid-flight heads-up instead of waiting on the timeout.
LONG_TOOL_THRESHOLD = 300.0
def __init__(self, agent_id: str, jsonl_dir: Path, baseline: dict, now: float, pid: object) -> None:
self.agent_id = agent_id
self.jsonl_dir = jsonl_dir
self.baseline = baseline
self.pid = pid
self.last_activity_at = now
self.in_flight_since: float | None = None
self.stall_reported = False
self.long_tool_reported = False
def observe(self, now: float) -> None:
"""Evaluate one poll tick: reset on activity, else report a stall past threshold."""
size_grew = _has_jsonl_activity(self.jsonl_dir, self.baseline)
inflight = False if size_grew else _last_entry_is_inflight_tool(self.jsonl_dir)
if size_grew or inflight:
self._mark_active(now, size_grew, inflight)
elif not self.stall_reported and (now - self.last_activity_at) >= self.STALL_THRESHOLD:
self._report_stall(now)
def _mark_active(self, now: float, size_grew: bool, inflight: bool) -> None:
if size_grew:
self.baseline = _snapshot_jsonl_sizes(self.jsonl_dir)
self.last_activity_at = now
if self.stall_reported:
_stdout_event(f"[watchdog.resumed] {self.agent_id}: JSONL activity resumed — stall cleared")
_stderr(f"[watchdog.agent] {self.agent_id}: activity resumed")
logger.info("[watchdog.agent] activity resumed agent_id=%s", self.agent_id)
self.stall_reported = False
if inflight:
self._track_inflight(now)
else:
self.in_flight_since = None
self.long_tool_reported = False
def _track_inflight(self, now: float) -> None:
if self.in_flight_since is None:
self.in_flight_since = now
inflight_secs = int(now - self.in_flight_since)
if self.long_tool_reported or inflight_secs < self.LONG_TOOL_THRESHOLD:
return
_stdout_event(
f"[watchdog.longtool] {self.agent_id}: one tool call running {inflight_secs}s "
f"(PID {self.pid} alive) — likely a heavy op, but may be a hung tool. "
f"Check, or kill+resume if stuck."
)
logger.info(
"[watchdog.agent] long-running tool agent_id=%s inflight=%ss pid=%s",
self.agent_id,
inflight_secs,
self.pid,
)
self.long_tool_reported = True
def _report_stall(self, now: float) -> None:
idle_secs = int(now - self.last_activity_at)
_stdout_event(
f"[watchdog.stall] {self.agent_id}: STALLED — no JSONL activity for {idle_secs}s "
f"(PID {self.pid} alive, no in-flight tool call). Agent may be stuck — "
f"check, or kill+resume."
)
_stderr(
f"[watchdog.agent] {self.agent_id}: STALLED — no JSONL activity for {idle_secs}s "
f"(PID {self.pid} still alive)"
)
logger.info(
"[watchdog.agent] stall detected agent_id=%s idle=%ss pid=%s",
self.agent_id,
idle_secs,
self.pid,
)
self.stall_reported = True
def watch_agent(
agent_id: str,
timeout_seconds: int = 600,
@@ -383,10 +569,7 @@ def watch_agent(
_stderr(f"[watchdog.agent] {agent_id}: lock present, monitor PID={initial_pid}")
jsonl_dir = _get_jsonl_projects_dir(branch_path)
jsonl_baseline = _snapshot_jsonl_sizes(jsonl_dir)
last_activity_at = time.monotonic()
stall_reported = False
stall_threshold = 120.0
tracker = StallTracker(agent_id, jsonl_dir, _snapshot_jsonl_sizes(jsonl_dir), time.monotonic(), initial_pid)
while True:
elapsed = time.monotonic() - started_at
@@ -441,26 +624,7 @@ def watch_agent(
"handle": handle,
}
if _has_jsonl_activity(jsonl_dir, jsonl_baseline):
jsonl_baseline = _snapshot_jsonl_sizes(jsonl_dir)
last_activity_at = time.monotonic()
if stall_reported:
_stderr(f"[watchdog.agent] {agent_id}: activity resumed")
logger.info("[watchdog.agent] activity resumed agent_id=%s", agent_id)
stall_reported = False
elif not stall_reported and (time.monotonic() - last_activity_at) >= stall_threshold:
idle_secs = int(time.monotonic() - last_activity_at)
_stderr(
f"[watchdog.agent] {agent_id}: STALLED — no JSONL activity for {idle_secs}s "
f"(PID {initial_pid} still alive)"
)
logger.info(
"[watchdog.agent] stall detected agent_id=%s idle=%ss pid=%s",
agent_id,
idle_secs,
initial_pid,
)
stall_reported = True
tracker.observe(time.monotonic())
time.sleep(poll_interval)
finally:
@@ -207,6 +207,190 @@ def test_watch_agent_return_keys():
assert set(result.keys()) == expected
# ─────────────────────────────────────────────────────────────────────────────
# #634 — in-flight tool detection + stall surfaced to stdout
# ─────────────────────────────────────────────────────────────────────────────
def _write_jsonl(projects_dir: Path, *lines: dict, name: str = "session.jsonl") -> Path:
"""Write JSONL entries (one dict per line) into a projects dir."""
projects_dir.mkdir(parents=True, exist_ok=True)
f = projects_dir / name
f.write_text("".join(json.dumps(ln) + "\n" for ln in lines), encoding="utf-8")
return f
def test_last_entry_is_inflight_tool_true_for_assistant_tool_use(tmp_path):
"""Last line = assistant message with a tool_use block → in-flight tool call."""
proj = tmp_path / "proj"
_write_jsonl(
proj,
{"type": "user", "message": {"role": "user", "content": [{"type": "text", "text": "go"}]}},
{
"type": "assistant",
"message": {"role": "assistant", "content": [{"type": "tool_use", "id": "a", "name": "Bash", "input": {}}]},
},
)
assert agent_handler._last_entry_is_inflight_tool(proj) is True
def test_last_entry_is_inflight_tool_false_for_text_and_results(tmp_path):
"""Assistant text-only, a returned tool_result, malformed, and empty all → False."""
proj = tmp_path / "proj"
_write_jsonl(
proj, {"type": "assistant", "message": {"role": "assistant", "content": [{"type": "text", "text": "done"}]}}
)
assert agent_handler._last_entry_is_inflight_tool(proj) is False
_write_jsonl(
proj,
{
"type": "user",
"message": {"role": "user", "content": [{"type": "tool_result", "tool_use_id": "a", "content": "ok"}]},
},
)
assert agent_handler._last_entry_is_inflight_tool(proj) is False
(proj / "session.jsonl").write_text("{not valid json\n", encoding="utf-8")
assert agent_handler._last_entry_is_inflight_tool(proj) is False
# Nonexistent dir and empty dir → False.
assert agent_handler._last_entry_is_inflight_tool(tmp_path / "nope") is False
(tmp_path / "empty").mkdir()
assert agent_handler._last_entry_is_inflight_tool(tmp_path / "empty") is False
def test_last_entry_is_inflight_tool_picks_newest_file(tmp_path):
"""With multiple JSONLs, only the most-recently-modified one decides."""
proj = tmp_path / "proj"
old = _write_jsonl(
proj,
{"message": {"role": "assistant", "content": [{"type": "tool_use", "name": "Bash"}]}},
name="old.jsonl",
)
new = _write_jsonl(
proj,
{"message": {"role": "assistant", "content": [{"type": "text", "text": "hi"}]}},
name="new.jsonl",
)
os.utime(old, (1, 1))
os.utime(new, (2, 2))
assert agent_handler._last_entry_is_inflight_tool(proj) is False # newest = text-only
os.utime(old, (3, 3)) # old is now newest and holds the tool_use
assert agent_handler._last_entry_is_inflight_tool(proj) is True
def test_stalltracker_reports_stall_after_threshold(monkeypatch, capsys):
"""No activity past STALL_THRESHOLD → a [watchdog.stall] line on stdout."""
monkeypatch.setattr(agent_handler, "_has_jsonl_activity", lambda *a, **kw: False)
monkeypatch.setattr(agent_handler, "_last_entry_is_inflight_tool", lambda *a, **kw: False)
t = agent_handler.StallTracker("@x", Path("/nope"), {}, now=0.0, pid=123)
t.observe(now=60.0) # below threshold
assert "[watchdog.stall]" not in capsys.readouterr().out
assert t.stall_reported is False
t.observe(now=agent_handler.StallTracker.STALL_THRESHOLD) # at threshold
assert "[watchdog.stall]" in capsys.readouterr().out
assert t.stall_reported is True
def test_stalltracker_inflight_tool_prevents_stall(monkeypatch, capsys):
"""An in-flight tool call resets the idle timer every tick → never a stall."""
monkeypatch.setattr(agent_handler, "_has_jsonl_activity", lambda *a, **kw: False)
monkeypatch.setattr(agent_handler, "_last_entry_is_inflight_tool", lambda *a, **kw: True)
t = agent_handler.StallTracker("@x", Path("/nope"), {}, now=0.0, pid=123)
for now in (60.0, 120.0, 180.0, 240.0):
t.observe(now=now)
out = capsys.readouterr().out
assert "[watchdog.stall]" not in out
assert t.stall_reported is False
def test_stalltracker_long_tool_advisory(monkeypatch, capsys):
"""One tool call held in-flight past LONG_TOOL_THRESHOLD → advisory, not a stall."""
monkeypatch.setattr(agent_handler, "_has_jsonl_activity", lambda *a, **kw: False)
monkeypatch.setattr(agent_handler, "_last_entry_is_inflight_tool", lambda *a, **kw: True)
t = agent_handler.StallTracker("@x", Path("/nope"), {}, now=0.0, pid=123)
t.observe(now=0.0) # first in-flight tick → anchors in_flight_since
t.observe(now=agent_handler.StallTracker.LONG_TOOL_THRESHOLD)
out = capsys.readouterr().out
assert "[watchdog.longtool]" in out
assert "[watchdog.stall]" not in out
assert t.long_tool_reported is True
def test_stalltracker_resume_clears_stall(monkeypatch, capsys):
"""After a stall, real activity emits [watchdog.resumed] and clears the flag."""
signals = {"size": False}
monkeypatch.setattr(agent_handler, "_has_jsonl_activity", lambda *a, **kw: signals["size"])
monkeypatch.setattr(agent_handler, "_last_entry_is_inflight_tool", lambda *a, **kw: False)
monkeypatch.setattr(agent_handler, "_snapshot_jsonl_sizes", lambda *a, **kw: {})
t = agent_handler.StallTracker("@x", Path("/nope"), {}, now=0.0, pid=123)
t.observe(now=agent_handler.StallTracker.STALL_THRESHOLD) # stall
assert "[watchdog.stall]" in capsys.readouterr().out
assert t.stall_reported is True
signals["size"] = True # activity resumes
t.observe(now=agent_handler.StallTracker.STALL_THRESHOLD + 5)
out = capsys.readouterr().out
assert "[watchdog.resumed]" in out
assert t.stall_reported is False
def _fake_clock_sleep(agent_module, monkeypatch, lock_file, unlink_at=200.0, step=60.0):
"""Patch monotonic + sleep with a fake clock that advances `step`s per sleep
and unlinks the dispatch lock once the clock passes `unlink_at` (loop exit)."""
clock = {"t": 0.0}
monkeypatch.setattr(agent_module.time, "monotonic", lambda: clock["t"])
def fake_sleep(_seconds):
"""Advance the fake clock and drop the lock once past unlink_at (loop exit)."""
clock["t"] += step
if clock["t"] >= unlink_at:
lock_file.unlink(missing_ok=True)
monkeypatch.setattr(agent_module.time, "sleep", fake_sleep)
def test_watch_agent_surfaces_stall_to_stdout(monkeypatch, tmp_path, capsys):
"""End-to-end: a genuine stall reaches STDOUT so the Monitor wrapper relays it."""
branch_path = _build_fake_branch(tmp_path)
lock_file = _write_lock(branch_path, pid=os.getpid())
monkeypatch.setattr(agent_handler, "_find_repo_root", lambda *a, **kw: tmp_path)
monkeypatch.setattr(agent_handler, "_pid_alive", lambda pid: True)
monkeypatch.setattr(agent_handler, "_has_jsonl_activity", lambda *a, **kw: False)
monkeypatch.setattr(agent_handler, "_last_entry_is_inflight_tool", lambda *a, **kw: False)
_fake_clock_sleep(agent_handler, monkeypatch, lock_file)
result = agent_handler.watch_agent("@fakebranch", timeout_seconds=100000, poll_interval=0.01)
out = capsys.readouterr().out
assert "[watchdog.stall]" in out
assert result["woke"] is True
def test_watch_agent_inflight_tool_no_false_stall(monkeypatch, tmp_path, capsys):
"""End-to-end: a long in-flight tool call must NOT surface a stall (#634 part 1)."""
branch_path = _build_fake_branch(tmp_path)
lock_file = _write_lock(branch_path, pid=os.getpid())
monkeypatch.setattr(agent_handler, "_find_repo_root", lambda *a, **kw: tmp_path)
monkeypatch.setattr(agent_handler, "_pid_alive", lambda pid: True)
monkeypatch.setattr(agent_handler, "_has_jsonl_activity", lambda *a, **kw: False)
monkeypatch.setattr(agent_handler, "_last_entry_is_inflight_tool", lambda *a, **kw: True)
_fake_clock_sleep(agent_handler, monkeypatch, lock_file)
result = agent_handler.watch_agent("@fakebranch", timeout_seconds=100000, poll_interval=0.01)
out = capsys.readouterr().out
assert "[watchdog.stall]" not in out
assert result["woke"] is True
# ─────────────────────────────────────────────────────────────────────────────
# Integration tests (require live ai_mail dispatch — skipped by default)
# ─────────────────────────────────────────────────────────────────────────────