diff --git a/bin/fm-pending-reply-lib.sh b/bin/fm-pending-reply-lib.sh index 42bd2d5de4c..ce5ba3349b0 100755 --- a/bin/fm-pending-reply-lib.sh +++ b/bin/fm-pending-reply-lib.sh @@ -1445,20 +1445,60 @@ fm_pending_reply_tick_one() { # [secondmate- return 0 } +# Print, one per line, the records among ... that the tick has work +# for, reading every record once in a single awk process instead of forking +# per record. A resolved record needs work only while an escalation it opened +# is still unclosed: for every other resolved record the tick's per-record path +# (fm_pending_reply_close_escalation) is a no-op that still pays a lock and +# several forks, and records are never pruned, so that cost grew with the store. +# Every other record - any phase but resolved, no phase at all, or one awk +# cannot read - is selected, so the per-record path still decides it. Values +# follow fm_pending_reply_get: the last line for a key wins. +_fm_pending_reply_select_needing_work() { # ... + [ "$#" -gt 0 ] || return 0 + printf '%s\n' "$@" | LC_ALL=C awk ' + { + path = $0 + phase = ""; escalated = ""; closed = "" + while ((rc = (getline line < path)) > 0) { + if (index(line, "phase=") == 1) phase = substr(line, 7) + else if (index(line, "escalated_epoch=") == 1) escalated = substr(line, 17) + else if (index(line, "escalation_closed_epoch=") == 1) closed = substr(line, 25) + } + close(path) + if (rc < 0 || phase != "resolved" || (escalated != "" && closed == "")) print path + } + ' +} + # Scan every pending record for this parent state. Safe to call every poll. # Never scrapes secondmate conversation; uses only parent status, backend busy -# state, and optional secondmate-home wrong-home path checks. +# state, and optional secondmate-home wrong-home path checks. Records are +# selected in one pass first (_fm_pending_reply_select_needing_work), so a +# settled record costs no lock and no fork, and the per-record path below runs, +# unchanged, only for the records that selection returns. fm_pending_reply_tick() { # local state=$1 dir rec corr task_id phase delivered meta backend target label busy sm_home harness remote_host local observation observation_task found i - local -a observation_tasks=() observation_values=() + local -a observation_tasks=() observation_values=() records=() selected=() dir=$(fm_pending_reply_dir "$state") [ -d "$dir" ] || return 0 for rec in "$dir"/*; do [ -f "$rec" ] || continue - case "$(basename "$rec")" in + case "${rec##*/}" in .*) continue ;; esac + case "$rec" in + # A newline would split this path in the selection's input, so such a + # record skips selection and always takes the per-record path. + *$'\n'*) selected+=("$rec") ;; + *) records+=("$rec") ;; + esac + done + while IFS= read -r rec; do + selected+=("$rec") + done < <(_fm_pending_reply_select_needing_work ${records[@]+"${records[@]}"}) + for rec in ${selected[@]+"${selected[@]}"}; do corr=$(fm_pending_reply_get "$rec" corr_id) [ -n "$corr" ] || corr=$(basename "$rec") task_id=$(fm_pending_reply_get "$rec" task_id) diff --git a/bin/fm-procevent-remote-reply.sh b/bin/fm-procevent-remote-reply.sh index 9794f336962..7d545c0508b 100755 --- a/bin/fm-procevent-remote-reply.sh +++ b/bin/fm-procevent-remote-reply.sh @@ -19,6 +19,8 @@ # ingests it, acknowledges the captured generation, then registers the next # cursor-anchored source. `relisten` tells that runner to poll again in the same # process, still holding the claim, after an empty window and after that re-arm. +# A window the remote job worker preempted is reported to the runner as an empty +# window, so it relistens too (see JOB_PREEMPTED below). # A continuity break is escalated and not re-armed, so the registration is dropped # and the runner stops. The runner does not refresh the owner lease. # @@ -95,7 +97,7 @@ DOCUMENT_LOCAL_FAILURE=2 . "$SCRIPT_DIR/fm-pending-reply-lib.sh" die() { printf 'error: %s\n' "$1" >&2; exit 1; } -usage() { sed -n '2,64p' "$0" | sed 's/^# \{0,1\}//'; exit 2; } +usage() { sed -n '2,66p' "$0" | sed 's/^# \{0,1\}//'; exit 2; } sha256_file() { if command -v shasum >/dev/null 2>&1; then @@ -256,6 +258,14 @@ cmd_arm() { # honest watermark, and bin/fm-pending-reply-lib.sh consumes it so a missing # correlated report is judged only against a channel known to have caught up. WINDOW_CLOSED_EMPTY=75 +# The remote job worker's exit when it preempted this long-poll to run another +# job for the same home (bin/fm-remote-job-lib.sh header), such as the watcher's +# per-cycle liveness probe. The read is cursor-anchored and non-destructive, so a +# preempted window loses nothing: it is a window that closed early, and the +# runner relistens exactly as after WINDOW_CLOSED_EMPTY instead of reading it as +# a failed read that releases the listener's claim. It proves nothing about the +# channel being caught up, so it records no watermark. +JOB_PREEMPTED=76 cmd_source() { local id=${1:-} started rc=0 @@ -266,6 +276,8 @@ cmd_source() { "$REMOTE_LOG" "$CURSOR_OFFSET" "$CURSOR_HASH" "$WAIT_SECONDS" < /dev/null || rc=$? if [ "$rc" -eq "$WINDOW_CLOSED_EMPTY" ]; then fm_pending_reply_note_remote_channel_caught_up "$STATE" "$id" "$started" || true + elif [ "$rc" -eq "$JOB_PREEMPTED" ]; then + rc=$WINDOW_CLOSED_EMPTY fi return "$rc" } diff --git a/bin/fm-wake-lib.sh b/bin/fm-wake-lib.sh index 78da8a58600..60a9d289090 100755 --- a/bin/fm-wake-lib.sh +++ b/bin/fm-wake-lib.sh @@ -141,6 +141,18 @@ fm_poll_derived_grace() { printf '%s\n' "$derived" } +# fm_watcher_stall_bound [poll-seconds] +# Hard bound on a live watcher holder's beacon age: FM_WATCHER_STALL_BOUND, +# defaulting to 3x the watcher's stale grace (FM_WATCHER_STALE_GRACE, else +# FM_GUARD_GRACE, else fm_poll_derived_grace). Under it a live holder with a +# stale beacon is a slow cycle; at or past it bin/fm-watch.sh evicts that holder +# and bin/fm-watch-arm.sh stops following it, so both read this one definition. +fm_watcher_stall_bound() { + local poll=${1:-${FM_POLL:-15}} grace + grace=${FM_WATCHER_STALE_GRACE:-${FM_GUARD_GRACE:-$(fm_poll_derived_grace "$poll")}} + printf '%s\n' "${FM_WATCHER_STALL_BOUND:-$((grace * 3))}" +} + # fm_watcher_lock_unheld # True when the watcher lock or its symlinked owner directory is absent, or when # the existing lock records no pid at all. Any non-empty pid remains held here; diff --git a/bin/fm-watch-arm.sh b/bin/fm-watch-arm.sh index ec95458406f..1932c0b9414 100755 --- a/bin/fm-watch-arm.sh +++ b/bin/fm-watch-arm.sh @@ -32,11 +32,20 @@ # watcher: FAILED - cycle ended without an actionable reason # - a clean cycle ended with no wake and no # verified healthy successor +# watcher: FAILED - attached watcher pid= stalled (beacon s at or past hard bound s) +# - the followed holder is alive but its beacon +# reached the stall bound # It NEVER reports started/attached/healthy off a stale beacon or a dead/reused pid: a # stale-beacon or dead-pid holder either self-heals (the fresh child steals the # dead lock per the singleton self-eviction/steal path and is confirmed) or this # returns the FAILED line. On started it waits the child and propagates the wake -# reason; on attached it stays live across identity-matched successors. A cycle +# reason; on attached it stays live across identity-matched successors. Once +# attached, a stale beacon alone does not end the followed cycle: while that +# holder is alive and the lock still names it under the same identity, the arm +# keeps following it, as a started arm waits out a slow child, until the lock +# changes or the beacon reaches fm_watcher_stall_bound (bin/fm-wake-lib.sh), the +# age at which the watcher's own re-arm evicts it; there it fails with the +# stalled-holder line so its owner's retry replaces the holder. A cycle # that ends with no reason line and no healthy successor is resolved against the # watcher's identity-bound delivery record: a matching record reports that wake # and exits 0, and only a cycle that delivered nothing is the typed nonzero @@ -117,6 +126,9 @@ esac CONFIRM_TIMEOUT=${FM_ARM_CONFIRM_TIMEOUT:-$ARM_CONFIRM_DEFAULT} # Poll interval while attached to an existing healthy watcher. ATTACH_POLL=${FM_ARM_ATTACH_POLL:-0.5} +# The beacon age at which the watcher's own re-arm evicts a live holder; an +# attached arm follows a slow holder up to it (attach_and_wait). +STALL_BOUND=$(fm_watcher_stall_bound) CYCLE_LOG="$STATE/.watch-cycle-exits.log" CYCLE_LOG_LOCK="$STATE/.watch-cycle-exits.lock" CYCLE_LOG_MAX_BYTES=${FM_WATCH_CYCLE_LOG_MAX_BYTES:-262144} @@ -343,12 +355,28 @@ close_unobserved_cycle() { return 1 } +# True while is alive and this home's watcher lock still names it under +# the identity this arm's current cycle attached to, whatever its beacon age. +attached_holder_live() { + local pid=$1 lock_pid + lock_pid=$(cat "$WATCH_LOCK/pid" 2>/dev/null || true) + [ "$lock_pid" = "$pid" ] || return 1 + fm_pid_alive "$pid" || return 1 + fm_watcher_lock_matches_pid "$STATE" "$WATCH" "$pid" "$FM_HOME" || return 1 + [ "$FM_WATCHER_MATCHED_IDENTITY" = "$cycle_watcher_identity" ] +} + # Stay alive across identity-matched healthy holders. If one cycle ends, attach # to a verified successor. With no successor, report the wake that cycle durably # delivered, or fail loudly - never a clean empty completion that an adapter could # mistake for a no-op. +# A stale beacon alone does not end the followed cycle: while the holder is alive +# and the lock still names it under the same identity, it is a slow cycle, which +# a started arm tolerates by waiting on its child, so this arm keeps following it. +# Only at the stall bound, where the watcher's own re-arm evicts a live holder, +# does it fail with the typed stalled-holder line so its owner's retry replaces it. attach_and_wait() { - local attached_pid=$1 + local attached_pid=$1 age while :; do if healthy_watcher; then if [ "$HEALTHY_PID" != "$attached_pid" ] || [ "$HEALTHY_IDENTITY" != "$cycle_watcher_identity" ]; then @@ -360,6 +388,16 @@ attach_and_wait() { sleep "$ATTACH_POLL" continue fi + if attached_holder_live "$attached_pid"; then + age=$(fm_path_age "$BEAT") + if [ "$age" -lt "$STALL_BOUND" ]; then + sleep "$ATTACH_POLL" + continue + fi + cycle_log_append unknown unknown attached-holder-stalled none + echo "watcher: FAILED - attached watcher pid=$attached_pid stalled (beacon ${age}s at or past hard bound ${STALL_BOUND}s)" + return 1 + fi if wait_for_healthy_successor; then cycle_log_append unknown unknown attached-cycle-ended "attached:$HEALTHY_PID" attached_pid=$HEALTHY_PID diff --git a/bin/fm-watch.sh b/bin/fm-watch.sh index 7fa311fd2a2..96bae225fa5 100755 --- a/bin/fm-watch.sh +++ b/bin/fm-watch.sh @@ -279,7 +279,9 @@ WATCHER_STALE_GRACE=${FM_WATCHER_STALE_GRACE:-${FM_GUARD_GRACE:-$(fm_poll_derive # for inspection (the grace above); at or past it the re-arm evicts the holder # instead, because a watcher whose beacon has stalled that long is not polling # and nothing else would ever replace it (evict_stalled_holder below). -WATCHER_STALL_BOUND=${FM_WATCHER_STALL_BOUND:-$((WATCHER_STALE_GRACE * 3))} +# fm_watcher_stall_bound (bin/fm-wake-lib.sh) owns the derivation, shared with +# the arm that follows this watcher. +WATCHER_STALL_BOUND=$(fm_watcher_stall_bound "$POLL") HEARTBEAT=${FM_HEARTBEAT:-600} # base seconds between heartbeat scans HEARTBEAT_MAX=${FM_HEARTBEAT_MAX:-7200} # heartbeat backoff cap CHECK_INTERVAL=${FM_CHECK_INTERVAL:-300} # seconds between *.check.sh sweeps diff --git a/docs/configuration.md b/docs/configuration.md index 209da6093a2..c97571fc862 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -2326,7 +2326,7 @@ FM_CLAUDE_AUTOARM_SYNC_WAIT_MS=800 # milliseconds the --claude turn-end guard FM_CLAUDE_AUTOARM_EPOCH_FRESH=15 # seconds a recorded auto-arm outcome remains eligible for the current event epoch's recovery or failure decision FM_CLAUDE_TURNEND_BLOCK_BUDGET=3 # consecutive --claude guard re-blocks before the verified one-time attended fail-open; safely below Claude Code's 8-block override FM_ARM_CONFIRM_TIMEOUT=10 # seconds fm-watch-arm waits to confirm a fresh watcher before reporting FAILED; default 30 on Git Bash/MSYS -FM_ARM_ATTACH_POLL=0.5 # seconds between checks while fm-watch-arm is attached to an existing healthy watcher cycle +FM_ARM_ATTACH_POLL=0.5 # seconds between checks while fm-watch-arm follows an attached watcher cycle (bin/fm-watch-arm.sh header) FM_OPENCODE_ARM_READY_TIMEOUT_MS=12000 # milliseconds the OpenCode primary watcher plugin waits for an arm attempt to report started, healthy, wake, or failure; default 35000 on Windows to stay above the MSYS confirm budget FM_PI_ARM_READY_TIMEOUT_MS=12000 # milliseconds the Pi watcher extension waits for a successor arm to report started or attached; default 35000 on Windows to stay above the MSYS confirm budget FM_WATCH_ARM_RETIRE_TIMEOUT_MS=1000 # milliseconds Pi/OpenCode wait for an unready successor arm to exit before abandoning retries @@ -2335,8 +2335,8 @@ FM_WATCH_REARM_RETRY_MAX_MS=4000 # Pi/OpenCode adapter cap for exponential con FM_WATCH_REARM_RETRY_LIMIT=5 # Pi/OpenCode adapter launch-failure retries before surfacing restoration failure FM_WATCH_CYCLE_LOG_MAX_BYTES=262144 # size cap for the arm-owned watcher lifecycle ledger FM_WATCH_CYCLE_LOG_KEEP_LINES=1000 # newest complete lifecycle rows considered when the ledger is capped -FM_WATCHER_STALE_GRACE=300 # defaults to FM_GUARD_GRACE if set, else the poll-derived grace (docs/turnend-guard.md "Guard grace and the poll cadence"); seconds a live watcher lock may have a stale beacon before re-arm errors -FM_WATCHER_STALL_BOUND= # defaults to 3x FM_WATCHER_STALE_GRACE; a live holder whose beacon is stale past this hard bound is evicted with TERM and replaced by the re-arm rather than refused (docs/turnend-guard.md, bin/fm-watch.sh header) +FM_WATCHER_STALE_GRACE=300 # defaults to FM_GUARD_GRACE if set, else the poll-derived grace (docs/turnend-guard.md "Guard grace and the poll cadence"); seconds before a fresh arm refuses a live holder's stale beacon (attached arms: FM_WATCHER_STALL_BOUND) +FM_WATCHER_STALL_BOUND= # live-holder stall bound; default and arm/re-arm behavior: docs/turnend-guard.md "Guard grace and the poll cadence" FM_SIGNAL_GRACE=30 # seconds to coalesce nearby status and turn-end signals into one wake FM_WATCHER_CLEANUP_LOCK_BOUND= # optional watcher EXIT marker-lock wait; default and validation: docs/watcher-continuity.md FM_TURNEND_CHURN_ABSORB_SECS=900 # longest one endpoint's bare turn-ends may be deferred on pane-churn evidence alone; only consulted when config/turnend-churn-absorb is present diff --git a/docs/remote-secondmates.md b/docs/remote-secondmates.md index 8c2d36ecb5b..42d496acb39 100644 --- a/docs/remote-secondmates.md +++ b/docs/remote-secondmates.md @@ -99,6 +99,7 @@ It distinguishes preemption from a wait window that closes with no data: - Only a genuinely quiet window proves channel freshness. - Either outcome can re-arm without losing data. +- The parent's reply listener polls again under the same claim after either one, so a same-home command such as the per-cycle liveness probe never tears the listener down; [`bin/fm-procevent-remote-reply.sh`](../bin/fm-procevent-remote-reply.sh) owns that mapping. ### Cancelled and orphaned jobs diff --git a/docs/turnend-guard.md b/docs/turnend-guard.md index 7f26471851f..23af69a0eb1 100644 --- a/docs/turnend-guard.md +++ b/docs/turnend-guard.md @@ -231,6 +231,9 @@ Once the live holder's beacon is stale past `FM_WATCHER_STALL_BOUND` (default th A watcher wedged mid-cycle can therefore no longer refuse every replacement indefinitely. `bin/fm-watch.sh`'s header owns the exact wording and the survives-TERM fallback. +Below that bound a stale beacon alone does not end an attached arm's watch of a live, identity-matched holder; a changed lock can end it sooner. +At the bound the arm reports a typed stalled-holder failure so its owner's retry can replace the holder. +`fm_watcher_stall_bound` in `bin/fm-wake-lib.sh` owns the shared derivation; `bin/fm-watch-arm.sh`'s header owns the exact attached-arm close behavior. The auto-arm hook additionally exports its resolved `FM_GUARD_GRACE` when it forks `bin/fm-watch-arm.sh`. The arm wrapper and the watcher it may start then judge staleness with the exact same value the hook just judged it with, whether that value came from an operator override or the poll-derived default. diff --git a/tests/fm-pending-reply.test.sh b/tests/fm-pending-reply.test.sh index cd31fbaf552..1a14b49a658 100755 --- a/tests/fm-pending-reply.test.sh +++ b/tests/fm-pending-reply.test.sh @@ -1086,6 +1086,94 @@ test_tick_skips_terminal_and_reuses_target_observation() { pass "tick skips terminal records and reuses target observations" } +# Records are never pruned, so a home accumulates thousands of settled ones. The +# tick selects the records it has work for in one pass and leaves every settled +# record alone: it must not block on a settled record's per-record lock that a +# live foreign process holds, and it still does the work the selected records need. +test_tick_leaves_settled_records_alone() { + local home state settled closed open_esc awaiting rec i copy holder tick_pid ticked=0 open lib + local sums_before sums_after holder_lock_pid + home=$(setup_parent settled-store) + state="$home/state" + # Reset the fixture clock after isolated subshell tests. + # shellcheck disable=SC2031 + export FM_PENDING_REPLY_NOW=5200 + # A resolved record that never escalated, and one whose escalation closed. + settled=$(fm_pending_reply_create "$home" "$state" hibit "settled request") + fm_pending_reply_mark_delivered "$state" "$settled" + printf 'done [corr=%s]: settled reply\n' "$settled" >> "$state/hibit.status" + fm_pending_reply_try_resolve "$state" "$settled" || fail "settled fixture should resolve" + closed=$(fm_pending_reply_create "$home" "$state" hibit "closed escalation") + fm_pending_reply_mark_delivered "$state" "$closed" + rec=$(fm_pending_reply_path "$state" "$closed") + fm_pending_reply_set "$rec" phase escalated + fm_pending_reply_set "$rec" escalated_epoch 5100 + printf 'blocked [key=pending-reply-%s]: pending-reply-missed: task=hibit pending-reply-id=%s request=closed escalation\n' \ + "$closed" "$closed" >> "$state/hibit.status" + printf 'done [corr=%s]: late reply\n' "$closed" >> "$state/hibit.status" + fm_pending_reply_try_resolve "$state" "$closed" || fail "closed-escalation fixture should resolve" + [ -n "$(fm_pending_reply_get "$rec" escalation_closed_epoch)" ] || fail "fixture escalation did not close" + # Many settled copies, as a long-lived home accumulates. + i=0 + while [ "$i" -lt 300 ]; do + copy=$(printf '%016x' $((0x5e7700000000 + i))) + for rec in "$settled" "$closed"; do + sed "s/^corr_id=.*/corr_id=$copy/" "$(fm_pending_reply_path "$state" "$rec")" \ + > "$(fm_pending_reply_path "$state" "$copy")" + copy=$(printf '%016x' $((0x5e7780000000 + i))) + done + i=$((i + 1)) + done + # Work the tick still owes: a resolved record whose escalation close did not + # land, and a delivered request whose correlated report is in the parent status. + open_esc=$(fm_pending_reply_create "$home" "$state" esc "open escalation") + fm_pending_reply_mark_delivered "$state" "$open_esc" + rec=$(fm_pending_reply_path "$state" "$open_esc") + printf 'blocked [key=pending-reply-%s]: pending-reply-missed: task=esc pending-reply-id=%s request=open escalation\n' \ + "$open_esc" "$open_esc" > "$state/esc.status" + printf 'done [corr=%s]: late reply\n' "$open_esc" >> "$state/esc.status" + fm_pending_reply_set "$rec" escalated_epoch 5150 + fm_pending_reply_set "$rec" resolved_via status + fm_pending_reply_set "$rec" phase resolved + awaiting=$(fm_pending_reply_create "$home" "$state" open "awaiting report") + fm_pending_reply_mark_delivered "$state" "$awaiting" + printf 'done [corr=%s]: the report\n' "$awaiting" > "$state/open.status" + sums_before=$(cd "$(fm_pending_reply_dir "$state")" && cksum 00005e77* "$settled" "$closed") + [ "$(printf '%s\n' "$sums_before" | wc -l | tr -d ' ')" -eq 602 ] || fail "settled fixture store is incomplete" + + # A live foreign process holds one settled record's per-record lock. + lib="$ROOT/bin/fm-wake-lib.sh" + bash -c '. "$1"; fm_lock_acquire_wait "$2" && : > "$3"; exec sleep 300' _ \ + "$lib" "$state/.pending-reply-00005e7700000000.lock" "$home/held" & + holder=$! + for _ in $(seq 1 100); do [ -e "$home/held" ] && break; sleep 0.1; done + [ -e "$home/held" ] || { kill "$holder" 2>/dev/null; fail "foreign holder never took the lock"; } + + fm_pending_reply_tick "$state" & + tick_pid=$! + for _ in $(seq 1 600); do + case "$(ps -p "$tick_pid" -o stat= 2>/dev/null)" in ''|Z*) ticked=1; break ;; esac + sleep 0.1 + done + [ "$ticked" = 1 ] || kill -TERM "$tick_pid" 2>/dev/null + wait "$tick_pid" 2>/dev/null || true + holder_lock_pid=$(cat "$state/.pending-reply-00005e7700000000.lock/pid" 2>/dev/null || true) + kill -TERM "$holder" 2>/dev/null + wait "$holder" 2>/dev/null || true + + [ "$ticked" = 1 ] || fail "the tick blocked on a settled record's foreign-held lock" + [ "$holder_lock_pid" = "$holder" ] || fail "the tick disturbed the foreign holder's lock (pid=$holder_lock_pid)" + sums_after=$(cd "$(fm_pending_reply_dir "$state")" && cksum 00005e77* "$settled" "$closed") + [ "$sums_before" = "$sums_after" ] || fail "the tick rewrote settled records" + [ -n "$(fm_pending_reply_get "$(fm_pending_reply_path "$state" "$open_esc")" escalation_closed_epoch)" ] \ + || fail "the tick did not close the resolved record's open escalation" + open=$(status_open_decisions "$state/esc.status") + [ -z "$open" ] || fail "the resolved record's escalation stayed open: $open" + [ "$(phase_of "$state" "$awaiting")" = resolved ] \ + || fail "the tick did not resolve the awaiting record from its correlated report" + pass "the tick leaves settled records alone and still does the selected records' work" +} + test_correlations_reuse_only_for_matching_open_task() { local dir fb log home state got corr1 corr2 corr3 rec dir="$TMP_ROOT/corr-reuse"; mkdir -p "$dir" @@ -1629,6 +1717,7 @@ test_busy_idle_observation_via_backend_abstraction test_unknown_backend_state_uses_capture_fallback test_kimi_capture_fallback_uses_recorded_harness test_tick_skips_terminal_and_reuses_target_observation +test_tick_leaves_settled_records_alone test_correlations_reuse_only_for_matching_open_task test_tick_end_to_end_missed_then_escalate test_failed_send_discards_undelivered_expectation diff --git a/tests/fm-remote-reply.test.sh b/tests/fm-remote-reply.test.sh index 0ea7a425da0..47f75a75593 100755 --- a/tests/fm-remote-reply.test.sh +++ b/tests/fm-remote-reply.test.sh @@ -847,11 +847,56 @@ set +e wait "$PREEMPTED_SOURCE" preempted_rc=$? set -e -[ "$preempted_rc" -eq "$FM_REMOTE_JOB_PREEMPTED_EXIT" ] \ - || fail "the reply poll did not expose remote-job preemption: $preempted_rc" +[ "$preempted_rc" -eq 75 ] \ + || fail "a preempted reply poll did not report a closed window: $preempted_rc" assert_absent "$PARENT/state/remote-replies/ios.caught-up" \ "a preempted reply poll published a caught-up watermark" -pass "a preempted reply poll cannot publish channel freshness" +pass "a preempted reply poll reports a closed window without publishing channel freshness" + +# The per-cycle liveness probe is a non-preemptible job for the same remote home, +# so the job worker preempts the listener's long-poll on every watcher cycle. +# That must not cost the listener: it keeps its claim and polls again, and the +# watcher's reconcile has nothing to relaunch. +: > "$TMP_ROOT/preempted-polls" +FM_REMOTE_REPLY_POLL_LOG="$TMP_ROOT/preempted-polls" FM_PROCEVENT_LAUNCH_FLOOR_SECONDS=1 \ + remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" >/dev/null 2>&1 & +PREEMPTED_RUNNER=$! +wait_for "$CLAIMS/$SID.claim" || fail "the preempted-listener case never claimed the source" +HELD_PID=$(sed -n '2p' "$CLAIMS/$SID.claim") +running_poll='' +for _ in $(seq 1 100); do + for job in "$TMP_ROOT"/remote-jobs/jobs/job-*; do + [ -d "$job" ] || continue + if [ "$(fm_remote_job_read_state "$job" 2>/dev/null || true)" = running ]; then + running_poll=$job + break 2 + fi + done + sleep 0.05 +done +[ -n "$running_poll" ] || fail "the listener's poll did not begin running before preemption" +remote_env "$ROOT/bin/fm-on.sh" ios fm-remote-file.sh get data/reply/report.md 262144 >/dev/null +polls=0 +for _ in $(seq 1 120); do + polls=$(wc -l < "$TMP_ROOT/preempted-polls" | tr -d ' ') + [ "$polls" -ge 2 ] && break + sleep 0.25 +done +[ "$polls" -ge 2 ] || fail "the preempted listener did not poll again" +case "$(ps -p "$PREEMPTED_RUNNER" -o stat= 2>/dev/null)" in + ''|Z*) fail "a preempted poll ended the reply listener" ;; +esac +[ "$(reply_owner)" = live ] || fail "a preempted poll released the listener's claim" +[ "$(sed -n '2p' "$CLAIMS/$SID.claim")" = "$HELD_PID" ] \ + || fail "a preempted poll replaced the reply listener" +reconcile_out=$(remote_env "$ROOT/bin/fm-procevent.sh" reconcile) +assert_contains "$reconcile_out" 'started=0' \ + "reconcile relaunched a listener after a preempted poll" +[ "$(sed -n '2p' "$CLAIMS/$SID.claim")" = "$HELD_PID" ] \ + || fail "reconcile replaced the preempted listener" +stop_reply_listener || fail "the preempted listener did not stop" +wait "$PREEMPTED_RUNNER" 2>/dev/null || true +pass "a preempted reply poll keeps its listener and reconcile launches nothing" # A quiet window is the one moment this channel can prove it is NOT behind, and # the parent's pending-reply guard needs that proof: a remote report that exists diff --git a/tests/fm-watch-arm.test.sh b/tests/fm-watch-arm.test.sh index b640a60ffa4..a2c13f946f6 100755 --- a/tests/fm-watch-arm.test.sh +++ b/tests/fm-watch-arm.test.sh @@ -48,9 +48,9 @@ SEED_PID= ARM_PID= # Start the real watcher as the singleton holder. -start_seed_watcher() { # - local state=$1 fakebin=$2 out=$3 i - PATH="$fakebin:$PATH" FM_STATE_OVERRIDE="$state" FM_POLL=5 FM_SIGNAL_GRACE=1 \ +start_seed_watcher() { # [poll-seconds] + local state=$1 fakebin=$2 out=$3 poll=${4:-5} i + PATH="$fakebin:$PATH" FM_STATE_OVERRIDE="$state" FM_POLL="$poll" FM_SIGNAL_GRACE=1 \ FM_CHECK_INTERVAL=999999 FM_HEARTBEAT=999999 "$WATCH" > "$out" & SEED_PID=$! i=0 @@ -267,6 +267,129 @@ test_attached_arm_still_fails_on_a_wake_it_did_not_deliver() { pass "watch-arm: a cycle that delivered no wake of its own still fails loudly" } +# A slow cycle is not an ended cycle. The holder is frozen past the grace plus +# the successor confirmation window, which is where an attached arm used to +# declare the cycle over and fail while the holder was alive and still held the +# lock; the owner's retry then hit that live holder's refusal (the auto-arm +# FAILED notice), or, if the holder beat again first, nothing followed it at all. +test_attached_arm_follows_a_slow_live_holder() { + local dir state fakebin out armout status i arm_followed armout_frozen ledger_frozen + dir=$(make_case attached-slow-holder) + state="$dir/state" + fakebin="$dir/fakebin" + out="$dir/watch.out" + armout="$dir/arm.out" + start_seed_watcher "$state" "$fakebin" "$out" 1 + PATH="$fakebin:$PATH" FM_STATE_OVERRIDE="$state" FM_ARM_ATTACH_POLL=0.1 \ + FM_ARM_CONFIRM_TIMEOUT=1 FM_GUARD_GRACE=4 FM_WATCHER_STALL_BOUND=600 "$WATCH_ARM" > "$armout" & + ARM_PID=$! + i=0 + while [ "$i" -lt 80 ]; do + grep -qF "watcher: attached pid=$SEED_PID" "$armout" 2>/dev/null && break + sleep 0.1 + i=$((i + 1)) + done + grep -qF "watcher: attached pid=$SEED_PID" "$armout" \ + || fail "arm did not attach to the live watcher: $(cat "$armout")" + + kill -STOP "$SEED_PID" + # Grace 4s plus the 1s confirmation window plus its rounding second is where + # the old arm gave up; hold the holder well past that. Observe while frozen, + # but resume before asserting so a failure never strands a stopped watcher. + i=0 + while [ "$(FM_STATE_OVERRIDE="$state" bash -c '. "$1"; fm_path_age "$2"' _ "$ROOT/bin/fm-wake-lib.sh" "$state/.last-watcher-beat")" -lt 10 ] \ + && [ "$i" -lt 300 ]; do + sleep 0.1 + i=$((i + 1)) + done + arm_followed=0 + is_live_non_zombie "$ARM_PID" && arm_followed=1 + armout_frozen=$(cat "$armout") + ledger_frozen=$(cat "$state/.watch-cycle-exits.log" 2>/dev/null || true) + kill -CONT "$SEED_PID" + [ "$arm_followed" = 1 ] || fail "attached arm ended while its holder was alive: $armout_frozen" + assert_not_contains "$armout_frozen" 'watcher: FAILED' \ + "attached arm failed a live holder's slow cycle" + assert_not_contains "$ledger_frozen" 'reason=attached-cycle-ended' \ + "attached arm closed a cycle that had not ended" + + # The holder resumes and delivers a wake: the arm that kept following it + # reports that wake, so nothing is lost. + printf 'needs-decision: which export format?\n' > "$state/demo.status" + wait_for_exit "$SEED_PID" 150 + grep -q '^signal:' "$out" || fail "resumed holder did not surface the signal wake: $(cat "$out")" + wait_for_exit "$ARM_PID" 150 + status=$? + ! grep -qF 'watcher: FAILED' "$armout" \ + || fail "attached arm failed after its holder resumed: $(cat "$armout")" + grep -q '^signal:' "$armout" \ + || fail "attached arm did not report the resumed holder's wake: $(cat "$armout")" + expect_code 0 "$status" "an attached arm that followed a slow holder must close with its wake" + pass "watch-arm: an attached arm keeps following a slow live holder and reports its wake" +} + +# The stall bound is where following ends. A live holder whose beacon reaches it +# is what the watcher's own re-arm evicts, so the attached arm stops there with +# the typed stalled-holder line, and its owner's retry replaces the holder +# instead of being refused. +test_attached_arm_hands_a_stalled_holder_to_its_replacement() { + local dir state fakebin armout rearmout holder identity status + dir=$(make_case attached-stalled-holder) + state="$dir/state" + fakebin="$dir/fakebin" + armout="$dir/arm.out" + rearmout="$dir/rearm.out" + # A live process the lock records under its real identity, which never beats: + # the shape of a watcher wedged mid-cycle that still answers TERM. + sleep 300 & + holder=$! + identity=$(FM_STATE_OVERRIDE="$state" bash -c '. "$1"; fm_pid_identity "$2"' _ "$ROOT/bin/fm-wake-lib.sh" "$holder") \ + || fail "could not identify the fake holder" + mkdir -p "$state/.watch.lock" + printf '%s\n' "$holder" > "$state/.watch.lock/pid" + printf '%s\n' "$dir" > "$state/.watch.lock/fm-home" + printf '%s\n' "$WATCH" > "$state/.watch.lock/watcher-path" + printf '%s\n' "$identity" > "$state/.watch.lock/pid-identity" + : > "$state/.last-watcher-beat" + + PATH="$fakebin:$PATH" FM_HOME="$dir" FM_STATE_OVERRIDE="$state" FM_ARM_ATTACH_POLL=0.1 \ + FM_ARM_CONFIRM_TIMEOUT=1 FM_GUARD_GRACE=5 FM_WATCHER_STALL_BOUND=12 "$WATCH_ARM" > "$armout" & + ARM_PID=$! + wait_for_exit "$ARM_PID" 300 + status=$? + grep -qF "watcher: attached pid=$holder" "$armout" \ + || fail "arm did not attach to the fresh holder: $(cat "$armout")" + grep -E "^watcher: FAILED - attached watcher pid=$holder stalled \(beacon [0-9]+s at or past hard bound 12s\)\$" "$armout" >/dev/null \ + || fail "attached arm did not report the stalled holder: $(cat "$armout")" + ! grep -qF 'cycle ended without an actionable reason' "$armout" \ + || fail "attached arm gave up on the live holder before the stall bound: $(cat "$armout")" + [ "$status" -ne 0 ] && [ "$status" -ne 124 ] \ + || fail "stalled-holder close did not exit nonzero (status $status)" + grep -q 'reason=attached-holder-stalled' "$state/.watch-cycle-exits.log" \ + || fail "the stalled-holder close was not classified in the lifecycle ledger" + is_live_non_zombie "$holder" || fail "the attached arm signalled the holder it follows" + + # The owner's retry: a fresh arm reaches the watcher's eviction path, and the + # replacement surfaces an ordinary wake instead of the refusal that used to + # end in the auto-arm FAILED notice. + PATH="$fakebin:$PATH" FM_HOME="$dir" FM_STATE_OVERRIDE="$state" \ + FM_POLL=1 FM_SIGNAL_GRACE=0 FM_CHECK_INTERVAL=999999 FM_HEARTBEAT=999999 \ + FM_ARM_CONFIRM_TIMEOUT=10 FM_GUARD_GRACE=5 FM_WATCHER_STALL_BOUND=12 "$WATCH_ARM" > "$rearmout" 2>&1 & + ARM_PID=$! + wait_for_exit "$ARM_PID" 300 + status=$? + grep -qF "watcher: replaced stalled pid $holder " "$rearmout" \ + || fail "the retry did not replace the stalled holder: $(cat "$rearmout")" + ! grep -qF 'watcher: FAILED' "$rearmout" \ + || fail "the retry failed instead of replacing the stalled holder: $(cat "$rearmout")" + grep -Eq '^(signal|stale|check):' "$rearmout" \ + || fail "the replacement surfaced no ordinary wake: $(cat "$rearmout")" + expect_code 0 "$status" "the retry that replaced a stalled holder must close with an ordinary wake" + wait_for_pid_gone "$holder" 50 || fail "the stalled holder survived its replacement" + wait "$holder" 2>/dev/null || true + pass "watch-arm: an attached arm hands a holder stalled past the bound to its owner's replacement" +} + test_rearm_resurfaces_durable_queue_and_remote_open_decision() { local dir home state fakebin result armout drainout status watcher_pid sequence generation decision_recovery_arm decision_successor dir=$(make_case rearm-resurface) @@ -1202,6 +1325,8 @@ test_watcher_exits_when_its_state_directory_is_removed test_watcher_exits_when_its_home_is_removed test_reaper_stops_a_tracked_watcher test_attached_arm_still_fails_on_a_wake_it_did_not_deliver +test_attached_arm_follows_a_slow_live_holder +test_attached_arm_hands_a_stalled_holder_to_its_replacement test_rearm_resurfaces_durable_queue_and_remote_open_decision test_slow_rearm_recovery_is_still_surfaced test_marker_publish_failure_retains_recovery_evidence