Skip to content
Merged
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
46 changes: 43 additions & 3 deletions bin/fm-pending-reply-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -1445,20 +1445,60 @@ fm_pending_reply_tick_one() { # <state-dir> <corr_id> <busy_state> [secondmate-
return 0
}

# Print, one per line, the records among <record-path>... 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() { # <record-path>...
[ "$#" -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() { # <state-dir>
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)
Expand Down
14 changes: 13 additions & 1 deletion bin/fm-procevent-remote-reply.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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.
#
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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"
}
Expand Down
12 changes: 12 additions & 0 deletions bin/fm-wake-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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 <state>
# 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;
Expand Down
42 changes: 40 additions & 2 deletions bin/fm-watch-arm.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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=<N> stalled (beacon <age>s at or past hard bound <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
Expand Down Expand Up @@ -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}
Expand Down Expand Up @@ -343,12 +355,28 @@ close_unobserved_cycle() {
return 1
}

# True while <pid> 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
Expand All @@ -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
Expand Down
4 changes: 3 additions & 1 deletion bin/fm-watch.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 3 additions & 3 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
1 change: 1 addition & 0 deletions docs/remote-secondmates.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
3 changes: 3 additions & 0 deletions docs/turnend-guard.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading