Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 65 additions & 6 deletions bin/backends/herdr-eventwait.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,14 +35,61 @@
A non-zero exit tells the bash caller to fall back to plain polling for this
cycle (the permanent fail-closed backstop), never to go silent.
"""
import errno
import io
import json
import os
import select
import socket
import sys
import time

CONNECT_TIMEOUT = 5.0
ACK_TIMEOUT = 5.0
RECV_CHUNK = 65536
# Per-write select wait when the stdout pipe is full. Small so the overall
# deadline stays responsive while a slow consumer (the bash caller's
# line-by-line drain) frees pipe space.
WRITE_POLL = 0.25


def _emit_bounded(data, deadline):
"""Write ALL of data to stdout, honoring an absolute monotonic deadline.

stdout is set non-blocking when it is a real pipe, so a full pipe raises
BlockingIOError instead of parking the process inside the kernel write
syscall - the failure mode the 2026-08-21 quiet-fleet incident exposed: a
subscribed backlog replay can fill the pipe faster than the bash caller
drains it, and a blocking write then outlives the wait budget
indefinitely, starving the watcher's liveness beacon. A stdout without a
real file descriptor (tests capturing into StringIO) falls back to a
single plain write, which cannot block. Returns True when everything was
written, False when the budget expired first (the caller ends its wait;
unwritten lines are a safe drop because the poll loop is the permanent
backstop).
"""
try:
fd = sys.stdout.fileno()
except (AttributeError, ValueError, OSError, io.UnsupportedOperation):
sys.stdout.write(data.decode("utf-8", "replace"))
return True
while data:
remaining = deadline - time.monotonic()
if remaining <= 0:
return False
_, writable, _ = select.select([], [fd], [], min(remaining, WRITE_POLL))
if not writable:
continue
try:
written = os.write(fd, data)
except BlockingIOError:
continue
except OSError as exc:
if exc.errno == errno.EINTR:
continue
raise
data = data[written:]
return True


def _read_line(sock, buf, deadline):
Expand Down Expand Up @@ -121,11 +168,23 @@ def main(argv):
if result.get("type") != "subscription_started":
return 3

sys.stdout.write("@subscribed\n")
sys.stdout.flush()

# Stream projected events until the deadline or the server closes.
# Bounded stdout writes for the whole streaming phase (see _emit_bounded):
# a blocking write on a full pipe could otherwise outlive the deadline.
try:
os.set_blocking(sys.stdout.fileno(), False)
except OSError:
pass
if not _emit_bounded(b"@subscribed\n", deadline):
return 0

# Stream projected events until the deadline or the server closes. The
# deadline is re-checked on EVERY iteration and every stdout write is
# budget-bounded, so a stream that stays saturated (a pending-backlog
# replay, a fast-flapping agent) can never make this loop outlive the
# wait budget - the watcher's heartbeat depends on the wait returning.
while True:
if time.monotonic() >= deadline:
return 0
line, buf, outcome = _read_line(sock, buf, deadline)
if line is None:
return 0 if outcome == "timeout" else 4
Expand All @@ -142,8 +201,8 @@ def main(argv):
_clean(data.get("agent_status") or ""),
_clean(data.get("agent") or ""),
)
sys.stdout.write("\t".join(fields) + "\n")
sys.stdout.flush()
if not _emit_bounded(("\t".join(fields) + "\n").encode("utf-8"), deadline):
return 0


if __name__ == "__main__":
Expand Down
54 changes: 40 additions & 14 deletions bin/backends/herdr.sh
Original file line number Diff line number Diff line change
Expand Up @@ -3132,7 +3132,7 @@ fm_backend_herdr_list_live() { # <session>
# --- native event push: pane.agent_status_changed subscriber -----------------
#
# The push half of the immediate blocked-state escalation (AGENTS.md section 8,
# docs/herdr-backend.md "Native pane.agent_status_changed push escalation").
# docs/herdr-backend.md "Push events and polling fallback").
# fm_backend_herdr_wait_transition is the watcher's bounded wait primitive for
# herdr homes: instead of a blind sleep, it blocks on herdr's native event
# stream and returns the instant a subscribed pane transitions to `blocked`, so
Expand Down Expand Up @@ -3267,8 +3267,18 @@ fm_backend_herdr_clear_transition() { # <state_dir> <window>
# fm_backend_herdr_wait_transition: the bounded event wait. Blocks up to
# <timeout_secs> for one of <pane_window...> ("<session>:<pane_id>") to reach a
# fresh `blocked` edge, then prints the normalized record and returns 0.
# Returns 1 on a clean timeout (the reader ran the full budget, no fresh
# actionable edge - the caller has effectively already slept and just continues)
# Every enforced side of the wait is bounded to the budget (reader lifetime,
# subscription ack, and stream drain): the reader enforces the deadline on
# every stream iteration and bounds every stdout write, and this side bounds
# the ack and drain with `read -t` and stops the reader when the budget
# expires - a saturated or backlogged event stream can never make the call
# outlive its poll budget, because an unbounded wait here starves the
# watcher's liveness beacon (2026-08-21 quiet-fleet crash-loop,
# negotiation-os mate home). The level reconcile runs inside the budget
# window but is bounded only by its per-window `herdr agent get` CLI reads,
# which carry no deadline of their own.
# Returns 1 on a clean timeout (the budget elapsed with no fresh actionable
# edge - the caller has effectively already slept and just continues)
# and 2 when the event path is unusable (not capable, socket unresolved, reader
# failed to run/subscribe - the caller sleeps the budget itself, the fail-closed
# backstop). See the header block above for the full contract.
Expand Down Expand Up @@ -3312,6 +3322,21 @@ fm_backend_herdr_wait_transition() { # <session> <timeout_secs> <state_dir> <pa
rm -rf "$fifo_dir" 2>/dev/null || true
return 2
fi

# Overall budget for this whole wait, measured from the reader spawn: the
# reader enforces the same deadline internally (every iteration and every
# stdout write), the subscription-ack and drain reads below use it as their
# `read -t` bound, and the reader is stopped the moment it expires - the
# enforced sides of the wait (reader lifetime, ack, drain) can never outlive
# the caller's poll cadence. The level reconcile between them consumes the
# budget but its per-window `herdr agent get` reads carry no deadline of
# their own. Unbounded waits here starve the watcher's liveness beacon
# (2026-08-21 quiet-fleet crash-loop: a saturated or backlogged event stream
# kept the drain alive past the 300s heartbeat grace while the watcher sat
# healthy-but-frozen in this call).
local wait_started wait_budget drain_remaining
wait_started=$(date +%s)
wait_budget=$(( timeout + 2 ))
"${reader[@]}" "$sock" "$timeout" "${pane_ids[@]}" > "$fifo" 2>/dev/null &
reader_pid=$!
if ! exec 9< "$fifo"; then
Expand All @@ -3320,7 +3345,7 @@ fm_backend_herdr_wait_transition() { # <session> <timeout_secs> <state_dir> <pa
rm -rf "$fifo_dir" 2>/dev/null || true
return 2
fi
if ! IFS= read -r -u 9 line || [ "$line" != "@subscribed" ]; then
if ! IFS= read -r -t "$wait_budget" -u 9 line || [ "$line" != "@subscribed" ]; then
rc=2
fi

Expand All @@ -3345,14 +3370,20 @@ fm_backend_herdr_wait_transition() { # <session> <timeout_secs> <state_dir> <pa
done
fi

# Drain stream edges until a fresh blocked edge or the timeout. The reader is
# a subprocess of this call (NOT a second watcher), and is killed the instant
# a blocked edge is found.
# Drain stream edges until a fresh blocked edge or the overall budget. The
# reader is a subprocess of this call (NOT a second watcher), and is stopped
# unconditionally below the moment the drain ends - by a blocked edge, the
# reader's own deadline exit, or the budget expiring under a stream that
# saturates the drain (the poll loop stays the fail-closed backstop for any
# edges left undelivered).
# Split each raw projected line (pane_id\tworkspace_id\tagent_status\tagent)
# with `cut`, NOT `IFS=$'\t' read`: a tab is IFS-whitespace, so `read` would
# collapse an empty middle field (e.g. an absent workspace_id) and shift the
# status into the wrong column. `cut` preserves empty fields.
while [ "$rc" -eq 1 ] && IFS= read -r line <&9; do
while [ "$rc" -eq 1 ]; do
drain_remaining=$(( wait_budget - ( $(date +%s) - wait_started ) ))
[ "$drain_remaining" -gt 0 ] || break
IFS= read -r -t "$drain_remaining" line <&9 || break
[ -n "$line" ] || continue
pane_id=$(printf '%s' "$line" | cut -f1)
ws=$(printf '%s' "$line" | cut -f2)
Expand All @@ -3366,12 +3397,7 @@ fm_backend_herdr_wait_transition() { # <session> <timeout_secs> <state_dir> <pa
break
fi
done
if [ "$rc" -eq 0 ]; then
kill "$reader_pid" 2>/dev/null || true
fi
if [ "$rc" -eq 2 ]; then
kill "$reader_pid" 2>/dev/null || true
fi
kill "$reader_pid" 2>/dev/null || true
# No actionable edge: distinguish a clean full-budget wait (reader exit 0 ->
# return 1, caller already waited) from a reader error (connect/subscribe
# failure, exit non-zero -> return 2, caller sleeps and counts toward the
Expand Down
4 changes: 2 additions & 2 deletions bin/fm-backend.sh
Original file line number Diff line number Diff line change
Expand Up @@ -913,8 +913,8 @@ fm_backend_agent_alive() { # <backend> <target>
# and for those backends replaces its blind `sleep POLL` with a bounded wait on
# fm_backend_wait_transition. Every push-capable backend reuses the shared
# normalized-transition shape and policy table (bin/fm-transition-lib.sh); today
# only herdr implements the surface (docs/herdr-backend.md "Native
# pane.agent_status_changed push escalation"). A backend with no native push
# only herdr implements the surface (docs/herdr-backend.md "Push events and
# polling fallback"). A backend with no native push
# reports has-push false and returns 2 from the dispatchers below, so the
# watcher falls back to its poll loop - the permanent fail-closed backstop.

Expand Down
6 changes: 5 additions & 1 deletion docs/herdr-backend.md
Original file line number Diff line number Diff line change
Expand Up @@ -281,8 +281,12 @@ The watcher maps the pane back to the task and skips secondmate endpoints, decla
The push path only shortens latency.
Polling runs every cycle and remains the permanent fallback when protocol 16, the event schema, Python, connection, subscription, or repeated reader execution is unavailable.
There is still one watcher process; the event reader is a bounded child of that watcher.
Every enforced side of the wait - reader lifetime, subscription ack, and stream drain - is bounded by the caller's poll budget: the reader re-checks its deadline on each stream iteration and bounds every stdout write (a saturated backlog replay can otherwise park the process inside a blocking pipe write), and the bash side bounds its ack and drain reads with `read -t` and stops its reader the moment the budget expires under a stream that outruns it.
The mid-wait level reconcile consumes the budget but is bounded only by its per-window `herdr agent get` CLI reads, which carry no deadline of their own.
Edges dropped at the deadline are safe: the poll loop is the permanent backstop.
An unbounded wait here starved the watcher's liveness beacon and crash-looped the away daemon's restarts on quiet fleets (2026-08-21 incident, negotiation-os mate home).

`tests/fm-backend-herdr-eventwait-smoke.test.sh`, `tests/fm-transition-lib.test.sh`, and `tests/fm-supervision-events.test.sh` cover capability, subscribe-then-reconcile ordering, dedupe, exemptions, and polling fallback.
`tests/fm-backend-herdr-eventwait-smoke.test.sh`, `tests/fm-backend-herdr-eventwait.test.py` (including the saturated-stream budget regression), `tests/fm-backend-herdr.test.sh` (including the watcher-level saturated-wait beat and restart-refuse regressions), `tests/fm-transition-lib.test.sh`, and `tests/fm-supervision-events.test.sh` cover capability, subscribe-then-reconcile ordering, dedupe, exemptions, budget bounds, watcher-level beat continuity, and polling fallback.

## Away-mode supervisor support

Expand Down
3 changes: 3 additions & 0 deletions docs/watcher-continuity.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,8 @@ The file is size-capped through `FM_WATCH_CYCLE_LOG_MAX_BYTES` and `FM_WATCH_CYC

The default 300-second grace is unchanged.
Only the watcher process touches `state/.last-watcher-beat`; no helper process can make a wedged watcher appear healthy.
The watcher's own waits must respect the beacon: every blocking wait inside a cycle (today the Herdr native event wait in `bin/fm-watch.sh`) is bounded by the poll budget on every enforced side, so a legitimately-blocked watcher still beats while the poll loop remains the delivery backstop.
A wait that outlives the grace makes a healthy watcher read dead, which the away daemon then cannot restart past the singleton lock (2026-08-21 quiet-fleet incident; budget bounds live in `docs/herdr-backend.md` "Push events and polling fallback").

## Regression coverage

Expand All @@ -84,6 +86,7 @@ The same suite covers ordinary same-process session replacement for `/new`, `/re
`tests/fm-watch-arm.test.sh` covers durable queue replay, real remote parent-replies ingestion into the authoritative status log, decision-only OPEN DECISIONS recovery, interrupted handling replay, generation-bound acknowledgement, a persistent live successor after recovery, a watcher close inside the handling window that must leave the printed acknowledgement valid, and the self-healing moved-generation acknowledgement that consumes its handled rows and names its remedy.
`tests/fm-watch-recovery-loop.test.sh` covers the once-per-generation announcement bound with the real Pi extension against a refused handling handshake, and a handling successor that must surface a real crew event instead of going blind.
`tests/fm-watcher-lock.test.sh` covers verified-successor attach, recovery publication before stale-lock removal, the typed self-eviction failure, bounded and successor-linked lifecycle rows, and a SIGSTOP counterfactual that distinguishes a live PID from a stale beacon before classifying termination.
`tests/fm-backend-herdr.test.sh` runs the real `bin/fm-watch.sh` loop against a saturating fake herdr event stream under a shrunken grace window and proves the beacon keeps advancing across bounded event waits while a second watcher invocation refuses cleanly against the fresh-beat singleton (the 2026-08-21 quiet-fleet crash-loop regressions).
`tests/fm-subagent-pretool-check.test.sh` proves Claude retains only the non-status Bash seatbelts.
`tests/fm-claude-stop-autoarm.test.sh` covers the auto-arm's scope, stale and live session owners, unchanged AFK and need boundaries, single-flight, bounded failure retries, benign live-watcher cycle ends, one-notice failure episodes, and exit-2 translation.
It also covers abandoned single-flight claims: a claim the ledger shows already finished, and one whose recorded pid-identity no longer matches its live pid while the ledger still reads arming or is absent entirely, are both reclaimed so a lapsed home re-arms, while an identity-matched claim still arming, one the ledger does not name, and the guard's own terminal check keep the gate closed ([`turnend-guard.md`](turnend-guard.md) owns that boundary).
Expand Down
91 changes: 91 additions & 0 deletions tests/fm-backend-herdr-eventwait.test.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,97 @@ def test_main_does_not_signal_readiness_before_valid_ack(self):
self.assertEqual(result, 3)
self.assertEqual(stdout.getvalue(), "")

def test_main_saturated_stream_stays_within_budget(self):
# 2026-08-21 quiet-fleet crash-loop regression: a stream that keeps
# producing events (a pending-backlog replay, a fast-flapping agent)
# with a consumer slower than the stream must not make the reader
# outlive its deadline - neither in the stream loop nor inside a
# blocking stdout write on a full pipe. Drive the REAL reader process
# against a saturating fake server and a throttled stdout consumer and
# assert it exits at the deadline (rc 0, a clean bounded wait).
import json as _json
import os
import subprocess
import sys as _sys
import tempfile
import threading

timeout = 2
tmpdir = tempfile.mkdtemp(prefix="fm-eventwait-budget.")
self.addCleanup(
lambda: __import__("shutil").rmtree(tmpdir, ignore_errors=True)
)
sock_path = os.path.join(tmpdir, "s.sock")
server = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
server.bind(sock_path)
server.listen(1)
self.addCleanup(server.close)
stop = threading.Event()

def serve():
try:
conn, _ = server.accept()
with conn:
conn.recv(65536)
conn.sendall(
b'{"id":"x","result":{"type":"subscription_started"}}\n'
)
event = (
_json.dumps(
{
"event": "pane.agent_status_changed",
"data": {
"pane_id": "w1:p2",
"workspace_id": "w1",
"agent_status": "working",
"agent": "claude",
},
}
)
+ "\n"
).encode()
while not stop.is_set():
conn.sendall(event)
except OSError:
pass

thread = threading.Thread(target=serve, daemon=True)
thread.start()

proc = subprocess.Popen(
[_sys.executable, str(READER_PATH), sock_path, str(timeout), "w1:p2"],
stdout=subprocess.PIPE,
stderr=subprocess.DEVNULL,
)

def consume_slowly():
while True:
line = proc.stdout.readline()
if not line:
return
time.sleep(0.02)

consumer = threading.Thread(target=consume_slowly, daemon=True)
consumer.start()

deadline = time.monotonic() + timeout + 6
while time.monotonic() < deadline:
if proc.poll() is not None:
break
time.sleep(0.1)

try:
self.assertIsNotNone(
proc.poll(),
"reader still alive long past its wait budget under a "
"saturated stream (quiet-fleet crash-loop regression)",
)
self.assertEqual(proc.returncode, 0)
finally:
stop.set()
proc.kill()
proc.wait()


if __name__ == "__main__":
unittest.main()
Loading