diff --git a/.agents/skills/afk/SKILL.md b/.agents/skills/afk/SKILL.md index 98b4d684782..17ec40f52ad 100644 --- a/.agents/skills/afk/SKILL.md +++ b/.agents/skills/afk/SKILL.md @@ -150,7 +150,7 @@ The daemon still clears its buffer only on the backend's `empty` success verdict The daemon wraps `fm-watch.sh`, runs the watcher as a child, presents every durable wake after each actionable watcher close, classifies each presented record in bash, and acknowledges the presented generation only after routing completes. It self-handles the routine majority without consuming a firstmate turn. -Captain-relevant events, plus a bounded recheck of a declared external wait that is still declared, escalate to firstmate's context as one pre-read, single-line, batched digest. +Captain-relevant events, plus a bounded recheck of a declared external wait that is still declared, escalate to firstmate's context as pre-read, single-line, batched digests. The captain-relevant verb set, declared-wait vocabulary, status-span classifier, and presentation-marker contract live in shared `bin/fm-classify-lib.sh`, while each supervisor owns its routing and fleet scan as a consumer of that policy. While `state/.afk` exists the daemon owns the watcher, so the watcher reverts to one-shot and lets the daemon do the triage - the two never run their triage at the same time. @@ -175,10 +175,12 @@ Classify each wake this way: - An unknown wake reason escalates fail-safe, while status-read uncertainty follows the shared one-report-without-position-advance contract referenced under Dedupe below. Escalations are buffered up to `FM_ESCALATE_BATCH_SECS` (default 90s; 0 = -immediate) and flushed as one single-line digest prefixed with the current +immediate) and flushed in oldest-first, single-line digests prefixed with the current operational prefix, carrying pre-read status summaries and a recommended action. The single-line format makes the submission unambiguous across harnesses, and the operational prefix lets firstmate distinguish it from a real captain message. +Each digest has a fixed 1,000-byte bound so it stays below every transport limit it must pass; only confirmed deliveries leave the buffer, and queued events follow in later batches. +Signal events from different status files are buffered separately with their source paths, so truncation points to the source log even when event text names another status file. ### Injection hardening @@ -214,7 +216,7 @@ the operational prefix lets firstmate distinguish it from a real captain message text firstmate sees is clean. - **Portable singleton lock** - the daemon uses the repo's portable lock helper (`fm-wake-lib.sh`) instead of `flock`, which is absent on macOS. -- **Dedupe across signal/stale/scan** - all three paths use the shared status presentation markers defined by `bin/fm-classify-lib.sh`, so a successfully classified span is not re-escalated by another path in the same digest. +- **Dedupe across signal/stale/scan** - all three paths use the shared status presentation markers defined by `bin/fm-classify-lib.sh`, so a successfully classified span is not re-escalated by another path. Never treat a reported unreadable state as classified; the shared library header owns that marker contract, and the marker does not clear or suppress possible-wedge aging for a nonterminal progress line. - **Auto-discovered supervisor pane** - the daemon resolves its own BACKEND (tmux vs herdr) and TARGET independently, mirroring diff --git a/bin/fm-afk-return.sh b/bin/fm-afk-return.sh index 08dc5f86b7d..3f1bbb34a38 100755 --- a/bin/fm-afk-return.sh +++ b/bin/fm-afk-return.sh @@ -688,7 +688,15 @@ EOF append_evidence wedge "$wedge" "$evidence" fi if [ -s "$STATE/.subsuper-escalations" ]; then - escalations=$(cat "$STATE/.subsuper-escalations" 2>/dev/null || true) + escalations=$( + while IFS= read -r record || [ -n "$record" ]; do + if [[ $record == @status-log=*$'\t'* ]]; then + printf '%s\n' "${record#*$'\t'}" + else + printf '%s\n' "$record" + fi + done < "$STATE/.subsuper-escalations" + ) append_evidence escalation "$escalations" "$evidence" fi diff --git a/bin/fm-supervise-daemon.sh b/bin/fm-supervise-daemon.sh index 472d19a20cb..a7136d7e9b8 100755 --- a/bin/fm-supervise-daemon.sh +++ b/bin/fm-supervise-daemon.sh @@ -9,8 +9,7 @@ # token-efficient replacement for the prior always-inject daemon: routine # signal/stale/heartbeat wakes cost zero firstmate context; only done/ # needs-decision/blocked/failed/persistent-wedge/check-output events and a -# declared-wait recheck reach the LLM, and even then as one pre-read digest per -# batch window. +# declared-wait recheck reach the LLM, and even then as bounded pre-read digests. # # PRESENCE-GATING (the /afk contract). The daemon is the away-mode engine: it # injects ONLY when the durable away-mode flag state/.afk is present. Invoking @@ -224,6 +223,8 @@ WEDGE_ALARM_NOTIFIER_PID= INJECT_FAIL_SLEEP_DEFAULT=30 INJECT_CONFIRM_RETRIES_DEFAULT=3 INJECT_CONFIRM_SLEEP_DEFAULT=0.5 +INJECT_MAX_BYTES=1000 +ESCALATION_ITEM_MAX_BYTES=800 CRASH_THRESHOLD_DEFAULT=10 CRASH_WINDOW_DEFAULT=60 CRASH_BACKOFF_DEFAULT=60 @@ -365,6 +366,7 @@ classify_signal() { # marker=$(_seen_status_path "$state" "$task") status_presentation_marker_reported_matches "$marker" "$sig" && continue distilled="${distilled}$(basename "$f"): unreadable status span | " + [ -z "${FM_ESCALATION_ITEMS_FILE:-}" ] || printf '%s\t%s\n' "$f" "$(basename "$f"): unreadable status span" >> "$FM_ESCALATION_ITEMS_FILE" || return 1 [ -n "${FM_STATUS_SPAN_ENDPOINT_FILE:-}" ] \ && printf 'ERROR\t%s\t%s\n' "$task" "$sig" >> "$FM_STATUS_SPAN_ENDPOINT_FILE" rel=1 @@ -377,12 +379,14 @@ classify_signal() { # if [ "$rc" -eq 0 ]; then event=${rest#*$'\t'} distilled="${distilled}$(basename "$f"): ${event} | " + [ -z "${FM_ESCALATION_ITEMS_FILE:-}" ] || printf '%s\t%s\n' "$f" "$(basename "$f"): $event" >> "$FM_ESCALATION_ITEMS_FILE" || return 1 rel=1 continue fi last=$(last_status_line "$f") [ -n "$last" ] || continue distilled="${distilled}$(basename "$f"): ${last} | " + [ -z "${FM_ESCALATION_ITEMS_FILE:-}" ] || printf '%s\t%s\n' "$f" "$(basename "$f"): $last" >> "$FM_ESCALATION_ITEMS_FILE" || return 1 # Nothing captain-relevant is left ahead of the recorded offset. When the log # nonetheless ends on a captain-relevant line, this signal is a re-notification # of something already escalated, not a routine one; position is the whole @@ -421,7 +425,7 @@ classify_stale() { # [ ] if [ "$rc" -eq 0 ]; then rest=${record#*$'\t'} event=${rest#*$'\t'} - printf 'escalate|stale + actionable status: %s' "$event" + printf 'escalate|%s.status: stale + actionable status: %s' "$task" "$event" return fi if [ -n "$last" ] && status_is_paused_or_captain_held "$last"; then @@ -473,7 +477,8 @@ classify_unknown() { # # --- stale marker + escalation buffer (stateful, but via explicit state dir) - # Marker: state/.subsuper-stale- contains the epoch first seen idle. -# Buffer: state/.subsuper-escalations one distilled line per escalation. +# Buffer: state/.subsuper-escalations one distilled line per event; records +# with a source log carry @status-log= before the text. # Seen: state/.subsuper-seen-status- last reported file signature and # classified byte offset, so failures and events do not re-fire while # unread bytes remain recoverable. @@ -692,28 +697,127 @@ stale_window_is_busy() { # [ "${verdict%% *}" = busy ] } -escalate_add() { # - local state=$1 item=$2 buf +escalate_add() { # [source-status-log] + local state=$1 item=$2 source=${3:-} buf over record buf="$state/.subsuper-escalations" [ -s "$buf" ] || _now > "${buf}.since" + record=$item + [ -z "$source" ] || record="@status-log=$source"$'\t'"$item" + over=$(( $(_byte_len "$record") - ESCALATION_ITEM_MAX_BYTES )) + [ "$over" -le 0 ] || item=$(_escalation_item_truncate "$item" "$over" "$source") + [ -z "$source" ] || item="@status-log=$source"$'\t'"$item" printf '%s\n' "$item" >> "$buf" } -# Flush the escalation buffer as ONE batched, single-line digest to the -# supervisor pane. Returns 0 on successful inject (or empty buffer), non-zero on +# --- digest byte bound --------------------------------------------------------- +# One inject is typed as a single argument to the backend's send command, so it +# must stay below every transport ceiling it can meet: Linux refuses any single +# exec argument of 128 KiB or more, tmux refuses a command of about 16 KB, and a +# Claude composer on Herdr can drop the head of a typed burst above about 1,020 +# characters. A digest over a ceiling never reaches the pane, and because the +# buffer is kept on failure every retry would resend the same batch forever. +# escalate_flush therefore sends at most INJECT_MAX_BYTES of typed text per +# inject and leaves the rest buffered for later batches. + +# Byte length of , independent of the caller's locale. +_byte_len() ( # + LC_ALL=C + printf '%s' "${#1}" +) + +# The longest prefix of of at most bytes that does not end inside +# a UTF-8 sequence. +_cut_bytes() ( # + LC_ALL=C + s=$1 + [ "${#s}" -gt "$2" ] || { printf '%s' "$s"; exit 0; } + s=${s:0:$2} + t=$s + c=0 + while :; do + case "${t: -1}" in [$'\x80'-$'\xbf']) t=${t%?}; c=$((c + 1)) ;; *) break ;; esac + done + case "${t: -1}" in + [$'\xc0'-$'\xdf']) need=1 ;; + [$'\xe0'-$'\xef']) need=2 ;; + [$'\xf0'-$'\xf7']) need=3 ;; + *) need=$c ;; + esac + # Drop the last character only when the cut left it incomplete. + [ "$c" -eq "$need" ] || s=${t%?} + printf '%s' "$s" +) + +# Shorten one buffered item by at least bytes and end it with a marker +# naming the dropped byte count and, for a status-log event, the log that still +# holds the full text. +_escalation_item_truncate() ( # + item=$1 over=$2 source=$3 + LC_ALL=C + [ -z "$source" ] || source="; full text in $source" + # The marker is sized with the item's full length, which bounds the digits of + # the count actually dropped. + marker=" ... [+${#item} bytes truncated$source]" + keep=$(( ${#item} - over - ${#marker} )) + [ "$keep" -gt 0 ] || keep=0 + head=$(_cut_bytes "$item" "$keep") + printf '%s ... [+%s bytes truncated%s]' "$head" "$(( ${#item} - ${#head} ))" "$source" +) + +_escalation_digest() { # + printf 'Supervisor escalate (%s event(s)): %s (pre-read; re-arm not needed — watcher daemon-managed)' "$1" "$2" +} + +# Flush the oldest buffered escalations that fit one inject as a single-line +# digest to the supervisor pane. A first item too large to fit alone is +# truncated, so every flush delivers at least one item. Returns 0 on successful +# inject (or empty buffer) after removing only the delivered lines, non-zero on # inject failure (buffer preserved for retry / catch-up). escalate_flush() { # - local state=$1 buf item n msg + local state=$1 buf budget envelope envelope_bytes record source item joined='' try msg='' over taken=0 buf="$state/.subsuper-escalations" [ -s "$buf" ] || return 0 - n=$(wc -l < "$buf" 2>/dev/null || echo 0) - # Join buffered items with the literal " | " separator into one digest line. - msg=$(awk 'NR>1{printf " | "} {printf "%s",$0} END{print ""}' "$buf" 2>/dev/null) - # Single-line wrapper: no embedded newlines (inject_msg also collapses as a - # safety net, but keeping the source single-line makes the intent explicit). - msg=$(printf 'Supervisor escalate (%s event(s)): %s (pre-read; re-arm not needed — watcher daemon-managed)' "$n" "$msg") - if inject_msg "$msg" "$state"; then : > "$buf"; rm -f "${buf}.since" "$state/.subsuper-inject-wedged"; return 0; fi - return 1 + [ -f "$buf" ] || return 1 + budget=$INJECT_MAX_BYTES + # inject_msg wraps the digest in the typed envelope, which counts too. + fm_operational_input_encode away-supervisor x envelope || return 1 + envelope_bytes=$(( $(_byte_len "$envelope") - 1 )) + while IFS= read -r record || [ -n "$record" ]; do + source= + item=$record + if [[ $record == @status-log=*$'\t'* ]]; then + source=${record%%$'\t'*} + source=${source#@status-log=} + item=${record#*$'\t'} + fi + try=$(_escalation_digest "$((taken + 1))" "${joined:+$joined | }$item") + over=$(( envelope_bytes + $(_byte_len "$try") - budget )) + if [ "$over" -gt 0 ]; then + [ "$taken" -eq 0 ] || break + item=$(_escalation_item_truncate "$item" "$over" "$source") + try=$(_escalation_digest 1 "$item") + fi + # Join items with the literal " | " separator into one digest line. + joined=${joined:+$joined | }$item + msg=$try + taken=$((taken + 1)) + done < "$buf" + inject_msg "$msg" "$state" || return 1 + if ! tail -n +"$((taken + 1))" "$buf" > "${buf}.tmp" 2>/dev/null; then + rm -f "${buf}.tmp" + log "inject delivered but escalation buffer update failed: remainder copy" + return 1 + fi + if ! mv -f "${buf}.tmp" "$buf"; then + rm -f "${buf}.tmp" + log "inject delivered but escalation buffer update failed: remainder replacement" + return 1 + fi + # Delivery works again, so the max-defer clock restarts for any remainder and + # the next batch goes after the normal batch window. + if [ -s "$buf" ]; then _now > "${buf}.since"; else rm -f "${buf}.since"; fi + rm -f "$state/.subsuper-inject-wedged" + return 0 } # --- backend-independent active wedge alert --------------------------------- @@ -1185,7 +1289,7 @@ housekeeping() { # ident=$(status_observed_signature "$f") status_presentation_marker_reported_matches "$(_seen_status_path "$state" "$task")" "$ident" \ && continue - if escalate_add "$state" "$(basename "$f"): unreadable status span (catch-all scan)"; then + if escalate_add "$state" "$(basename "$f"): unreadable status span (catch-all scan)" "$f"; then status_presentation_marker_report "$(_seen_status_path "$state" "$task")" "$ident" || true fi continue @@ -1195,11 +1299,11 @@ housekeeping() { # rest=${record#*$'\t'}; ident=${rest%%$'\t'*} if [ "$rc" -eq 0 ]; then event=${rest#*$'\t'} - if escalate_add "$state" "$(basename "$f"): $event (catch-all scan)"; then + if escalate_add "$state" "$(basename "$f"): $event (catch-all scan)" "$f"; then mark_status_seen "$state" "$task" "$endpoint" "$ident" || true fi elif ! mark_status_seen "$state" "$task" "$endpoint" "$ident"; then - escalate_add "$state" "$(basename "$f"): status position commit failed (catch-all scan)" + escalate_add "$state" "$(basename "$f"): status position commit failed (catch-all scan)" "$f" fi done fi @@ -1296,7 +1400,11 @@ inject_msg() { # [state] if [ "$verdict" = empty ]; then return 0 # Backend confirmed the submit. fi - log "inject failed: submit unconfirmed after $retries retries (verdict=$verdict, text may be in composer)" + if [ "$verdict" = send-failed ]; then + log "inject failed: backend text send or submit key failed (verdict=send-failed, $(_byte_len "$msg") bytes, text may be in composer)" + else + log "inject failed: submit unconfirmed after $retries retries (verdict=$verdict, $(_byte_len "$msg") bytes, text may be in composer)" + fi return 1 } @@ -1335,8 +1443,9 @@ is_wake_reason() { # # is populated, suppression markers commit, and the digest names the decision # instead of "unknown wake:". handle_wake() { # - local reason=$1 state=$2 decision action distilled task last stale_detail + local reason=$1 state=$2 decision action distilled task last stale_detail source='' item buffered local capture="$state/.subsuper-classified-end.$$" span_record='' span_rc='' endpoint ident rest sig marker + local items_file="$state/.subsuper-classified-items.$$" local kind="" arg="" classification_failed=0 span_failure_repeat=0 : > "$capture" || return 1 if should_force_self "$reason"; then @@ -1351,7 +1460,9 @@ handle_wake() { # needs-decision:*) arg="${reason#needs-decision: }" ;; *) arg="${reason#signal: }" ;; esac - decision=$(FM_STATUS_SPAN_ENDPOINT_FILE="$capture" classify_signal "$arg" "$state") ;; + : > "$items_file" || { rm -f "$capture"; return 1; } + decision=$(FM_STATUS_SPAN_ENDPOINT_FILE="$capture" FM_ESCALATION_ITEMS_FILE="$items_file" classify_signal "$arg" "$state") \ + || { rm -f "$capture" "$items_file"; return 1; } ;; stale:*) kind=stale; arg="${reason#stale: }"; stale_detail="${arg#"$arg"}" case "$arg" in *" ("*) stale_detail="${arg#*" ("}"; arg="${arg%% \(*}" ;; esac task=$(window_to_task "$arg" "$state") @@ -1381,6 +1492,7 @@ handle_wake() { # decision="self|unreadable status span already reported for $task" else decision=$(classify_stale "$arg" "$state" "$span_record" "$span_rc") + [ "$span_rc" != 0 ] || source="$state/$task.status" fi # An enriched wedge reason carries the watcher's own escalation count # and its "do not re-absorb on the run-step/pane state alone" demand, @@ -1397,8 +1509,10 @@ handle_wake() { # *) case "$stale_detail" in idle\ *s,\ possible\ wedge,\ escalation\ *) last=$(last_status_line "$state/$task.status") - status_is_paused_or_captain_held "$last" \ - || decision="escalate|${reason#stale: }" + if ! status_is_paused_or_captain_held "$last"; then + decision="escalate|${reason#stale: }" + source= + fi ;; esac ;; esac ;; @@ -1417,7 +1531,16 @@ handle_wake() { # case "$action" in escalate) log "escalate: $reason -> $distilled" - if escalate_add "$state" "$distilled"; then + buffered=0 + if [ "$kind" = signal ]; then + while IFS=$'\t' read -r source item; do + escalate_add "$state" "$item" "$source" || { buffered=1; break; } + done < "$items_file" + [ -s "$items_file" ] || buffered=1 + else + escalate_add "$state" "$distilled" "$source" || buffered=1 + fi + if [ "$buffered" -eq 0 ]; then # A terminal-stale escalate must not leave a persistence marker behind, or # housekeeping re-escalates the same pane as a false wedge later. [ "$kind" = "stale" ] && stale_marker_remove "$arg" "$state" @@ -1473,7 +1596,7 @@ handle_wake() { # if [ "$action" = self ] && { [ "$kind" = signal ] || [ "$kind" = stale ]; }; then mark_escalated_seen "$state" "$capture" || classification_failed=1 fi - rm -f "$capture" + rm -f "$capture" "$items_file" [ "$classification_failed" -eq 0 ] } diff --git a/docs/architecture.md b/docs/architecture.md index fe0048ae206..f5ce4e8ea26 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -196,7 +196,7 @@ The daemon's declared-wait window ages against the crew's own latest status line A wake already decorated as a possible wedge does not override the daemon's own declared-wait verdict either, so a declaration keeps its pane on the recheck cadence instead of the wedge cadence. In away mode, seen-status dedupe does not clear possible-wedge aging for nonterminal progress, so housekeeping still re-escalates an unchanged idle pane at the configured bound. Away-mode housekeeping has no worktree-write deferral of its own, so while `state/.afk` exists a quiet crew that is writing its own worktree still escalates as a possible wedge at that bound. -The daemon escalates captain-relevant events, plus a bounded recheck for a declared external wait that is still declared, as one batched, single-line digest using the canonical `away-supervisor` kind from `bin/fm-operational-input.sh` so firstmate can distinguish it structurally from real messages; captain-held transfers remain silent until return while the posture record exists. +The daemon escalates captain-relevant events and bounded declared-wait rechecks through the batching contract in the [AFK skill](../.agents/skills/afk/SKILL.md); its injections use the canonical `away-supervisor` kind from `bin/fm-operational-input.sh` so firstmate can distinguish them structurally from real messages, while captain-held transfers remain silent until return while the posture record exists. Its supervisor injection path supports tmux and herdr panes, with `FM_SUPERVISOR_BACKEND` and `FM_SUPERVISOR_TARGET` resolved independently from the task-spawn backend. Pane existence, busy checks, composer checks, capture, and verified submit route through `bin/fm-backend.sh`: tmux keeps the same submit core used by the tmux send backend, while herdr uses native agent-state submit confirmation on idle baselines, a composer empty fallback when native stays idle, and a pre-Enter rendered-footer transition when that baseline is unavailable. The retries-exhausted queued-Enter decision is owned by `fm_composer_queued_enter_verdict` in `bin/fm-composer-lib.sh`; tmux and herdr provide only their backend-specific busy signals. diff --git a/tests/fm-afk-inject-e2e.test.sh b/tests/fm-afk-inject-e2e.test.sh index 65de2e6e1af..41a89e003bf 100755 --- a/tests/fm-afk-inject-e2e.test.sh +++ b/tests/fm-afk-inject-e2e.test.sh @@ -1,6 +1,6 @@ #!/usr/bin/env bash # tests/fm-afk-inject-e2e.test.sh - private-socket end-to-end test for the afk -# daemon's injection path. It covers three operator-visible injection contracts: +# daemon's injection path. It covers four operator-visible injection contracts: # # Scenario A (human-partial-input): a partial line is typed into the # supervisor pane with NO Enter, then an escalation fires. The daemon must @@ -15,6 +15,10 @@ # A captain-relevant status must deliver exactly ONE sentinel-prefixed, # single-line digest with no duplicate or spurious user submission. # +# Scenario D (oversized buffer): buffered items over tmux's command limit and +# the kernel's single-argument limit must still drain, as bounded digests +# that truncate the oversized items and deliver every event. +# # Isolation: all test tmux runs on a dedicated socket (tmux -L afk-e2e-). # A tmux shim first on PATH redirects the daemon's bare `tmux` calls to the # private socket. The daemon points at a throwaway state dir (FM_STATE_OVERRIDE) @@ -421,8 +425,61 @@ test_scenario_c() { pass "Scenario C: a normal captain status injects exactly one clean single-line sentinel digest" } +# --- Scenario D: oversized buffer drains in bounded digests ----------------- +# One buffered item over the kernel's 128 KiB single-argument limit and one over +# tmux's ~16 KB command limit: joined into one digest, real tmux (or exec) +# refuses the send and the buffer never drains. Each flush must instead send at +# most 1,000 bytes and keep the rest for the next batch. + +test_scenario_d() { + local big mid i=0 injections line text bytes + reset_state + big=$(head -c 150000 /dev/zero | tr '\0' 'x') + mid=$(head -c 20000 /dev/zero | tr '\0' 'y') + : > "$STATE_DIR/big-d1.status" + : > "$STATE_DIR/mid-d2.status" + escalate_add "$STATE_DIR" "event A: done: PR https://example.test/pr/401" + escalate_add "$STATE_DIR" "big-d1.status: done: $big | extra context" "$STATE_DIR/big-d1.status" + escalate_add "$STATE_DIR" "mid-d2.status: done: $mid" "$STATE_DIR/mid-d2.status" + escalate_add "$STATE_DIR" "event B: done: PR https://example.test/pr/402" + afk_enter "$STATE_DIR" + start_daemon + + while [ -s "$STATE_DIR/.subsuper-escalations" ] && [ "$i" -lt 150 ]; do + sleep 0.2 + i=$((i + 1)) + done + [ ! -s "$STATE_DIR/.subsuper-escalations" ] \ + || fail "Scenario D: oversized buffer did not drain: $(grep 'inject' "$STATE_DIR/.supervise-daemon.log" | tail -3)" + sleep 1 + + injections=$(grep -c $'\tinjection$' "$LOG_FILE" || true) + [ "$injections" -ge 2 ] || fail "Scenario D: expected multiple bounded digests, got $injections" + grep -F 'more queued' "$LOG_FILE" >/dev/null && fail "Scenario D: digest announced a queued count" + while IFS= read -r line; do + text=$(printf '%s' "$line" | cut -f2) + bytes=$(LC_ALL=C; printf '%s' "${#text}") + [ "$bytes" -le 1000 ] || fail "Scenario D: a submitted digest was $bytes bytes" + done < "$LOG_FILE" + grep -F 'event A: done: PR https://example.test/pr/401' "$LOG_FILE" >/dev/null \ + || fail "Scenario D: first small event not delivered" + grep -F "big-d1.status: done: xxx" "$LOG_FILE" | grep -F "bytes truncated; full text in $STATE_DIR/big-d1.status]" >/dev/null \ + || fail "Scenario D: over-128KiB item not delivered truncated with its log pointer" + grep -F "mid-d2.status: done: yyy" "$LOG_FILE" | grep -F "bytes truncated; full text in $STATE_DIR/mid-d2.status]" >/dev/null \ + || fail "Scenario D: over-16KB item not delivered truncated with its log pointer" + grep -F 'event B: done: PR https://example.test/pr/402' "$LOG_FILE" >/dev/null \ + || fail "Scenario D: last small event not delivered" + if grep -q 'inject failed' "$STATE_DIR/.supervise-daemon.log"; then + fail "Scenario D: an inject failed: $(grep 'inject failed' "$STATE_DIR/.supervise-daemon.log" | head -1)" + fi + + stop_daemon + pass "Scenario D: an oversized escalation buffer drains through real tmux in bounded digests" +} + test_scenario_a test_scenario_b test_scenario_c +test_scenario_d echo "all e2e injection tests passed" diff --git a/tests/fm-afk-inject-herdr-e2e.test.sh b/tests/fm-afk-inject-herdr-e2e.test.sh index e761336e7b4..a5a3b043893 100755 --- a/tests/fm-afk-inject-herdr-e2e.test.sh +++ b/tests/fm-afk-inject-herdr-e2e.test.sh @@ -523,9 +523,57 @@ test_scenario_d_max_defer() { pass "real herdr Scenario D: a persistently pending composer raises the max-defer wedge alarm, preserves the buffer, and never crashes the daemon" } +# --- Scenario E: oversized backlog and a misleading status excerpt ---------- +# Exercise the incident's large buffered send on the real Herdr transport. A +# status may quote another existing status filename without creating a second +# event or changing the recovery pointer for its own truncated text. +test_scenario_e_bounded_backlog() { + local big excerpt line text bytes injections + reset_state + big=$(head -c 150000 /dev/zero | tr '\0' 'x') + excerpt=$(head -c 2000 /dev/zero | tr '\0' 'z') + printf 'working: routine\n' > "$STATE_DIR/b.status" + : > "$STATE_DIR/big.status" + escalate_add "$STATE_DIR" "event A: done: PR https://example.test/pr/401" + escalate_add "$STATE_DIR" "big.status: done: $big" "$STATE_DIR/big.status" + escalate_add "$STATE_DIR" "event B: done: PR https://example.test/pr/402" + afk_enter "$STATE_DIR" + start_daemon + printf 'needs-decision [key=choose]: inspect excerpt | b.status: %s\n' "$excerpt" > "$STATE_DIR/a.status" + sleep 20 + + [ ! -s "$STATE_DIR/.subsuper-escalations" ] \ + || fail "Scenario E: Herdr backlog did not drain: $(tail -3 "$STATE_DIR/.supervise-daemon.log")" + injections=$(grep -c $'\tinjection$' "$LOG_FILE" || true) + [ "$injections" -ge 2 ] || fail "Scenario E: expected multiple Herdr submissions, got $injections" + grep -F 'event A: done: PR https://example.test/pr/401' "$LOG_FILE" >/dev/null \ + || fail "Scenario E: oldest event was not delivered" + grep -F "full text in $STATE_DIR/big.status]" "$LOG_FILE" >/dev/null \ + || fail "Scenario E: 150 KB event lost its status-log pointer" + grep -F 'event B: done: PR https://example.test/pr/402' "$LOG_FILE" >/dev/null \ + || fail "Scenario E: later event was not delivered" + grep -F "full text in $STATE_DIR/a.status]" "$LOG_FILE" >/dev/null \ + || fail "Scenario E: quoted b.status changed the a.status recovery pointer" + grep -F "full text in $STATE_DIR/b.status]" "$LOG_FILE" >/dev/null \ + && fail "Scenario E: quoted b.status became a false event" + while IFS= read -r line; do + text=$(printf '%s' "$line" | cut -f2) + bytes=$(LC_ALL=C; printf '%s' "${#text}") + [ "$bytes" -le 1000 ] || fail "Scenario E: Herdr submission was $bytes bytes" + printf 'Herdr submitted (%s bytes): %s\n' "$bytes" "$text" + done < "$LOG_FILE" + grep -q 'inject failed' "$STATE_DIR/.supervise-daemon.log" \ + && fail "Scenario E: Herdr refused a bounded digest" + + printf 'real Herdr submitted %s bounded digests; oldest and later events arrived; truncated a.status pointed to itself\n' "$injections" + stop_daemon + pass "real herdr Scenario E: oversized backlog drains and a quoted existing status stays one event" +} + test_scenario_a test_scenario_b test_scenario_c +test_scenario_e_bounded_backlog test_scenario_d_max_defer echo "all real-herdr afk injection e2e tests passed" diff --git a/tests/fm-afk-return.test.sh b/tests/fm-afk-return.test.sh index 0ce90aba151..120aa863d76 100755 --- a/tests/fm-afk-return.test.sh +++ b/tests/fm-afk-return.test.sh @@ -114,7 +114,8 @@ test_return_gate_owns_remediation_and_reports_catchup_to_bearings() { printf '\n## Done\n' } > "$dir/home/data/backlog.md" date +%s > "$dir/home/state/.afk" - printf 'repair-task.status: blocked synthetic dependency\n' > "$dir/home/state/.subsuper-escalations" + printf '@status-log=%s\trepair-task.status: blocked synthetic dependency\n' \ + "$dir/home/state/repair-task.status" > "$dir/home/state/.subsuper-escalations" printf 'fm away-mode inject WEDGED: 4555s undelivered\n' > "$dir/home/state/.subsuper-inject-wedged" { printf '1784074271\t2\tsignal\trepair-task.status\tsignal: synthetic status\n' diff --git a/tests/fm-daemon.test.sh b/tests/fm-daemon.test.sh index 510a4320988..939492ff46a 100755 --- a/tests/fm-daemon.test.sh +++ b/tests/fm-daemon.test.sh @@ -1432,6 +1432,207 @@ test_escalate_batches_into_one_digest() { pass "multiple escalations flush as a single batched digest" } +# An oversized buffer (the 2026-09-22 overnight shape: one catch-all span far +# over the kernel's single-argument limit) must drain in bounded batches: each +# typed digest fits the fixed byte budget, only delivered lines leave the buffer, +# and an item too long to fit alone is truncated with a pointer to its log. +test_escalate_flush_bounds_each_digest() { + local dir state fakebin sent capture big mid flushes=0 line bytes + dir=$(make_supercase batch-bounded) + state="$dir/state" + fakebin="$dir/fakebin" + sent="$dir/sent.log"; : > "$sent" + capture="$dir/pane.txt"; printf '\342\235\257 \n' > "$capture" + big=$(head -c 200000 /dev/zero | tr '\0' 'x') + mid=$(head -c 20000 /dev/zero | tr '\0' 'y') + : > "$state/big-t1.status" + : > "$state/mid-t2.status" + escalate_add "$state" "event A: done: PR 1" + escalate_add "$state" "event B: done: PR 2" + escalate_add "$state" "big-t1.status: done: $big | extra context" "$state/big-t1.status" + escalate_add "$state" "mid-t2.status: done: $mid" "$state/mid-t2.status" + escalate_add "$state" "event C: done: PR 3" + [ "$(wc -l < "$state/.subsuper-escalations")" -eq 5 ] \ + || fail "combined signal did not preserve each task as a separate buffered item" + while IFS= read -r line; do + bytes=$(LC_ALL=C; printf '%s' "${#line}") + [ "$bytes" -le "$ESCALATION_ITEM_MAX_BYTES" ] || fail "a buffered item was $bytes bytes" + done < "$state/.subsuper-escalations" + grep -F "full text in $state/big-t1.status]" "$state/.subsuper-escalations" >/dev/null \ + || fail "first task lost its status-log pointer when buffered" + grep -F "full text in $state/mid-t2.status]" "$state/.subsuper-escalations" >/dev/null \ + || fail "second task lost its status-log pointer when buffered" + afk_enter "$state" + + PATH="$fakebin:$PATH" FM_FAKE_TMUX_PANE_ALIVE=1 FM_FAKE_TMUX_SENT="$sent" \ + FM_FAKE_TMUX_CAPTURE="$capture" escalate_flush "$state" \ + || fail "first bounded flush failed" + grep -F 'event A: done: PR 1 | event B: done: PR 2' "$sent" >/dev/null \ + || fail "first batch did not carry the oldest items in order" + grep -F 'more queued' "$sent" >/dev/null && fail "digest announced an unrequested queued count" + [ -s "$state/.subsuper-escalations" ] || fail "partial flush dropped the queued remainder" + [ -e "$state/.subsuper-escalations.since" ] || fail "partial flush dropped the remainder's first-append sidecar" + + while [ -s "$state/.subsuper-escalations" ] && [ "$flushes" -lt 5 ]; do + PATH="$fakebin:$PATH" FM_FAKE_TMUX_PANE_ALIVE=1 FM_FAKE_TMUX_SENT="$sent" \ + FM_FAKE_TMUX_CAPTURE="$capture" escalate_flush "$state" \ + || fail "bounded follow-up flush failed" + flushes=$((flushes + 1)) + done + [ ! -s "$state/.subsuper-escalations" ] || fail "bounded flushes did not drain the buffer" + [ ! -e "$state/.subsuper-escalations.since" ] || fail "drained buffer kept its first-append sidecar" + [ "$(grep -c '\[ENTER\]' "$sent")" -ge 2 ] || fail "expected multiple bounded digests" + grep -F "big-t1.status: done: xxx" "$sent" | grep -F "bytes truncated; full text in $state/big-t1.status]" >/dev/null \ + || fail "oversized item was not truncated with a pointer to its status log" + grep -F "mid-t2.status: done: yyy" "$sent" | grep -F "bytes truncated; full text in $state/mid-t2.status]" >/dev/null \ + || fail "second status in a combined signal lost its log pointer" + grep -F 'event C: done: PR 3' "$sent" >/dev/null || fail "last item was not delivered" + while IFS= read -r line; do + [ "$line" = '[ENTER]' ] && continue + bytes=$(LC_ALL=C; printf '%s' "${#line}") + [ "$bytes" -le "$INJECT_MAX_BYTES" ] || fail "a typed digest was $bytes bytes, over the $INJECT_MAX_BYTES-byte budget" + done < "$sent" + pass "an oversized escalation buffer drains in bounded batches and truncates an oversized item" +} + +test_escalate_add_ignores_false_status_boundaries() { + local dir state fakebin sent capture big line + dir=$(make_supercase false-status-boundary) + state="$dir/state"; fakebin="$dir/fakebin" + sent="$dir/sent.log"; : > "$sent" + capture="$dir/pane.txt"; printf '\342\235\257 \n' > "$capture" + big=$(head -c 2000 /dev/zero | tr '\0' 'z') + printf 'needs-decision [key=choose]: inspect excerpt | b.status: %s\n' "$big" > "$state/a.status" + printf 'working: routine\n' > "$state/b.status" + FM_ESCALATE_BATCH_SECS=999 handle_wake "signal: $state/a.status" "$state" + [ "$(wc -l < "$state/.subsuper-escalations")" -eq 1 ] \ + || fail "status text naming another existing task was split into another event" + line=$(cat "$state/.subsuper-escalations") + case "$line" in + *"full text in $state/a.status]"*) ;; + *) fail "signal lost its own status log pointer" ;; + esac + case "$line" in + *"full text in $state/b.status]"*) fail "signal pointed to the status named in its excerpt" ;; + esac + afk_enter "$state" + PATH="$fakebin:$PATH" FM_FAKE_TMUX_PANE_ALIVE=1 FM_FAKE_TMUX_SENT="$sent" \ + FM_FAKE_TMUX_CAPTURE="$capture" escalate_flush "$state" \ + || fail "single status signal flush failed" + grep -F "full text in $state/a.status]" "$sent" >/dev/null \ + || fail "single status signal lost its pointer during delivery" + grep -F "full text in $state/b.status]" "$sent" >/dev/null \ + && fail "single status signal delivered the excerpt as another task" + [ ! -s "$state/.subsuper-escalations" ] || fail "delivered signal stayed buffered" + printf 'needs-decision [key=c]: choose | d.status: %s\n' "$big" > "$state/c.status" + printf 'done: %s\n' "$big" > "$state/d.status" + FM_ESCALATE_BATCH_SECS=999 handle_wake "signal: $state/c.status $state/d.status" "$state" + [ "$(wc -l < "$state/.subsuper-escalations")" -eq 2 ] \ + || fail "two status files did not produce two buffered records" + grep -F "full text in $state/c.status]" "$state/.subsuper-escalations" >/dev/null \ + || fail "first status in a combined signal lost its pointer" + grep -F "full text in $state/d.status]" "$state/.subsuper-escalations" >/dev/null \ + || fail "second status in a combined signal lost its pointer" + escalate_add "$state" "missing.status: $big" + line=$(tail -1 "$state/.subsuper-escalations") + case "$line" in + *"full text in $state/missing.status]"*) fail "missing status received a recovery pointer" ;; + esac + pass "signal source paths define task boundaries and truncation pointers" +} + +test_escalate_flush_reports_buffer_update_failure() { + local failure dir state fakebin sent capture + for failure in copy replacement; do + dir=$(make_supercase "buffer-update-$failure") + state="$dir/state"; fakebin="$dir/fakebin" + sent="$dir/sent.log"; : > "$sent" + capture="$dir/pane.txt"; printf '\342\235\257 \n' > "$capture" + escalate_add "$state" "needs-decision: pick A" + afk_enter "$state" + ( + tail() { + if [ "$failure" = copy ] && [ "${1:-}" = -n ] && [ "${2:-}" = +2 ]; then return 1; fi + command tail "$@" + } + mv() { + if [ "$failure" = replacement ] && [ "${3:-}" = "$state/.subsuper-escalations" ]; then return 1; fi + command mv "$@" + } + if PATH="$fakebin:$PATH" FM_FAKE_TMUX_PANE_ALIVE=1 FM_FAKE_TMUX_SENT="$sent" \ + FM_FAKE_TMUX_CAPTURE="$capture" LOG="$dir/daemon.log" escalate_flush "$state"; then + exit 1 + fi + ) || fail "flush returned success after $failure failure" + grep -F "needs-decision: pick A" "$sent" >/dev/null \ + || fail "digest was not delivered before $failure failure" + grep -F "escalation buffer update failed: remainder $failure" "$dir/daemon.log" >/dev/null \ + || fail "$failure failure was not logged" + grep -Fx "needs-decision: pick A" "$state/.subsuper-escalations" >/dev/null \ + || fail "$failure failure discarded the original escalation buffer" + [ ! -e "$state/.subsuper-escalations.tmp" ] || fail "$failure failure left a partial remainder file" + done + pass "confirmed delivery reports copy and replacement failures" +} + +test_stale_actionable_status_truncation_keeps_log_pointer() { + local dir state fakebin sent capture big line bytes + dir=$(make_supercase stale-bounded) + state="$dir/state"; fakebin="$dir/fakebin" + sent="$dir/sent.log"; : > "$sent" + capture="$dir/pane.txt"; printf '\342\235\257 \n' > "$capture" + big=$(head -c 2000 /dev/zero | tr '\0' 'z') + printf 'blocked [key=approval]: %s\n' "$big" > "$state/stale-t1.status" + FM_ESCALATE_BATCH_SECS=999 handle_wake 'stale: sess:fm-stale-t1' "$state" + line=$(cat "$state/.subsuper-escalations") + bytes=$(LC_ALL=C; printf '%s' "${#line}") + [ "$bytes" -le "$ESCALATION_ITEM_MAX_BYTES" ] || fail "stale status was not capped when buffered" + case "$line" in + *"full text in $state/stale-t1.status]"*) ;; + *) fail "stale status lost its log pointer when buffered" ;; + esac + afk_enter "$state" + PATH="$fakebin:$PATH" FM_FAKE_TMUX_PANE_ALIVE=1 FM_FAKE_TMUX_SENT="$sent" \ + FM_FAKE_TMUX_CAPTURE="$capture" escalate_flush "$state" \ + || fail "stale status flush failed" + grep -F "full text in $state/stale-t1.status]" "$sent" >/dev/null \ + || fail "stale status lost its log pointer during flush" + printf '@status-log=%s\tstale-t1.status: stale + actionable status: blocked [key=approval]: %s\n' "$state/stale-t1.status" "$big" \ + >> "$state/.subsuper-escalations" + PATH="$fakebin:$PATH" FM_FAKE_TMUX_PANE_ALIVE=1 FM_FAKE_TMUX_SENT="$sent" \ + FM_FAKE_TMUX_CAPTURE="$capture" escalate_flush "$state" \ + || fail "older oversized stale status flush failed" + [ "$(grep -Fc "full text in $state/stale-t1.status]" "$sent")" -eq 2 ] \ + || fail "flush-time truncation lost the stale status log pointer" + pass "stale actionable status keeps its recovery log through buffering and delivery" +} + +test_escalate_flush_truncates_on_a_character_boundary() { + local out + out=$(_cut_bytes 'ab'$'\xc3\xa9''cd' 3) + [ "$out" = 'ab' ] || fail "cut inside a two-byte character kept a partial sequence: $(printf '%s' "$out" | od -An -tx1)" + out=$(_cut_bytes 'ab'$'\xc3\xa9''cd' 4) + [ "$out" = 'ab'$'\xc3\xa9' ] || fail "cut after a complete character dropped it" + pass "digest truncation never splits a UTF-8 character" +} + +test_escalate_flush_send_refusal_is_logged_honestly() { + local dir state fakebin sent + dir=$(make_bordered_case send-refused) + state="$dir/state"; fakebin="$dir/fakebin" + sent="$dir/sent.log"; : > "$sent" + escalate_add "$state" "needs-decision: pick A" + afk_enter "$state" + if PATH="$fakebin:$PATH" FM_FAKE_COMPOSER="$dir/composer" FM_FAKE_SENT="$sent" FM_FAKE_SEND_FAIL=1 \ + FM_INJECT_CONFIRM_SLEEP=0.05 LOG="$dir/daemon.log" escalate_flush "$state"; then + fail "escalate_flush succeeded although the backend refused the send" + fi + grep -E 'inject failed: backend text send or submit key failed \(verdict=send-failed, [0-9]+ bytes, text may be in composer\)' "$dir/daemon.log" >/dev/null \ + || fail "send failure log did not cover both text and submit failures: $(cat "$dir/daemon.log")" + [ -s "$state/.subsuper-escalations" ] || fail "buffer lost after a refused send" + pass "a refused backend send is logged as such, with the digest size" +} + test_escalate_batch_age_uses_first_append() { local dir state fakebin sent capture dir=$(make_supercase batch-age) @@ -2818,6 +3019,12 @@ test_housekeeping_herdr_idle_busy_record_clears_stale test_housekeeping_herdr_resumed_stale_cleared test_housekeeping_orca_persistent_stale_resolves_terminal test_escalate_batches_into_one_digest +test_escalate_flush_bounds_each_digest +test_escalate_add_ignores_false_status_boundaries +test_escalate_flush_reports_buffer_update_failure +test_stale_actionable_status_truncation_keeps_log_pointer +test_escalate_flush_truncates_on_a_character_boundary +test_escalate_flush_send_refusal_is_logged_honestly test_escalate_batch_age_uses_first_append test_heartbeat_scan_dedup test_handle_wake_routes_self_and_escalate