diff --git a/bin/fm-backlog-handoff.sh b/bin/fm-backlog-handoff.sh index fa729c9d1b6..879de6053db 100755 --- a/bin/fm-backlog-handoff.sh +++ b/bin/fm-backlog-handoff.sh @@ -549,7 +549,7 @@ remote_deliver_outbox() { # mv -f -- "$counter_tmp" "$counter" \ || { rm -f -- "$snapshot" "$counter_tmp"; return 1; } remote_rel="state/handoff/$id.outbox.md" - if ! "$SCRIPT_DIR/fm-on.sh" "$id" fm-remote-file.sh put "$remote_rel" 1048576 \ + if ! "$SCRIPT_DIR/fm-on.sh" --stdin "$id" fm-remote-file.sh put "$remote_rel" 1048576 \ "$bytes" "$hash" "$generation" < "$snapshot"; then rm -f -- "$snapshot" echo "error: handoff transfer to $id was unavailable or completion is unknown; outbox preserved at $outbox" >&2 diff --git a/bin/fm-on.sh b/bin/fm-on.sh index 5e24f2cef1d..eff02f7c350 100755 --- a/bin/fm-on.sh +++ b/bin/fm-on.sh @@ -2,7 +2,7 @@ # Execute one tracked Firstmate command in a configured remote secondmate home. # # Usage: -# fm-on.sh [args...] +# fm-on.sh [--stdin] [args...] # # Routes come only from remote records in data/secondmates.md. A record names an # SSH config alias, remote Firstmate code root, and remote FM_HOME. A host alias @@ -11,11 +11,14 @@ # bin/fm-*.sh namespace. No per-command table exists. # # argv is encoded as one NUL-delimited stream and passed through the fixed -# fm-remote-entrypoint.sh. stdin remains the caller's stdin, stdout and stderr -# remain separate, and ssh's exit status is returned unchanged. OpenSSH never -# receives an auto-retry instruction here. Exit 255 therefore means unavailable -# transport or unknown remote completion and must be reconciled by the semantic -# caller, never blindly repeated by this layer. +# fm-remote-entrypoint.sh. The remote command's stdin is /dev/null by default, +# because remote staging captures stdin to EOF and an open caller stream would +# block staging indefinitely; a payload caller passes --stdin to forward its +# own stream as the job's bounded input. stdout and stderr remain separate, and +# ssh's exit status is returned unchanged. OpenSSH never receives an auto-retry +# instruction here. Exit 255 therefore means unavailable transport or unknown +# remote completion and must be reconciled by the semantic caller, never +# blindly repeated by this layer. # # The SSH alias keeps normal public-key and strict host-key policy in ~/.ssh. # This command explicitly disables agent forwarding, forwarding setup, and @@ -42,12 +45,17 @@ PROTOCOL=1 . "$SCRIPT_DIR/fm-secondmate-registry-lib.sh" die() { printf 'error: %s\n' "$1" >&2; exit 1; } -usage() { sed -n '2,23p' "$0" | sed 's/^# \{0,1\}//'; exit 2; } +usage() { sed -n '2,25p' "$0" | sed 's/^# \{0,1\}//'; exit 2; } encode_base64() { base64 | tr -d '\n' } +STDIN_MODE=closed +if [ "${1:-}" = --stdin ]; then + STDIN_MODE=caller + shift +fi [ "$#" -ge 2 ] || usage ROUTE=$1 COMMAND=$2 @@ -103,10 +111,15 @@ case "$ALIVE_COUNT_MAX" in ''|*[!0-9]*) die "FM_SSH_ALIVE_COUNT_MAX must be a po [ "$ALIVE_INTERVAL" -gt 0 ] || die "FM_SSH_ALIVE_INTERVAL must be a positive integer: $ALIVE_INTERVAL" [ "$ALIVE_COUNT_MAX" -gt 0 ] || die "FM_SSH_ALIVE_COUNT_MAX must be a positive integer: $ALIVE_COUNT_MAX" -"$SSH_BIN" \ - -o ForwardAgent=no \ - -o ClearAllForwardings=yes \ - -o 'SendEnv=-*' \ - -o "ServerAliveInterval=$ALIVE_INTERVAL" \ - -o "ServerAliveCountMax=$ALIVE_COUNT_MAX" \ +SSH_ARGS=( + -o ForwardAgent=no + -o ClearAllForwardings=yes + -o 'SendEnv=-*' + -o "ServerAliveInterval=$ALIVE_INTERVAL" + -o "ServerAliveCountMax=$ALIVE_COUNT_MAX" -- "$HOST" fm-remote-entrypoint.sh "$PROTOCOL" "$ROOT_B64" "$HOME_B64" "$ARGV_B64" +) +if [ "$STDIN_MODE" = caller ]; then + exec "$SSH_BIN" "${SSH_ARGS[@]}" +fi +exec "$SSH_BIN" "${SSH_ARGS[@]}" < /dev/null diff --git a/bin/fm-remote-entrypoint.sh b/bin/fm-remote-entrypoint.sh index 6763e8c955d..4549ff6e9ca 100755 --- a/bin/fm-remote-entrypoint.sh +++ b/bin/fm-remote-entrypoint.sh @@ -19,6 +19,15 @@ # disconnect remains unknown completion to fm-on.sh, which preserves OpenSSH's # exit 255 behavior. The shared library header owns job fields, bounds, PATH, # LaunchAgent contract, and worker environment. +# +# A staged job whose caller goes away is cancelled rather than abandoned: any +# exit after staging and before the published result marks the job cancelled +# (signal traps cover a delivered HUP/TERM/PIPE/INT, and the exit trap covers a +# failed bounded wait), and while waiting this process probes its parent about +# once per second, so an ssh channel that dies without delivering any signal - +# sshd exiting and reparenting this process - also cancels the job. The worker +# then skips or stops the cancelled job instead of running it to completion for +# nobody. set -eu PROTOCOL=1 @@ -74,7 +83,36 @@ sha256_file() { # [ "$#" -eq 4 ] || die "remote entrypoint expects protocol, root, home, and argv" [ "$1" = "$PROTOCOL" ] || die "incompatible remote protocol: local=$1 remote=$PROTOCOL" TMP=$(mktemp -d "${TMPDIR:-/tmp}/fm-remote-entrypoint.XXXXXX") || die "cannot create protocol staging directory" 70 -trap 'rm -rf -- "$TMP"' EXIT + +JOB_ID= +JOB_COMPLETED=0 +ACCOUNT_HOME= +ENTRYPOINT_PPID=$(ps -o ppid= -p $$ 2>/dev/null | tr -d ' ' || true) + +# The recorded parent is the ssh session process; when it disappears this +# process is reparented and the caller is provably gone. An unreadable probe +# never cancels: only an observed parent change does. +# shellcheck disable=SC2329 # Invoked by fm_remote_job_wait through FM_REMOTE_JOB_DISCONNECT_PROBE. +entrypoint_caller_connected() { + local current + case "$ENTRYPOINT_PPID" in ''|*[!0-9]*) return 0 ;; esac + current=$(ps -o ppid= -p $$ 2>/dev/null | tr -d ' ' || true) + case "$current" in ''|*[!0-9]*) return 0 ;; esac + [ "$current" = "$ENTRYPOINT_PPID" ] +} + +# shellcheck disable=SC2329 # Invoked through the EXIT trap below. +entrypoint_cleanup() { + rm -rf -- "$TMP" + if [ -n "$JOB_ID" ] && [ "$JOB_COMPLETED" -eq 0 ] && [ -n "$ACCOUNT_HOME" ]; then + fm_remote_job_cancel "$ACCOUNT_HOME" "$JOB_ID" 2>/dev/null || true + fi +} +trap entrypoint_cleanup EXIT +trap 'exit 129' HUP +trap 'exit 130' INT +trap 'exit 141' PIPE +trap 'exit 143' TERM decode_text "remote root" "$2" "$TMP/root" decode_text "remote home" "$3" "$TMP/home" @@ -138,11 +176,14 @@ if ! fm_remote_job_ensure_worker "$ROOT" "$ACCOUNT_HOME"; then die "${FM_REMOTE_JOB_ERROR:-remote job worker is unavailable; run fm-on.sh fm-remote-doctor.sh --fix}" fi if ! JOB_ID=$(fm_remote_job_stage "$ACCOUNT_HOME" "$ROOT" "$HOME_PATH" "$COMMAND" "${ARGV[@]:1}"); then + JOB_ID= die "${FM_REMOTE_JOB_ERROR:-cannot stage remote job}" 70 fi +FM_REMOTE_JOB_DISCONNECT_PROBE=entrypoint_caller_connected if ! fm_remote_job_wait "$ACCOUNT_HOME" "$JOB_ID"; then die "${FM_REMOTE_JOB_ERROR:-remote job did not complete}" 70 fi +JOB_COMPLETED=1 cat "$FM_REMOTE_JOB_STDOUT" cat "$FM_REMOTE_JOB_STDERR" >&2 RESULT=$FM_REMOTE_JOB_EXIT diff --git a/bin/fm-remote-home-seed.sh b/bin/fm-remote-home-seed.sh index a679851cbc4..7deafc40dcf 100755 --- a/bin/fm-remote-home-seed.sh +++ b/bin/fm-remote-home-seed.sh @@ -242,7 +242,7 @@ if [ "$PREFLIGHT_RC" -ne 0 ]; then fi set +e -PROVISION_OUT=$("$SCRIPT_DIR/fm-on.sh" "$ID" fm-remote-home-provision.sh < "$TMP/manifest" 2>&1) +PROVISION_OUT=$("$SCRIPT_DIR/fm-on.sh" --stdin "$ID" fm-remote-home-provision.sh < "$TMP/manifest" 2>&1) PROVISION_RC=$? set -e if [ "$PROVISION_RC" -ne 0 ]; then diff --git a/bin/fm-remote-inherit-push.sh b/bin/fm-remote-inherit-push.sh index f0d6f416d4c..ed068622986 100755 --- a/bin/fm-remote-inherit-push.sh +++ b/bin/fm-remote-inherit-push.sh @@ -80,7 +80,7 @@ while IFS= read -r rel; do [ -f "$snapshot" ] && [ ! -L "$snapshot" ] || die "inherited source snapshot is unsafe: $source" bytes=$(LC_ALL=C wc -c < "$snapshot" | tr -d ' ') hash=$(sha256_file "$snapshot") || die "cannot hash inherited source: $source" - "$SCRIPT_DIR/fm-on.sh" "$ID" fm-remote-inherit.sh put "$rel" "$bytes" "$hash" "$GENERATION" < "$snapshot" + "$SCRIPT_DIR/fm-on.sh" --stdin "$ID" fm-remote-inherit.sh put "$rel" "$bytes" "$hash" "$GENERATION" < "$snapshot" else # This loop's heredoc is its control stream, not remote command input. "$SCRIPT_DIR/fm-on.sh" "$ID" fm-remote-inherit.sh absent "$rel" 0 "$EMPTY_HASH" "$GENERATION" < /dev/null diff --git a/bin/fm-remote-job-lib.sh b/bin/fm-remote-job-lib.sh index 25d7bb73b40..f6ac2ad9b99 100755 --- a/bin/fm-remote-job-lib.sh +++ b/bin/fm-remote-job-lib.sh @@ -7,23 +7,53 @@ # isolated tests), the bounded job record, worker installation, and the remote # runtime PATH. # -# A job directory is mode 0700 and contains root, home, argv (NUL-delimited), -# stdin, stdout, stderr, queue_deadline, timeout, deadline, exit, and state. -# Stage writes state=queued last. The worker atomically claims a job with -# .claim, establishes its execution deadline, changes state to running, writes -# bounded stdout/stderr and exit, then publishes state=done last. Callers wait -# for done, relay stdout and stderr separately, then reap only their completed -# record. Input, argv, stdout, and stderr are each capped at 1048576 bytes. +# A published job directory is mode 0700 and contains root, home, argv +# (NUL-delimited), stdin, seq, stdout, stderr, queue_deadline, timeout, and +# state; deadline and exit are added as execution advances, cancel is an +# optional caller-cancellation marker, and .claim may hold owner, owner_start, +# supervisor, supervisor_start, group, group_start, and armed records while +# work executes. +# Stage writes state=queued last. seq is a queue-wide monotonic staging +# sequence reserved atomically by its persistent .seq-claims directory; the +# counter is only a forward-moving allocation hint. If the bounded hint walk +# is exhausted, allocation rescans the claims for the maximum and continues +# above it. Expired claims are reaped by an independently hourly-rate-limited +# sweep. seq is the worker's FIFO ordering key within a home, with the job id +# as the deterministic tiebreak. +# FIFO is defined over completed stagings: a stage that returns before another +# begins executes first; concurrently overlapping stagings have no relative +# ordering contract. +# The worker atomically claims a job with .claim, establishes its execution +# deadline, changes state to running, writes bounded stdout/stderr and exit, +# then publishes state=done last. Callers wait for done, relay stdout and +# stderr separately, then reap only their completed record. Input, argv, +# stdout, and stderr are each capped at 1048576 bytes. # -# The worker executes one job at a time, so a deliberately long-blocking poll -# would serialize every short interactive command behind its wait window. +# The worker serves one lane per staged home: jobs for the same home run +# strictly FIFO in seq order while lanes for different homes run concurrently, +# so one home's long job never delays another home's commands. Within a lane a +# deliberately long-blocking poll would still serialize that home's short +# interactive commands behind its wait window. # fm_remote_job_command_preemptible names the read-only long-poll class # (fm-remote-delta-read.sh, the reply-log delta read). The worker preempts a -# running preemptible job as soon as a non-preemptible job is queued and -# publishes exit 76 with emptied stdout and stderr, distinct from the poll's -# exit 75 elapsed-window-with-no-data result. The delta read is non-destructive -# and cursor-anchored, so the caller's normal re-arm re-reads the same data and -# a preempted poll loses nothing. +# running preemptible job as soon as a non-preemptible job is queued for the +# same home and publishes exit 76 with emptied stdout and stderr, distinct from +# the poll's exit 75 elapsed-window-with-no-data result. The delta read is +# non-destructive and cursor-anchored, so the caller's normal re-arm re-reads +# the same data and a preempted poll loses nothing. +# +# A caller that disconnects before its job completes cancels it instead of +# abandoning it: fm_remote_job_cancel writes a cancel marker into the record, +# the worker skips a cancelled queued job and terminates a running cancelled +# job's process group, and whichever side observes terminal publication reaps +# the finalized record because no result consumer remains. fm_remote_job_wait +# honors an optional FM_REMOTE_JOB_DISCONNECT_PROBE function name. When set, +# the probe runs about once per second; a failure cancels the job and fails +# the wait. The staging entrypoint arms it with a parent-liveness probe so an +# ssh channel +# that dies without delivering a signal still cancels the abandoned job. +# Abandoned .stage.* staging litter older than +# FM_REMOTE_JOB_STAGE_REAP_SECONDS is reaped by the worker's stale sweep. # # The worker accepts only a tracked, non-symlink executable named fm-*.sh below # its configured FM_ROOT/bin. Every child receives env -i with the composed @@ -60,12 +90,16 @@ FM_REMOTE_JOB_TIMEOUT=${FM_REMOTE_JOB_TIMEOUT:-360} FM_REMOTE_JOB_WAIT_GRACE=${FM_REMOTE_JOB_WAIT_GRACE:-30} FM_REMOTE_JOB_POLL_SECONDS=${FM_REMOTE_JOB_POLL_SECONDS:-0.05} FM_REMOTE_JOB_REAP_SECONDS=${FM_REMOTE_JOB_REAP_SECONDS:-3600} +FM_REMOTE_JOB_STAGE_REAP_SECONDS=${FM_REMOTE_JOB_STAGE_REAP_SECONDS:-600} +FM_REMOTE_JOB_SEQ_CLAIM_REAP_SECONDS=86400 +FM_REMOTE_JOB_SEQ_CLAIM_REAP_INTERVAL=3600 # shellcheck disable=SC2034 # Shared protocol constant consumed by the worker and sourcing callers. FM_REMOTE_JOB_PREEMPTED_EXIT=76 FM_REMOTE_JOB_OPERATOR_PATH= FM_REMOTE_JOB_CHILD_PATH= FM_REMOTE_JOB_STATE= FM_REMOTE_JOB_JOBS= +FM_REMOTE_JOB_SEQ_CLAIMS= FM_REMOTE_JOB_ID= FM_REMOTE_JOB_STDOUT= FM_REMOTE_JOB_STDERR= @@ -96,6 +130,7 @@ fm_remote_job_validate_settings() { case "$FM_REMOTE_JOB_WAIT_GRACE" in ''|*[!0-9]*) return 1 ;; esac [ "$FM_REMOTE_JOB_WAIT_GRACE" -le 300 ] || return 1 case "$FM_REMOTE_JOB_REAP_SECONDS" in ''|*[!0-9]*|0) return 1 ;; esac + case "$FM_REMOTE_JOB_STAGE_REAP_SECONDS" in ''|*[!0-9]*|0) return 1 ;; esac return 0 } @@ -398,6 +433,10 @@ fm_remote_job_prepare_state() { # FM_REMOTE_JOB_ERROR="remote job queue is unsafe" return 1 } + FM_REMOTE_JOB_SEQ_CLAIMS=$(fm_remote_job_safe_child_dir "$FM_REMOTE_JOB_STATE" .seq-claims) || { + FM_REMOTE_JOB_ERROR="remote job sequence claims are unsafe" + return 1 + } fm_remote_job_safe_child_dir "$FM_REMOTE_JOB_STATE" logs >/dev/null || { FM_REMOTE_JOB_ERROR="remote job log directory is unsafe" return 1 @@ -423,6 +462,20 @@ fm_remote_job_regular_bounded() { # [ "$bytes" -le "$max" ] } +fm_remote_job_remove_claim_records() { # + local claim=$1 file + [ -d "$claim" ] && [ ! -L "$claim" ] || return 1 + for file in "$claim"/owner "$claim"/owner_start "$claim"/supervisor \ + "$claim"/supervisor_start "$claim"/group "$claim"/group_start "$claim"/armed \ + "$claim"/.owner.* "$claim"/.owner_start.* "$claim"/.supervisor.* \ + "$claim"/.supervisor_start.* "$claim"/.group.* "$claim"/.group_start.* \ + "$claim"/.armed.*; do + [ -e "$file" ] || [ -L "$file" ] || continue + fm_remote_job_regular_bounded "$file" 256 || return 1 + rm -f -- "$file" || return 1 + done +} + fm_remote_job_write_state() { # queued|running|done local job=$1 value=$2 tmp case "$value" in queued|running|done) ;; *) return 1 ;; esac @@ -444,9 +497,9 @@ fm_remote_job_read_state() { # case "$value" in queued|running|'done') printf '%s\n' "$value" ;; *) return 1 ;; esac } -fm_remote_job_read_number() { # queue_deadline|timeout|deadline +fm_remote_job_read_number() { # queue_deadline|timeout|deadline|seq local job=$1 field=$2 value - case "$field" in queue_deadline|timeout|deadline) ;; *) return 1 ;; esac + case "$field" in queue_deadline|timeout|deadline|seq) ;; *) return 1 ;; esac fm_remote_job_regular_bounded "$job/$field" 32 || return 1 value=$(tr -d '\n' < "$job/$field") case "$value" in ''|*[!0-9]*) return 1 ;; esac @@ -454,9 +507,9 @@ fm_remote_job_read_number() { # queue_deadline|timeout|deadline printf '%s\n' "$value" } -fm_remote_job_write_number() { # queue_deadline|timeout|deadline +fm_remote_job_write_number() { # queue_deadline|timeout|deadline|seq local job=$1 field=$2 value=$3 tmp - case "$field" in queue_deadline|timeout|deadline) ;; *) return 1 ;; esac + case "$field" in queue_deadline|timeout|deadline|seq) ;; *) return 1 ;; esac case "$value" in ''|*[!0-9]*|0) return 1 ;; esac [ -d "$job" ] && [ ! -L "$job" ] || return 1 tmp=$(umask 077; mktemp "$job/.$field.XXXXXX") || return 1 @@ -469,8 +522,96 @@ fm_remote_job_read_deadline() { # fm_remote_job_read_number "$1" deadline } +fm_remote_job_advance_seq_hint() { # + local value=$1 counter current tmp + counter="$FM_REMOTE_JOB_STATE/seq" + current=$(cat "$counter" 2>/dev/null || true) + case "$current" in ''|*[!0-9]*) current=0 ;; esac + [ "$value" -gt "$current" ] || return 0 + tmp=$(umask 077; mktemp "$FM_REMOTE_JOB_STATE/.seqhint.XXXXXX") || return 1 + printf '%s\n' "$value" > "$tmp" || { rm -f -- "$tmp"; return 1; } + chmod 600 "$tmp" || { rm -f -- "$tmp"; return 1; } + current=$(cat "$counter" 2>/dev/null || true) + case "$current" in ''|*[!0-9]*) current=0 ;; esac + if [ "$value" -gt "$current" ]; then + mv -f -- "$tmp" "$counter" || { rm -f -- "$tmp"; return 1; } + else + rm -f -- "$tmp" + fi +} + +fm_remote_job_next_seq() { # [stage-dir destination] + local stage=${1:-} destination=${2:-} counter value claim attempt=0 recovered=0 maximum entry + [ -n "$FM_REMOTE_JOB_STATE" ] && [ -n "$FM_REMOTE_JOB_SEQ_CLAIMS" ] || return 1 + counter="$FM_REMOTE_JOB_STATE/seq" + value=$(cat "$counter" 2>/dev/null || true) + case "$value" in ''|*[!0-9]*) value=0 ;; esac + while :; do + if [ "$attempt" -ge 100000 ]; then + [ "$recovered" -eq 0 ] || return 1 + maximum=0 + for entry in "$FM_REMOTE_JOB_SEQ_CLAIMS"/*; do + [ -d "$entry" ] && [ ! -L "$entry" ] || continue + entry=${entry##*/} + case "$entry" in ''|*[!0-9]*|0) continue ;; esac + [ "$entry" -le "$maximum" ] || maximum=$entry + done + value=$maximum + attempt=0 + recovered=1 + fi + attempt=$((attempt + 1)) + value=$((value + 1)) + claim="$FM_REMOTE_JOB_SEQ_CLAIMS/$value" + if (umask 077; mkdir "$claim") 2>/dev/null; then + chmod 700 "$claim" || return 1 + fm_remote_job_advance_seq_hint "$value" || true + if [ -n "$stage" ]; then + if ! fm_remote_job_write_number "$stage" seq "$value" \ + || ! fm_remote_job_write_state "$stage" queued \ + || ! mv -- "$stage" "$destination"; then + rm -f -- "$stage/state" "$stage/seq" + return 1 + fi + rm -f -- "$destination/.owner-pid" "$destination/.owner-start" || true + fi + printf '%s\n' "$value" + return 0 + fi + [ -d "$claim" ] && [ ! -L "$claim" ] || return 1 + done +} + +fm_remote_job_cancelled() { # + [ -f "$1/cancel" ] && [ ! -L "$1/cancel" ] +} + +# Mark a job cancelled on behalf of a disconnected or abandoning caller. The +# marker never rewrites state: the worker observes it, skips a cancelled queued +# job, and stops a running cancelled job's process group. The worker reaps after +# terminal publication; if publication already won the race, this function +# reaps instead. Cancelling a job that disappeared is a harmless no-op. +fm_remote_job_cancel() { # + local account_home=$1 id=$2 job state tmp + fm_remote_job_prepare_state "$account_home" || return 1 + job=$(fm_remote_job_job_dir "$id" 2>/dev/null) || return 0 + state=$(fm_remote_job_read_state "$job" 2>/dev/null || true) + if [ "$state" = 'done' ]; then + fm_remote_job_reap "$account_home" "$id" 2>/dev/null || true + return 0 + fi + tmp=$(umask 077; mktemp "$job/.cancel.XXXXXX") || return 1 + printf 'cancelled: caller disconnected or abandoned the job\n' > "$tmp" || { rm -f -- "$tmp"; return 1; } + chmod 600 "$tmp" || { rm -f -- "$tmp"; return 1; } + mv -f -- "$tmp" "$job/cancel" || return 1 + state=$(fm_remote_job_read_state "$job" 2>/dev/null || true) + if [ "$state" = 'done' ]; then + fm_remote_job_reap "$account_home" "$id" 2>/dev/null || true + fi +} + fm_remote_job_stage() { # [args...]; stdin is captured - local account_home=$1 root=$2 home=$3 command=$4 stage id destination bytes queue_deadline + local account_home=$1 root=$2 home=$3 command=$4 stage id destination bytes queue_deadline owner_start shift 4 fm_remote_job_prepare_state "$account_home" || return 1 root=$(fm_remote_job_canonical_existing_dir "$root") || { @@ -483,13 +624,20 @@ fm_remote_job_stage() { # [args...]; stdi } case "$command" in fm-*.sh) ;; *) FM_REMOTE_JOB_ERROR="remote job command is outside the fm-*.sh namespace"; return 1 ;; esac case "$command" in */*|*..*) FM_REMOTE_JOB_ERROR="remote job command contains a path or traversal"; return 1 ;; esac + owner_start=$(fm_remote_job_process_start "$$") || { + FM_REMOTE_JOB_ERROR="cannot establish remote job staging ownership" + return 1 + } stage=$(umask 077; mktemp -d "$FM_REMOTE_JOB_JOBS/.stage.XXXXXX") || { FM_REMOTE_JOB_ERROR="cannot stage remote job" return 1 } chmod 700 "$stage" || { rm -rf -- "$stage"; return 1; } queue_deadline=$(( $(date +%s) + FM_REMOTE_JOB_QUEUE_TIMEOUT )) - if ! printf '%s\n' "$root" > "$stage/root" || + if ! printf '%s\n' "$$" > "$stage/.owner-pid" || + ! printf '%s\n' "$owner_start" > "$stage/.owner-start" || + ! chmod 600 "$stage/.owner-pid" "$stage/.owner-start" || + ! printf '%s\n' "$root" > "$stage/root" || ! printf '%s\n' "$home" > "$stage/home" || ! printf '%s\n' "$queue_deadline" > "$stage/queue_deadline" || ! printf '%s\n' "$FM_REMOTE_JOB_TIMEOUT" > "$stage/timeout" || @@ -513,19 +661,23 @@ fm_remote_job_stage() { # [args...]; stdi : > "$stage/stdout" : > "$stage/stderr" chmod 600 "$stage/stdout" "$stage/stderr" || { rm -rf -- "$stage"; return 1; } - fm_remote_job_write_state "$stage" queued || { rm -rf -- "$stage"; return 1; } id="job-${stage##*/.stage.}" fm_remote_job_safe_id "$id" || { rm -rf -- "$stage"; return 1; } destination="$FM_REMOTE_JOB_JOBS/$id" [ ! -e "$destination" ] && [ ! -L "$destination" ] || { rm -rf -- "$stage"; return 1; } - mv -- "$stage" "$destination" || { rm -rf -- "$stage"; return 1; } + if ! fm_remote_job_next_seq "$stage" "$destination" >/dev/null; then + rm -rf -- "$stage" + FM_REMOTE_JOB_ERROR="cannot allocate and publish a remote job staging sequence" + return 1 + fi # shellcheck disable=SC2034 # Sourceable API consumed by callers that do not use command substitution. FM_REMOTE_JOB_ID=$id printf '%s\n' "$id" } -fm_remote_job_wait() { # +fm_remote_job_wait() { # ; honors FM_REMOTE_JOB_DISCONNECT_PROBE local account_home=$1 id=$2 job state queue_deadline execution_timeout wait_deadline exit_value + local now next_probe=0 fm_remote_job_prepare_state "$account_home" || return 1 job=$(fm_remote_job_job_dir "$id") || { FM_REMOTE_JOB_ERROR="remote job record disappeared or became unsafe" @@ -568,10 +720,19 @@ fm_remote_job_wait() { # queued|running) ;; *) FM_REMOTE_JOB_ERROR="remote job state is invalid"; return 1 ;; esac - if [ "$(date +%s)" -ge "$wait_deadline" ]; then + now=$(date +%s) + if [ "$now" -ge "$wait_deadline" ]; then FM_REMOTE_JOB_ERROR="remote job did not complete within its bounded wait" return 1 fi + if [ -n "${FM_REMOTE_JOB_DISCONNECT_PROBE:-}" ] && [ "$now" -ge "$next_probe" ]; then + next_probe=$((now + 1)) + if ! "$FM_REMOTE_JOB_DISCONNECT_PROBE"; then + fm_remote_job_cancel "$account_home" "$id" 2>/dev/null || true + FM_REMOTE_JOB_ERROR="remote job caller disconnected; the job was cancelled" + return 1 + fi + fi sleep "$FM_REMOTE_JOB_POLL_SECONDS" done } @@ -581,14 +742,14 @@ fm_remote_job_reap() { # ; only removes an exact completed re fm_remote_job_prepare_state "$account_home" || return 1 job=$(fm_remote_job_job_dir "$id") || return 1 [ "$(fm_remote_job_read_state "$job")" = 'done' ] || return 1 - for file in root home queue_deadline timeout deadline argv stdin stdout stderr exit state; do + for file in root home queue_deadline timeout deadline seq cancel argv stdin stdout stderr exit state .owner-pid .owner-start; do [ -e "$job/$file" ] || continue [ ! -L "$job/$file" ] || return 1 rm -f -- "$job/$file" || return 1 done if [ -e "$job/.claim" ] || [ -L "$job/.claim" ]; then [ -d "$job/.claim" ] && [ ! -L "$job/.claim" ] || return 1 - rm -f -- "$job/.claim/owner" "$job/.claim/supervisor" "$job/.claim/group" "$job/.claim/armed" || return 1 + fm_remote_job_remove_claim_records "$job/.claim" || return 1 rmdir "$job/.claim" || return 1 fi rmdir "$job" @@ -600,8 +761,18 @@ fm_remote_job_path_mtime() { # if [ "$(uname -s 2>/dev/null || true)" = Darwin ]; then stat -f %m "$1" 2>/dev/null; else stat -c %Y "$1" 2>/dev/null; fi } +fm_remote_job_stage_owner_alive() { # + local stage=$1 pid recorded_start actual_start + pid=$(fm_remote_job_read_single_line "$stage/.owner-pid" 64 2>/dev/null) || return 1 + case "$pid" in ''|*[!0-9]*) return 1 ;; esac + [ "$pid" -gt 1 ] || return 1 + recorded_start=$(fm_remote_job_read_single_line "$stage/.owner-start" 256 2>/dev/null) || return 1 + actual_start=$(fm_remote_job_process_start "$pid" 2>/dev/null) || return 1 + [ "$recorded_start" = "$actual_start" ] +} + fm_remote_job_reap_stale() { # - local account_home=$1 job id state mtime now + local account_home=$1 job id state mtime now stage claim value marker tmp reap_claims=0 fm_remote_job_prepare_state "$account_home" || return 1 now=$(date +%s) for job in "$FM_REMOTE_JOB_JOBS"/job-*; do @@ -615,6 +786,39 @@ fm_remote_job_reap_stale() { # [ $((now - mtime)) -ge "$FM_REMOTE_JOB_REAP_SECONDS" ] || continue fm_remote_job_reap "$account_home" "$id" || true done + marker="$FM_REMOTE_JOB_STATE/.seq-claims-reaped" + mtime=$(fm_remote_job_path_mtime "$marker" 2>/dev/null || true) + case "$mtime" in + ''|*[!0-9]*) reap_claims=1 ;; + *) [ $((now - mtime)) -lt "$FM_REMOTE_JOB_SEQ_CLAIM_REAP_INTERVAL" ] || reap_claims=1 ;; + esac + if [ "$reap_claims" -eq 1 ]; then + tmp=$(umask 077; mktemp "$FM_REMOTE_JOB_STATE/.seqreap.XXXXXX") || tmp= + if [ -n "$tmp" ] && printf '%s\n' "$now" > "$tmp" && chmod 600 "$tmp" \ + && mv -f -- "$tmp" "$marker"; then + for claim in "$FM_REMOTE_JOB_SEQ_CLAIMS"/*; do + [ -d "$claim" ] && [ ! -L "$claim" ] || continue + value=${claim##*/} + case "$value" in ''|*[!0-9]*|0) continue ;; esac + mtime=$(fm_remote_job_path_mtime "$claim" 2>/dev/null || true) + case "$mtime" in ''|*[!0-9]*) continue ;; esac + [ $((now - mtime)) -ge "$FM_REMOTE_JOB_SEQ_CLAIM_REAP_SECONDS" ] || continue + rmdir "$claim" 2>/dev/null || true + done + else + [ -z "$tmp" ] || rm -f -- "$tmp" + fi + fi + # Staging litter a killed caller left behind is reaped after its owner is no + # longer the process that created it and the stage has exceeded the age bound. + for stage in "$FM_REMOTE_JOB_JOBS"/.stage.*; do + [ -d "$stage" ] && [ ! -L "$stage" ] || continue + fm_remote_job_stage_owner_alive "$stage" && continue + mtime=$(fm_remote_job_path_mtime "$stage" 2>/dev/null || true) + case "$mtime" in ''|*[!0-9]*) continue ;; esac + [ $((now - mtime)) -ge "$FM_REMOTE_JOB_STAGE_REAP_SECONDS" ] || continue + rm -rf -- "$stage" + done } fm_remote_job_launchagent_paths() { # diff --git a/bin/fm-remote-job-worker.sh b/bin/fm-remote-job-worker.sh index 2d7528a427f..14598eb7670 100755 --- a/bin/fm-remote-job-worker.sh +++ b/bin/fm-remote-job-worker.sh @@ -15,6 +15,13 @@ # have been committed. The library header owns the exact record fields and # lifecycle. # +# The shared library header owns lane selection, FIFO, and caller-cancellation +# contracts. This serving loop implements each active lane as a tracked, +# top-level --lane process that claims one job, records itself as the claim's +# supervisor, and runs it to publication. Shutdown stops every tracked lane and +# its recorded command group, leaving interrupted records for the replacement +# worker's orphan recovery. +# # The worker is abandoned when its configured FM_ROOT stops being a genuine # Firstmate checkout - the state a pruned no-mistakes gate worktree, a returned # pooled worktree, or a removed test fixture root leaves behind. It can never @@ -50,13 +57,17 @@ FM_ROOT=${FM_ROOT_OVERRIDE:-$(CDPATH='' cd "$SCRIPT_DIR/.." && pwd -P)} # shellcheck source=bin/fm-remote-job-lib.sh . "$SCRIPT_DIR/fm-remote-job-lib.sh" -WORKER_ACTIVE_JOB= WORKER_LOCK= WORKER_LOCK_HELD=0 WORKER_RELEASE_OWNERSHIP=1 WORKER_SUPERVISED_PID= WORKER_PREEMPTIBLE=0 WORKER_PREEMPTED=0 +WORKER_LANE_HOME= +WORKER_LANE_HOMES=() +WORKER_LANE_PIDS=() +WORKER_LANE_STARTS=() +WORKER_LANE_JOBS=() worker_error() { printf 'remote-job-worker: %s\n' "$1" >&2; } @@ -135,7 +146,7 @@ worker_quarantined_execution_stopped() { # [ ! -e "$file" ] && [ ! -L "$file" ] && continue [ ! -L "$file" ] || return 1 pid=$(worker_read_process_id "$file") || return 1 - worker_process_or_group_alive "$kind" "$pid" && return 1 + worker_recorded_execution_alive "$job" "$kind" "$pid" && return 1 done done } @@ -248,6 +259,69 @@ worker_signal_process_or_group() { # process|group esac } +worker_supervisor_identity_status() { # + local job=$1 pid=$2 recorded_start actual_start + recorded_start=$(fm_remote_job_read_single_line "$job/.claim/supervisor_start" 256 2>/dev/null) || return 2 + actual_start=$(fm_remote_job_process_start "$pid" 2>/dev/null) || { + worker_process_or_group_alive process "$pid" && return 2 + return 1 + } + [ "$recorded_start" = "$actual_start" ] && return 0 + return 1 +} + +# A leaderless live group still belongs to the recorded execution: its PGID +# cannot be reused while any old member survives, so it remains safe to signal. +# A live leader whose start identity mismatches proves PID reuse and makes the +# recorded group stale; an unreadable live leader stays indeterminate so the +# stop loop retries rather than signaling or declaring the group dead. +worker_group_identity_status() { # + local job=$1 pid=$2 recorded_start actual_start file="$1/.claim/group_start" + [ -e "$file" ] || [ -L "$file" ] || return 3 + recorded_start=$(fm_remote_job_read_single_line "$file" 256 2>/dev/null) || return 2 + actual_start=$(fm_remote_job_process_start "$pid" 2>/dev/null) || { + kill -0 "$pid" 2>/dev/null && return 2 + worker_process_or_group_alive group "$pid" && return 0 + return 1 + } + [ "$recorded_start" = "$actual_start" ] && return 0 + return 1 +} + +worker_recorded_execution_alive() { # process|group + local job=$1 kind=$2 pid=$3 identity_status + if [ "$kind" = process ]; then + worker_supervisor_identity_status "$job" "$pid" + identity_status=$? + case "$identity_status" in + 0) ;; + 1) return 1 ;; + 2) worker_process_or_group_alive process "$pid"; return ;; + esac + else + worker_group_identity_status "$job" "$pid" + identity_status=$? + case "$identity_status" in + 0|3) ;; + 1) return 1 ;; + 2) worker_process_or_group_alive group "$pid"; return ;; + esac + fi + worker_process_or_group_alive "$kind" "$pid" +} + +worker_signal_recorded_execution() { # process|group + local job=$1 kind=$2 signal=$3 pid=$4 identity_status + if [ "$kind" = process ]; then + worker_supervisor_identity_status "$job" "$pid" || return 0 + else + worker_group_identity_status "$job" "$pid" + identity_status=$? + case "$identity_status" in 0|3) ;; *) return 0 ;; esac + fi + worker_signal_process_or_group "$kind" "$signal" "$pid" +} + worker_stop_recorded_execution() { # local job=$1 kind file pid attempt still_alive for kind in process group; do @@ -255,8 +329,8 @@ worker_stop_recorded_execution() { # [ ! -e "$file" ] && [ ! -L "$file" ] && continue [ ! -L "$file" ] || return 1 pid=$(worker_read_process_id "$file") || return 1 - worker_signal_process_or_group "$kind" TERM "$pid" - worker_signal_process_or_group "$kind" KILL "$pid" + worker_signal_recorded_execution "$job" "$kind" TERM "$pid" + worker_signal_recorded_execution "$job" "$kind" KILL "$pid" wait "$pid" 2>/dev/null || true done attempt=0 @@ -267,31 +341,51 @@ worker_stop_recorded_execution() { # case "$kind" in process) file="$job/.claim/supervisor" ;; group) file="$job/.claim/group" ;; esac [ -e "$file" ] || continue pid=$(worker_read_process_id "$file") || return 1 - worker_process_or_group_alive "$kind" "$pid" && still_alive=1 + if worker_recorded_execution_alive "$job" "$kind" "$pid"; then + still_alive=1 + worker_signal_recorded_execution "$job" "$kind" TERM "$pid" + worker_signal_recorded_execution "$job" "$kind" KILL "$pid" + fi done [ "$still_alive" -eq 1 ] || break sleep 0.01 done [ "$still_alive" -eq 0 ] || return 1 - rm -f -- "$job/.claim/supervisor" "$job/.claim/group" "$job/.claim/armed" + rm -f -- "$job/.claim/supervisor" "$job/.claim/supervisor_start" \ + "$job/.claim/group" "$job/.claim/group_start" "$job/.claim/armed" +} + +# Stop every tracked lane process and its recorded command execution. The lane +# is signalled first so it cannot dispatch further work, then the job's +# recorded supervisor and group are verified stopped; a job interrupted here +# stays running-with-a-dead-owner for the replacement worker's orphan recovery, +# exactly as a crashed single-process worker's job did. +worker_lane_identity_matches() { # + local pid=$1 start=$2 actual_start + [ -n "$start" ] || return 1 + actual_start=$(fm_remote_job_process_start "$pid" 2>/dev/null) || return 1 + [ "$actual_start" = "$start" ] } worker_stop_active_execution() { - local job=${WORKER_ACTIVE_JOB:-} owner owner_pid state - if [ -n "$job" ]; then - worker_stop_recorded_execution "$job" || return 1 - else - for job in "$FM_REMOTE_JOB_JOBS"/job-*; do - [ -d "$job" ] && [ ! -L "$job" ] || continue - state=$(fm_remote_job_read_state "$job" 2>/dev/null || true) - [ "$state" = running ] || continue - owner="$job/.claim/owner" - owner_pid=$(worker_read_process_id "$owner" 2>/dev/null || true) - [ "$owner_pid" = "${BASHPID:-$$}" ] || continue - worker_stop_recorded_execution "$job" || return 1 - done - fi - WORKER_ACTIVE_JOB= + local i=0 count=${#WORKER_LANE_PIDS[@]} job pid start failed=0 + while [ "$i" -lt "$count" ]; do + pid=${WORKER_LANE_PIDS[$i]} + start=${WORKER_LANE_STARTS[$i]} + job=${WORKER_LANE_JOBS[$i]} + if worker_lane_identity_matches "$pid" "$start"; then kill -TERM "$pid" 2>/dev/null || true; fi + if worker_lane_identity_matches "$pid" "$start"; then kill -KILL "$pid" 2>/dev/null || true; fi + wait "$pid" 2>/dev/null || true + if [ -d "$job" ] && [ ! -L "$job" ]; then + worker_stop_recorded_execution "$job" || failed=1 + fi + i=$((i + 1)) + done + WORKER_LANE_HOMES=() + WORKER_LANE_PIDS=() + WORKER_LANE_STARTS=() + WORKER_LANE_JOBS=() + [ "$failed" -eq 0 ] } # Ignore, rather than restore the default disposition for, the signals this @@ -333,21 +427,40 @@ worker_exit_cleanup() { } worker_claim() { # - local job=$1 claim + local job=$1 claim pid start pid_tmp start_tmp claim="$job/.claim" [ ! -e "$claim" ] && [ ! -L "$claim" ] || return 1 (umask 077; mkdir "$claim") || return 1 - printf '%s\n' "${BASHPID:-$$}" > "$claim/owner" || { rmdir "$claim" 2>/dev/null || true; return 1; } - chmod 600 "$claim/owner" || { rm -f -- "$claim/owner"; rmdir "$claim" 2>/dev/null || true; return 1; } + pid=${BASHPID:-$$} + start=$(fm_remote_job_process_start "$pid") || { rmdir "$claim" 2>/dev/null || true; return 1; } + pid_tmp=$(umask 077; mktemp "$claim/.owner.XXXXXX") || { rmdir "$claim" 2>/dev/null || true; return 1; } + start_tmp=$(umask 077; mktemp "$claim/.owner_start.XXXXXX") || { + rm -f -- "$pid_tmp" + rmdir "$claim" 2>/dev/null || true + return 1 + } + if ! printf '%s\n' "$pid" > "$pid_tmp" || ! printf '%s\n' "$start" > "$start_tmp" \ + || ! chmod 600 "$pid_tmp" "$start_tmp" || ! mv -f -- "$start_tmp" "$claim/owner_start" \ + || ! mv -f -- "$pid_tmp" "$claim/owner"; then + rm -f -- "$pid_tmp" "$start_tmp" "$claim/owner" "$claim/owner_start" + rmdir "$claim" 2>/dev/null || true + return 1 + fi } worker_claim_owner_alive() { # - local job=$1 claim="$1/.claim" owner pid + local job=$1 claim="$1/.claim" owner pid recorded_start actual_start [ -d "$claim" ] && [ ! -L "$claim" ] || return 1 owner="$claim/owner" fm_remote_job_regular_bounded "$owner" 64 || return 1 pid=$(tr -d '\n' < "$owner") case "$pid" in ''|*[!0-9]*) return 1 ;; esac + if [ -e "$claim/owner_start" ] || [ -L "$claim/owner_start" ]; then + recorded_start=$(fm_remote_job_read_single_line "$claim/owner_start" 256 2>/dev/null) || return 1 + actual_start=$(fm_remote_job_process_start "$pid" 2>/dev/null) || return 1 + [ "$recorded_start" = "$actual_start" ] + return + fi kill -0 "$pid" 2>/dev/null } @@ -357,15 +470,23 @@ worker_clear_dead_claim() { # worker_claim_owner_alive "$job" && return 1 [ -d "$claim" ] && [ ! -L "$claim" ] || return 1 [ ! -e "$claim/owner" ] || [ ! -L "$claim/owner" ] || return 1 - rm -f -- "$claim/owner" "$claim/supervisor" "$claim/group" "$claim/armed" || return 1 + fm_remote_job_remove_claim_records "$claim" || return 1 rmdir "$claim" } -worker_recover_orphaned_job() { # - local job=$1 file - worker_claim_owner_alive "$job" && return 1 +# Reclaim a running job this serving loop does not own: a record left by a +# crashed worker, whether its lane process died with it or survived it. The +# recorded execution is stopped either way - a surviving foreign lane is not +# supervised by any owner and a second lane for its home must never start +# beside it - and the record publishes unknown completion, exactly as a +# crashed single-process worker's job always has. +worker_reclaim_running_job() { # + local job=$1 file state worker_stop_recorded_execution "$job" || return 1 + state=$(fm_remote_job_read_state "$job" 2>/dev/null) || return 1 worker_clear_dead_claim "$job" || return 1 + [ "$state" = 'done' ] && return 0 + [ "$state" = running ] || return 1 for file in .stdout.pipe .stderr.pipe; do [ ! -e "$job/$file" ] && [ ! -L "$job/$file" ] || { [ ! -L "$job/$file" ] || return 1 @@ -394,7 +515,7 @@ worker_read_text() { # } worker_publish_result() { # - local job=$1 exit_status=$2 tmp + local job=$1 exit_status=$2 tmp account_home case "$exit_status" in ''|*[!0-9]*) exit_status=125 ;; esac [ "$exit_status" -le 255 ] || exit_status=125 for tmp in stdout stderr; do @@ -404,17 +525,23 @@ worker_publish_result() { # printf '%s\n' "$exit_status" > "$tmp" || { rm -f -- "$tmp"; return 1; } chmod 600 "$tmp" || { rm -f -- "$tmp"; return 1; } mv -f -- "$tmp" "$job/exit" || { rm -f -- "$tmp"; return 1; } - fm_remote_job_write_state "$job" 'done' + fm_remote_job_write_state "$job" 'done' || return 1 + if fm_remote_job_cancelled "$job"; then + account_home=$(worker_account_home 2>/dev/null || true) + if [ -n "$account_home" ]; then + fm_remote_job_reap "$account_home" "${job##*/}" 2>/dev/null || true + fi + fi } worker_run_with_timeout() { # [args...] - local job=$1 timeout=$2 group_file armed_file group_pid rc tmp deadline next_heartbeat attempt - local timed_out=0 heartbeat_failed=0 + local job=$1 timeout=$2 group_file group_start_file armed_file group_pid group_start + local group_tmp group_start_tmp rc tmp deadline next_check attempt timed_out=0 cancelled=0 WORKER_PREEMPTED=0 shift 2 group_file="$job/.claim/group" + group_start_file="$job/.claim/group_start" armed_file="$job/.claim/armed" - WORKER_ACTIVE_JOB=$job set -m ( while [ ! -f "$armed_file" ] || [ -L "$armed_file" ]; do @@ -425,43 +552,47 @@ worker_run_with_timeout() { # [args...] ) & group_pid=$! set +m - tmp=$(umask 077; mktemp "$job/.claim/.group.XXXXXX") || { + group_start=$(fm_remote_job_process_start "$group_pid") || { worker_signal_process_or_group group KILL "$group_pid" wait "$group_pid" 2>/dev/null || true - WORKER_ACTIVE_JOB= return 125 } - printf '%s\n' "$group_pid" > "$tmp" || { - rm -f -- "$tmp" + group_tmp=$(umask 077; mktemp "$job/.claim/.group.XXXXXX") || { worker_signal_process_or_group group KILL "$group_pid" wait "$group_pid" 2>/dev/null || true - WORKER_ACTIVE_JOB= return 125 } - if ! chmod 600 "$tmp" || ! mv -f -- "$tmp" "$group_file"; then - rm -f -- "$tmp" + group_start_tmp=$(umask 077; mktemp "$job/.claim/.group_start.XXXXXX") || { + rm -f -- "$group_tmp" + worker_signal_process_or_group group KILL "$group_pid" + wait "$group_pid" 2>/dev/null || true + return 125 + } + if ! printf '%s\n' "$group_pid" > "$group_tmp" \ + || ! printf '%s\n' "$group_start" > "$group_start_tmp" \ + || ! chmod 600 "$group_tmp" "$group_start_tmp" \ + || ! mv -f -- "$group_start_tmp" "$group_start_file" \ + || ! mv -f -- "$group_tmp" "$group_file"; then + rm -f -- "$group_tmp" "$group_start_tmp" "$group_file" "$group_start_file" worker_signal_process_or_group group KILL "$group_pid" wait "$group_pid" 2>/dev/null || true - WORKER_ACTIVE_JOB= return 125 fi tmp=$(umask 077; mktemp "$job/.claim/.armed.XXXXXX") || { worker_signal_process_or_group group KILL "$group_pid" wait "$group_pid" 2>/dev/null || true - rm -f -- "$group_file" - WORKER_ACTIVE_JOB= + rm -f -- "$group_file" "$group_start_file" return 125 } if ! chmod 600 "$tmp" || ! mv -f -- "$tmp" "$armed_file"; then rm -f -- "$tmp" worker_signal_process_or_group group KILL "$group_pid" wait "$group_pid" 2>/dev/null || true - rm -f -- "$group_file" - WORKER_ACTIVE_JOB= + rm -f -- "$group_file" "$group_start_file" return 125 fi deadline=$((SECONDS + timeout)) - next_heartbeat=$((SECONDS + 1)) + next_check=$((SECONDS + 1)) while worker_process_or_group_alive group "$group_pid"; do if [ "$SECONDS" -ge "$deadline" ]; then worker_signal_process_or_group group TERM "$group_pid" @@ -469,14 +600,19 @@ worker_run_with_timeout() { # [args...] timed_out=1 break fi - if [ "$SECONDS" -ge "$next_heartbeat" ]; then - if ! worker_write_heartbeat; then + if [ "$SECONDS" -ge "$next_check" ]; then + if fm_remote_job_cancelled "$job"; then worker_signal_process_or_group group TERM "$group_pid" + attempt=0 + while worker_process_or_group_alive group "$group_pid" && [ "$attempt" -lt 20 ]; do + attempt=$((attempt + 1)) + sleep 0.05 + done worker_signal_process_or_group group KILL "$group_pid" - heartbeat_failed=1 + cancelled=1 break fi - if [ "$WORKER_PREEMPTIBLE" -eq 1 ] && worker_preempting_waiter_exists; then + if [ "$WORKER_PREEMPTIBLE" -eq 1 ] && worker_preempting_waiter_exists "$WORKER_LANE_HOME"; then worker_signal_process_or_group group TERM "$group_pid" attempt=0 while worker_process_or_group_alive group "$group_pid" && [ "$attempt" -lt 20 ]; do @@ -487,16 +623,15 @@ worker_run_with_timeout() { # [args...] WORKER_PREEMPTED=1 break fi - next_heartbeat=$((SECONDS + 1)) + next_check=$((SECONDS + 1)) fi sleep "$FM_REMOTE_JOB_POLL_SECONDS" done wait "$group_pid" 2>/dev/null rc=$? - rm -f -- "$group_file" "$armed_file" - WORKER_ACTIVE_JOB= + rm -f -- "$group_file" "$group_start_file" "$armed_file" [ "$timed_out" -eq 0 ] || return 124 - [ "$heartbeat_failed" -eq 0 ] || return 125 + [ "$cancelled" -eq 0 ] || return 130 [ "$WORKER_PREEMPTED" -eq 0 ] || return "$FM_REMOTE_JOB_PREEMPTED_EXIT" return "$rc" } @@ -508,12 +643,17 @@ worker_job_command() { # ; the first argv element of a staged record printf '%s\n' "$first" } -worker_preempting_waiter_exists() { - local job state command +worker_preempting_waiter_exists() { # + local lane_home=$1 job state command job_home for job in "$FM_REMOTE_JOB_JOBS"/job-*; do [ -d "$job" ] && [ ! -L "$job" ] || continue state=$(fm_remote_job_read_state "$job" 2>/dev/null || true) [ "$state" = queued ] || continue + fm_remote_job_cancelled "$job" && continue + # Lanes are per home, so only a waiter for this lane's own home may + # preempt; another home's queue drains through its own lane. + job_home=$(worker_read_text "$job" home 8192 2>/dev/null || true) + [ "$job_home" = "$lane_home" ] || continue command=$(worker_job_command "$job" 2>/dev/null || true) fm_remote_job_command_preemptible "$command" || return 0 done @@ -544,6 +684,9 @@ worker_run_job() { # home=$(worker_read_text "$job" home 8192) || { worker_publish_result "$job" 126; return; } root=$(fm_remote_job_canonical_existing_dir "$root") || { worker_publish_result "$job" 126; return; } home=$(fm_remote_job_canonical_home "$home") || { worker_publish_result "$job" 126; return; } + # The lane key everywhere - dispatch and the preemption scan - is the staged + # home field's exact text, so this comparison value is read the same way. + WORKER_LANE_HOME=$(worker_read_text "$job" home 8192 2>/dev/null || true) [ "$root" = "$FM_ROOT" ] || { worker_publish_result "$job" 126; return; } [ -f "$root/AGENTS.md" ] && [ ! -L "$root/AGENTS.md" ] && [ -d "$root/bin" ] && [ ! -L "$root/bin" ] || { worker_publish_result "$job" 126; return; } @@ -634,8 +777,162 @@ worker_run_job() { # worker_publish_result "$job" "$rc" || worker_error "could not publish result for ${job##*/}" } +# Finalize a cancelled record nobody waits on: publish the interrupt result so +# the record is complete, then reap it because its caller is gone. +worker_finalize_cancelled() { # + local account_home=$1 job=$2 + : > "$job/stdout" 2>/dev/null || true + printf 'remote job cancelled after its caller disconnected\n' > "$job/stderr" 2>/dev/null || true + worker_publish_result "$job" 130 || return 1 + fm_remote_job_reap "$account_home" "${job##*/}" || true +} + +worker_lane_busy() { # + local home=$1 i=0 count=${#WORKER_LANE_HOMES[@]} + while [ "$i" -lt "$count" ]; do + [ "${WORKER_LANE_HOMES[$i]}" != "$home" ] || return 0 + i=$((i + 1)) + done + return 1 +} + +worker_lane_owns_job() { # + local job=$1 i=0 count=${#WORKER_LANE_JOBS[@]} + while [ "$i" -lt "$count" ]; do + [ "${WORKER_LANE_JOBS[$i]}" != "$job" ] || return 0 + i=$((i + 1)) + done + return 1 +} + +worker_reap_finished_lanes() { + local i=0 count=${#WORKER_LANE_PIDS[@]} pid start + local live_homes=() live_pids=() live_starts=() live_jobs=() + while [ "$i" -lt "$count" ]; do + pid=${WORKER_LANE_PIDS[$i]} + start=${WORKER_LANE_STARTS[$i]} + if worker_lane_identity_matches "$pid" "$start"; then + live_homes+=("${WORKER_LANE_HOMES[$i]}") + live_pids+=("$pid") + live_starts+=("$start") + live_jobs+=("${WORKER_LANE_JOBS[$i]}") + else + wait "$pid" 2>/dev/null || true + fi + i=$((i + 1)) + done + WORKER_LANE_HOMES=() + WORKER_LANE_PIDS=() + WORKER_LANE_STARTS=() + WORKER_LANE_JOBS=() + i=0 + count=${#live_pids[@]} + while [ "$i" -lt "$count" ]; do + WORKER_LANE_HOMES+=("${live_homes[$i]}") + WORKER_LANE_PIDS+=("${live_pids[$i]}") + WORKER_LANE_STARTS+=("${live_starts[$i]}") + WORKER_LANE_JOBS+=("${live_jobs[$i]}") + i=$((i + 1)) + done +} + +# One lane's whole execution of one job, run as a background lane process: +# claim, record this process as the claim supervisor, honor a cancel that +# arrived before running, establish the deadline, run to publication, and reap +# the record when its caller cancelled and can no longer reap it. +worker_lane_execute() { # + local account_home=$1 job=$2 timeout queue_deadline deadline + local supervisor_pid supervisor_start pid_tmp start_tmp + worker_claim "$job" || return 0 + supervisor_pid=${BASHPID:-$$} + supervisor_start=$(fm_remote_job_process_start "$supervisor_pid") || { + worker_publish_result "$job" 125 || true + return 0 + } + pid_tmp=$(umask 077; mktemp "$job/.claim/.supervisor.XXXXXX") || { + worker_publish_result "$job" 125 || true + return 0 + } + start_tmp=$(umask 077; mktemp "$job/.claim/.supervisor_start.XXXXXX") || { + rm -f -- "$pid_tmp" + worker_publish_result "$job" 125 || true + return 0 + } + if ! printf '%s\n' "$supervisor_pid" > "$pid_tmp" \ + || ! printf '%s\n' "$supervisor_start" > "$start_tmp" \ + || ! chmod 600 "$pid_tmp" "$start_tmp" \ + || ! mv -f -- "$start_tmp" "$job/.claim/supervisor_start" \ + || ! mv -f -- "$pid_tmp" "$job/.claim/supervisor"; then + rm -f -- "$pid_tmp" "$start_tmp" "$job/.claim/supervisor_start" + worker_publish_result "$job" 125 || true + return 0 + fi + if fm_remote_job_cancelled "$job"; then + worker_finalize_cancelled "$account_home" "$job" || true + return 0 + fi + queue_deadline=$(fm_remote_job_read_number "$job" queue_deadline 2>/dev/null || true) + case "$queue_deadline" in ''|*[!0-9]*) worker_publish_result "$job" 126 || true; return 0 ;; esac + if [ "$(date +%s)" -ge "$queue_deadline" ]; then + worker_publish_result "$job" 124 || true + return 0 + fi + timeout=$(fm_remote_job_read_number "$job" timeout 2>/dev/null || true) + case "$timeout" in ''|*[!0-9]*) worker_publish_result "$job" 126 || true; return 0 ;; esac + if [ "$timeout" -gt 3600 ]; then + worker_publish_result "$job" 126 || true + return 0 + fi + # The deadline is measured in whole seconds from a truncated clock read, so + # the +1 keeps the granted window at least the recorded timeout instead of + # silently shaving up to a second off it. + deadline=$(( $(date +%s) + timeout + 1 )) + fm_remote_job_write_number "$job" deadline "$deadline" || { + worker_publish_result "$job" 125 || true + return 0 + } + fm_remote_job_write_state "$job" running || { + worker_publish_result "$job" 125 || true + return 0 + } + worker_run_job "$account_home" "$job" + if fm_remote_job_cancelled "$job"; then + fm_remote_job_reap "$account_home" "${job##*/}" || true + fi +} + +# Each lane runs as its own top-level worker process (--lane), not a +# backgrounded subshell: a bash subshell does not reliably reap its dead +# children, and a zombie group leader keeps its process group signalable, so a +# subshell-hosted monitor loop can believe a finished command is still running +# until the job deadline. A top-level shell is the context the monitor loop +# has always run in. +worker_start_lane() { # + local job=$1 home=$2 lane_pid lane_start + "$SCRIPT_DIR/fm-remote-job-worker.sh" --lane "${job##*/}" & + lane_pid=$! + lane_start=$(fm_remote_job_process_start "$lane_pid" 2>/dev/null || true) + WORKER_LANE_HOMES+=("$home") + WORKER_LANE_PIDS+=("$lane_pid") + WORKER_LANE_STARTS+=("$lane_start") + WORKER_LANE_JOBS+=("$job") +} + +worker_lane_main() { # + local account_home job + fm_remote_job_safe_id "$1" || { worker_error "invalid lane job id"; exit 2; } + account_home=$(worker_account_home) || { worker_error "cannot resolve account home"; exit 1; } + FM_ROOT=$(fm_remote_job_canonical_existing_dir "$FM_ROOT") || { worker_error "configured FM_ROOT is unsafe"; exit 1; } + fm_remote_job_prepare_state "$account_home" || { worker_error "$FM_REMOTE_JOB_ERROR"; exit 1; } + job=$(fm_remote_job_job_dir "$1" 2>/dev/null) || exit 0 + worker_lane_execute "$account_home" "$job" +} + worker_process_once() { # - local account_home=$1 job id state queue_deadline timeout deadline + local account_home=$1 job id state queue_deadline home seq candidates='' + local reserved_index reserved_count home_reserved + local reserved_homes=() + worker_reap_finished_lanes for job in "$FM_REMOTE_JOB_JOBS"/job-*; do [ -d "$job" ] && [ ! -L "$job" ] || continue id=${job##*/} @@ -645,38 +942,59 @@ worker_process_once() { # state=$(fm_remote_job_read_state "$job" 2>/dev/null || true) case "$state" in queued) - worker_clear_dead_claim "$job" || continue + worker_lane_owns_job "$job" && continue + if ! worker_clear_dead_claim "$job"; then + if worker_claim_owner_alive "$job"; then + home=$(worker_read_text "$job" home 8192 2>/dev/null || true) + [ -n "$home" ] && reserved_homes+=("$home") + fi + continue + fi + if fm_remote_job_cancelled "$job"; then + worker_finalize_cancelled "$account_home" "$job" || true + continue + fi queue_deadline=$(fm_remote_job_read_number "$job" queue_deadline 2>/dev/null || true) case "$queue_deadline" in ''|*[!0-9]*) worker_publish_result "$job" 126 || true; continue ;; esac if [ "$(date +%s)" -ge "$queue_deadline" ]; then worker_publish_result "$job" 124 || true continue fi + home=$(worker_read_text "$job" home 8192 2>/dev/null || true) + [ -n "$home" ] || { worker_publish_result "$job" 126 || true; continue; } + # A record staged by an older library has no seq; order it ahead of + # sequenced work as the older job it is. + seq=$(fm_remote_job_read_number "$job" seq 2>/dev/null || true) + case "$seq" in ''|*[!0-9]*) seq=0 ;; esac + candidates="$candidates$seq"$'\t'"$id"$'\t'"$home"$'\n' ;; running) - worker_recover_orphaned_job "$job" || true + worker_lane_owns_job "$job" || worker_reclaim_running_job "$job" || true continue ;; *) continue ;; esac - worker_claim "$job" || continue - timeout=$(fm_remote_job_read_number "$job" timeout 2>/dev/null || true) - case "$timeout" in ''|*[!0-9]*) worker_publish_result "$job" 126 || true; continue ;; esac - if [ "$timeout" -gt 3600 ]; then - worker_publish_result "$job" 126 || true - continue - fi - deadline=$(( $(date +%s) + timeout )) - fm_remote_job_write_number "$job" deadline "$deadline" || { - worker_publish_result "$job" 125 || true - continue - } - fm_remote_job_write_state "$job" running || { - worker_publish_result "$job" 125 || true - continue - } - worker_run_job "$account_home" "$job" done + [ -n "$candidates" ] || return 0 + while IFS=$'\t' read -r seq id home; do + [ -n "$id" ] || continue + worker_lane_busy "$home" && continue + home_reserved=0 + reserved_index=0 + reserved_count=${#reserved_homes[@]} + while [ "$reserved_index" -lt "$reserved_count" ]; do + if [ "${reserved_homes[$reserved_index]}" = "$home" ]; then + home_reserved=1 + break + fi + reserved_index=$((reserved_index + 1)) + done + [ "$home_reserved" -eq 0 ] || continue + job=$(fm_remote_job_job_dir "$id" 2>/dev/null || true) + [ -n "$job" ] || continue + [ "$(fm_remote_job_read_state "$job" 2>/dev/null || true)" = queued ] || continue + worker_start_lane "$job" "$home" + done < <(printf '%s' "$candidates" | sort -t $'\t' -k1,1n -k2,2) } main() { @@ -795,6 +1113,10 @@ case "${1:-}" in [ "$#" -eq 1 ] || { worker_error "unexpected worker arguments"; exit 2; } main ;; + --lane) + [ "$#" -eq 2 ] || { worker_error "unexpected worker arguments"; exit 2; } + worker_lane_main "$2" + ;; '') if [ "$(fm_remote_job_platform)" = linux ]; then worker_supervise_linux; else main; fi ;; diff --git a/bin/fm-send.sh b/bin/fm-send.sh index b10f381ffd6..daa638fe7a7 100755 --- a/bin/fm-send.sh +++ b/bin/fm-send.sh @@ -131,7 +131,11 @@ # FM_PENDING_REPLY_EXISTING_CORR= resend command that preserves the body # and makes a later remote enqueue deduplicate onto that same record. An # unconfirmed fire-and-forget request exits 3 and names the same delivery id to -# retry. The remote host runs no re-ring ladder of its own: a swallowed ordinary +# retry. Every remote transport attempt is bounded by FM_SEND_REMOTE_BUDGET +# seconds (default 30, and any override must be a positive integer): a bound +# hit is completion-unknown and exits through this same unconfirmed contract +# instead of waiting out a busy remote queue. +# The remote host runs no re-ring ladder of its own: a swallowed ordinary # doorbell surfaces through the parent's pending-reply recovery and escalation, # whose recovery request re-rings the remote doorbell when it is enqueued; # fire-and-forget delivery deliberately arms neither mechanism. Internal @@ -228,6 +232,8 @@ fi . "$SCRIPT_DIR/fm-wake-lib.sh" # shellcheck source=bin/fm-task-inbox-lib.sh . "$SCRIPT_DIR/fm-task-inbox-lib.sh" +# shellcheck source=bin/fm-timeout-lib.sh +. "$SCRIPT_DIR/fm-timeout-lib.sh" FM_GUARD_CONTINUE_LINE='This is a supervision warning only; the requested message WILL still be sent.' "$SCRIPT_DIR/fm-guard.sh" || true @@ -648,7 +654,15 @@ if [ "${1:-}" = "--key" ]; then key=$2 semantic_key=$(fm_send_normalize_key "$key") if [ "$TARGET_BACKEND" = remote ]; then - if ! "$SCRIPT_DIR/fm-on.sh" "$TARGET_REMOTE_ID" fm-remote-secondmate-control.sh key "$TARGET_REMOTE_ID" "$key" < /dev/null; then + FM_SEND_REMOTE_BUDGET=${FM_SEND_REMOTE_BUDGET:-30} + case "$FM_SEND_REMOTE_BUDGET" in + ''|*[!0-9]*|0) + echo "error: FM_SEND_REMOTE_BUDGET must be a positive integer: $FM_SEND_REMOTE_BUDGET" >&2 + exit 1 + ;; + esac + if ! fm_run_timed "$FM_SEND_REMOTE_BUDGET" "$SCRIPT_DIR/fm-on.sh" "$TARGET_REMOTE_ID" \ + fm-remote-secondmate-control.sh key "$TARGET_REMOTE_ID" "$key" < /dev/null; then echo "error: key '$key' not sent to remote secondmate $TARGET_REMOTE_ID; completion may be unknown" >&2 exit 1 fi @@ -660,6 +674,15 @@ if [ "${1:-}" = "--key" ]; then fm_send_record_interrupt "$semantic_key" || exit 1 else MESSAGE=$* + if [ "$TARGET_BACKEND" = remote ]; then + FM_SEND_REMOTE_BUDGET=${FM_SEND_REMOTE_BUDGET:-30} + case "$FM_SEND_REMOTE_BUDGET" in + ''|*[!0-9]*|0) + echo "error: FM_SEND_REMOTE_BUDGET must be a positive integer: $FM_SEND_REMOTE_BUDGET" >&2 + exit 1 + ;; + esac + fi # The pre-marker answer text, kept for the closing resolved note so the # durable ledger records the plain answer without marker or corr bytes. RESOLVE_ANSWER_TEXT=$MESSAGE @@ -747,6 +770,10 @@ else # 255 is safe by that idempotence; a still-lost transport preserves a # reply-bearing request's expectation, while fire-and-forget reports the # delivery id that must be reused, because the record may have landed. + # Every transport attempt is bounded by FM_SEND_REMOTE_BUDGET seconds + # (default 30, overridable) so a busy remote queue cannot hold this send + # open indefinitely; a bound hit exits through the same + # unconfirmed-delivery contract. REMOTE_META_LOCK=$(fm_meta_lock_path "$TARGET_META") || exit 1 if ! fm_task_inbox_lock_acquire "$REMOTE_META_LOCK"; then if [ "$PENDING_REPLY_CREATED" = 1 ] && [ -n "$PENDING_REPLY_CORR" ]; then @@ -781,13 +808,22 @@ else remote_completion_unknown=0 REMOTE_SEND_ARGS=("$TARGET_REMOTE_ID" "$MESSAGE") [ -z "$FIRE_AND_FORGET_ID" ] || REMOTE_SEND_ARGS+=(fire-and-forget) - "$SCRIPT_DIR/fm-on.sh" "$TARGET_REMOTE_ID" fm-remote-secondmate-control.sh send \ - "${REMOTE_SEND_ARGS[@]}" < /dev/null || remote_rc=$? - if [ "$remote_rc" -eq 255 ]; then + # Each transport attempt is bounded by FM_SEND_REMOTE_BUDGET seconds. + # fm_run_timed's 124 means the attempt was killed at the bound with remote + # completion unknown - the enqueue may have landed - so it exits through + # the same unconfirmed-delivery contract as a lost transport, without a + # retry that would only wait out the same busy remote queue again. (A + # remote job's own timeout also relays as 124; treating it as unconfirmed + # stays safe because the remote enqueue deduplicates.) + fm_run_timed "$FM_SEND_REMOTE_BUDGET" "$SCRIPT_DIR/fm-on.sh" "$TARGET_REMOTE_ID" \ + fm-remote-secondmate-control.sh send "${REMOTE_SEND_ARGS[@]}" < /dev/null || remote_rc=$? + if [ "$remote_rc" -eq 124 ]; then + remote_completion_unknown=1 + elif [ "$remote_rc" -eq 255 ]; then remote_completion_unknown=1 remote_rc=0 - "$SCRIPT_DIR/fm-on.sh" "$TARGET_REMOTE_ID" fm-remote-secondmate-control.sh send \ - "${REMOTE_SEND_ARGS[@]}" < /dev/null || remote_rc=$? + fm_run_timed "$FM_SEND_REMOTE_BUDGET" "$SCRIPT_DIR/fm-on.sh" "$TARGET_REMOTE_ID" \ + fm-remote-secondmate-control.sh send "${REMOTE_SEND_ARGS[@]}" < /dev/null || remote_rc=$? fi fm_lock_release "$REMOTE_META_LOCK" if [ "$remote_rc" -ne 0 ] && [ "$remote_completion_unknown" -eq 1 ]; then @@ -800,6 +836,8 @@ else fi if [ "$remote_rc" -eq 255 ]; then echo "error: steer to remote secondmate $TARGET_REMOTE_ID is unconfirmed (transport lost twice; remote completion unknown). Only the correlation-reusing resend below is idempotent and lands on the same remote inbox record:" >&2 + elif [ "$remote_rc" -eq 124 ]; then + echo "error: steer to remote secondmate $TARGET_REMOTE_ID is unconfirmed (the remote transport did not complete within its ${FM_SEND_REMOTE_BUDGET}s budget; remote completion unknown). Only the correlation-reusing resend below is idempotent and lands on the same remote inbox record:" >&2 else echo "error: steer to remote secondmate $TARGET_REMOTE_ID is unconfirmed (the first transport attempt had unknown completion and the retry failed). Only the correlation-reusing resend below is idempotent and lands on the same remote inbox record:" >&2 fi diff --git a/bin/fm-test-run.sh b/bin/fm-test-run.sh index 9a1d4a8be1c..c89b6e60f83 100755 --- a/bin/fm-test-run.sh +++ b/bin/fm-test-run.sh @@ -172,6 +172,7 @@ family_for_basename() { ;; fm-backlog-handoff.test.sh|fm-on.test.sh|fm-remote-backlog-handoff.test.sh|\ fm-remote-doctor.test.sh|fm-remote-job.test.sh|fm-remote-job-orphan-reap.test.sh|\ + fm-remote-transport-lanes.test.sh|\ fm-remote-reply.test.sh|fm-remote-secondmate-lifecycle-e2e.test.sh|\ fm-remote-secondmate-trace-context.test.sh|\ fm-secondmate-harness.test.sh|fm-secondmate-lifecycle-e2e.test.sh|\ diff --git a/docs/remote-secondmates.md b/docs/remote-secondmates.md index 3099854056d..5a36fef02a3 100644 --- a/docs/remote-secondmates.md +++ b/docs/remote-secondmates.md @@ -32,8 +32,10 @@ The entrypoint authorizes that bootstrap with normal git tracking when git resol After setup, every other command verifies Firstmate's account-owned remote job worker, stages the encoded argv and stdin bytes, waits for its result, and relays stdout, stderr, and the exit status separately. On macOS the worker is `dev.firstmate.remote-job`, an Aqua-scoped LaunchAgent at `~/Library/LaunchAgents/dev.firstmate.remote-job.plist` with logs under `~/Library/Logs/`. After that bootstrap every non-doctor `fm-on.sh` target runs through that worker in the remote account's GUI session, never in the SSH process or a Herdr pane. -The worker runs one staged job at a time and preempts a running reply long-poll as soon as any command other than another reply long-poll is queued, so interactive commands and startup checks are never serialized behind a poll window. +The worker serves one lane per staged home: jobs for the same home follow the staging-order contract owned by [`bin/fm-remote-job-lib.sh`](../bin/fm-remote-job-lib.sh), while different homes' lanes run concurrently so one home's long job never delays another home's commands. +Within a home's lane the worker preempts a running reply long-poll as soon as any command other than another reply long-poll is queued for that home, so interactive commands and startup checks are never serialized behind a poll window. `bin/fm-remote-job-lib.sh` owns that preemption contract and distinguishes preemption from a wait window that closes with no data, so only a genuinely quiet window proves channel freshness while either outcome can re-arm without losing data. +A caller that disconnects or whose caller-side wait expires before its job completes cancels it instead of abandoning it: cancelled queued work is skipped, cancelled running work is stopped, and the finalized record is cleaned up, so retries never convoy behind abandoned work. Linux uses the same queue and worker protocol without the Aqua-session requirement. A worker stops itself once its configured code root stops being a Firstmate checkout, so a worker started from a worktree cannot outlive that worktree, and `bin/fm-remote-job-reap-orphans.sh` clears any worker already left behind that way without ever touching one whose checkout still exists. The remote account must provide the required toolchain, the selected worker runtime, the selected session backend, and credentials that work on that host. @@ -173,8 +175,9 @@ FM_HOME= bin/fm-send.sh fm- '' The [`fm-send.sh` header](../bin/fm-send.sh) owns the exact delivery-status contract. A routed request is delivered as a durable record in the remote home's steering inbox plus a best-effort doorbell, never by typing the payload into the pane; exit 0 means the record durably exists. -An unconfirmed transport (SSH exit 255) is retried identically once and preserves this ordinary reply-bearing request's pending-reply expectation for the record that may have landed. -If it remains unconfirmed, only the exact `FM_PENDING_REPLY_EXISTING_CORR=` resend command printed by `fm-send` is safe to run later because it preserves the request body and lets the remote enqueue deduplicate onto the same record; a plain rerun mints a different correlation and is not idempotent. +Every remote transport attempt is bounded by `FM_SEND_REMOTE_BUDGET`; that header owns the setting's default and validation contract. +An unconfirmed SSH transport (exit 255) is retried identically once, while a budget expiry is not retried because completion is unknown; either outcome preserves this ordinary reply-bearing request's pending-reply expectation for the record that may have landed. +If delivery remains unconfirmed, only the exact `FM_PENDING_REPLY_EXISTING_CORR=` resend command printed by `fm-send` is safe to run later because it preserves the request body and lets the remote enqueue deduplicate onto the same record; a plain rerun mints a different correlation and is not idempotent. When deduplication finds that the worker already moved the matching record into `handled/`, the resend exits successfully without ringing the doorbell again. The remote host runs no doorbell re-ring ladder of its own; a swallowed doorbell for an ordinary reply-bearing request surfaces through the parent's pending-reply recovery and escalation, whose recovery request rings the doorbell again when it is enqueued. `fm-peek.sh` and `fm-crew-state.sh` route remote-secondmate reads to the endpoint's host instead of consulting local worktree or backend state. @@ -250,6 +253,7 @@ bin/fm-test-run.sh tests/fm-secondmate-reconcile.test.sh bin/fm-test-run.sh tests/fm-peek-remote.test.sh bin/fm-test-run.sh tests/fm-crew-state.test.sh bin/fm-test-run.sh tests/fm-remote-job.test.sh +bin/fm-test-run.sh tests/fm-remote-transport-lanes.test.sh bin/fm-test-run.sh tests/fm-remote-doctor.test.sh bin/fm-test-run.sh tests/fm-project-origin.test.sh bin/fm-test-run.sh tests/fm-remote-reply.test.sh diff --git a/tests/fm-on.test.sh b/tests/fm-on.test.sh index 790a56d5038..cde6cb3ef49 100755 --- a/tests/fm-on.test.sh +++ b/tests/fm-on.test.sh @@ -119,7 +119,8 @@ fm_on() { # The pre-feature user path had no executable transport at all. The regression # exercises the adopted public surface end to end through a deterministic SSH -# process boundary rather than checking script source. +# process boundary rather than checking script source. A payload caller passes +# --stdin explicitly; without it the remote command's stdin is /dev/null. ARGV_ACTUAL="$REMOTE_HOME/argv.bin" ARGV_EXPECTED="$TMP_ROOT/argv-expected.bin" # shellcheck disable=SC2016 # Literal shell-looking argv is the injection probe. @@ -127,7 +128,7 @@ printf '%s\0' 'plain' 'two words' '$(touch /tmp/fm-on-injected)' '' $'line one\n printf 'payload one\npayload two\n' > "$TMP_ROOT/stdin" set +e # shellcheck disable=SC2016 # Literal shell-looking argv is the injection probe. -fm_on ios fm-probe-one.sh "$ARGV_ACTUAL" 23 \ +fm_on --stdin ios fm-probe-one.sh "$ARGV_ACTUAL" 23 \ 'plain' 'two words' '$(touch /tmp/fm-on-injected)' '' $'line one\nline two' \ < "$TMP_ROOT/stdin" > "$TMP_ROOT/stdout" 2> "$TMP_ROOT/stderr" rc=$? @@ -139,7 +140,21 @@ assert_grep 'stdin: payload one' "$TMP_ROOT/stdout" "remote stdin was not preser assert_grep 'stdin: payload two' "$TMP_ROOT/stdout" "remote stdin lost its second line" assert_grep 'stderr: separate' "$TMP_ROOT/stderr" "remote stderr was not preserved separately" assert_absent /tmp/fm-on-injected "shell-looking argv was interpreted" -pass "fm-on preserves argv, stdin, stdout, stderr, and exit status without shell interpretation" +pass "fm-on --stdin preserves argv, stdin, stdout, stderr, and exit status without shell interpretation" + +# Without --stdin the remote command must see EOF even when the caller's own +# stdin holds bytes: staging captures stdin to EOF, so an open caller stream +# must never reach it by default. +set +e +fm_on ios fm-probe-one.sh "$REMOTE_HOME/argv-default.bin" 0 'default-closed' \ + < "$TMP_ROOT/stdin" > "$TMP_ROOT/stdout-default" 2> "$TMP_ROOT/stderr-default" +rc=$? +set -e +[ "$rc" -eq 0 ] || fail "the default-closed invocation did not preserve exit status (got $rc)" +if grep -q 'stdin:' "$TMP_ROOT/stdout-default"; then + fail "caller stdin crossed the transport without --stdin: $(cat "$TMP_ROOT/stdout-default")" +fi +pass "fm-on defaults the remote command's stdin to /dev/null" # A vanished remote peer must become a bounded ssh failure instead of an # indefinite hang on a half-open TCP connection, so the existing no-result -> diff --git a/tests/fm-public-followup.test.sh b/tests/fm-public-followup.test.sh index 69a054abee8..a4f8d7bcad3 100755 --- a/tests/fm-public-followup.test.sh +++ b/tests/fm-public-followup.test.sh @@ -23,6 +23,7 @@ TEARDOWN="$ROOT/bin/fm-teardown.sh" PROMOTE="$ROOT/bin/fm-promote.sh" SESSION_START="$ROOT/bin/fm-session-start.sh" TMP_ROOT=$(fm_test_tmproot fm-public-followup) +PF_TEST_NOW=1787539200 command -v jq >/dev/null 2>&1 || { echo "skip: jq not found"; exit 0; } command -v tasks-axi >/dev/null 2>&1 || { echo "skip: tasks-axi not found"; exit 0; } @@ -97,7 +98,8 @@ run_pf() { # shift PATH="$home/fakebin:$PATH" FM_ROOT_OVERRIDE="$ROOT" FM_HOME="$home" \ FM_STATE_OVERRIDE="$home/state" FAKE_CURL_LOG="${FAKE_CURL_LOG:-}" \ - FAKE_FOLLOWUP_CODE="${FAKE_FOLLOWUP_CODE:-200}" "$PF" "$@" + FAKE_FOLLOWUP_CODE="${FAKE_FOLLOWUP_CODE:-200}" \ + FMX_NOW_OVERRIDE="${FMX_NOW_OVERRIDE:-$PF_TEST_NOW}" "$PF" "$@" } tasks_in() { # @@ -142,7 +144,7 @@ seed_commitment() { > "$home/state/x-inbox/$request.json" chmod 700 "$home/state/x-inbox" chmod 600 "$home/state/x-inbox/$request.json" - FM_HOME="$home" bash -c \ + FM_HOME="$home" FMX_NOW_OVERRIDE="$PF_TEST_NOW" bash -c \ ". '$ROOT/bin/fm-x-lib.sh'; fmx_context_registry_set '$home/state' '$request' '$platform' 1900" \ || fail "could not retain the private request context" @@ -172,7 +174,7 @@ seed_repro_commitment() { # /dev/null || fail "add failed" tasks_in "$home" public-followup bind-work "$obligation" --relation-file "$home/relation.json" >/dev/null \ || fail "bind-work failed" - FM_HOME="$home" bash -c \ + FM_HOME="$home" FMX_NOW_OVERRIDE="$PF_TEST_NOW" bash -c \ ". '$ROOT/bin/fm-x-lib.sh'; fmx_context_registry_set '$home/state' '$request' discord 2000" \ || fail "context retain failed" run_pf "$home" register "$obligation" --relation rel-code --work-home "$work_home" \ diff --git a/tests/fm-remote-job.test.sh b/tests/fm-remote-job.test.sh index a97b2edb3e8..fb6ea8ef99b 100755 --- a/tests/fm-remote-job.test.sh +++ b/tests/fm-remote-job.test.sh @@ -627,8 +627,11 @@ RECOVERY_REFUSED_RC=$? set -e [ "$RECOVERY_REFUSED_RC" -ne 0 ] || fail "quarantine recovery ignored a recorded live process" assert_present "$RECOVERY_STATE/worker.lock/quarantine" "a live recorded process lost quarantine protection" -kill "$QUARANTINED_PROCESS_PID" 2>/dev/null || true -wait "$QUARANTINED_PROCESS_PID" 2>/dev/null || true +printf '%s\n' "$QUARANTINED_PROCESS_PID" > "$RECOVERY_JOB/.claim/owner" +printf 'stale owner identity\n' > "$RECOVERY_JOB/.claim/owner_start" +printf 'stale supervisor identity\n' > "$RECOVERY_JOB/.claim/supervisor_start" +chmod 600 "$RECOVERY_JOB/.claim/owner" "$RECOVERY_JOB/.claim/owner_start" \ + "$RECOVERY_JOB/.claim/supervisor_start" HOME="$RECOVERY_HOME" FM_ROOT_OVERRIDE="$REMOTE_ROOT" FM_REMOTE_JOB_STATE_ROOT="$RECOVERY_STATE" \ FM_REMOTE_JOB_PLATFORM_OVERRIDE=Linux "$REMOTE_ROOT/bin/fm-remote-job-worker.sh" \ > "$TMP_ROOT/recovery-worker.out" 2> "$TMP_ROOT/recovery-worker.err" & @@ -637,12 +640,16 @@ for _ in $(seq 1 300); do [ -f "$RECOVERY_STATE/worker.ready" ] && break sleep 0.05 done -assert_present "$RECOVERY_STATE/worker.ready" "a stopped quarantined execution did not permit worker recovery" +assert_present "$RECOVERY_STATE/worker.ready" "a reused supervisor pid did not permit worker recovery" assert_absent "$RECOVERY_STATE/worker.lock/quarantine" "recovered worker retained stale quarantine" +kill -0 "$QUARANTINED_PROCESS_PID" 2>/dev/null \ + || fail "worker recovery signalled a process whose supervisor identity did not match" kill -TERM "$RECOVERY_WORKER_PID" wait "$RECOVERY_WORKER_PID" 2>/dev/null || true RECOVERY_WORKER_PID= -pass "quarantine clears only after recorded execution has stopped" +kill "$QUARANTINED_PROCESS_PID" 2>/dev/null || true +wait "$QUARANTINED_PROCESS_PID" 2>/dev/null || true +pass "quarantine recovery refuses unverifiable supervisors and ignores reused pids" # A replacement stops a Linux worker by signalling its whole isolated group, and # the supervisor in that group forwards a second stop signal to the same serving diff --git a/tests/fm-remote-transport-lanes.test.sh b/tests/fm-remote-transport-lanes.test.sh new file mode 100755 index 00000000000..e4815a7f7ea --- /dev/null +++ b/tests/fm-remote-transport-lanes.test.sh @@ -0,0 +1,425 @@ +#!/usr/bin/env bash +# Behavior tests for the remote transport's per-home lanes, caller-disconnect +# cancellation, stdin default, and staging-litter reaping. +# +# Pins, against the real worker and the real fm-on -> entrypoint transport +# (through the deterministic FM_SSH_BIN seam tests/fm-on.test.sh proves +# preserves exit status): +# T9: a job for home B completes while home A runs a long job, and two +# A-jobs execute strictly in stage order even when staged rapidly. +# T3: a caller killed mid-wait cancels its job - the worker never executes a +# cancelled queued job and terminates a running cancelled job's process +# group - and a caller whose parent dies without delivering a signal +# (the dead-ssh-channel shape) cancels the same way; afterwards a burst +# of short commands completes with no convoy. +# T6: a non-payload fm-on call with an OPEN stdin pipe completes instead of +# wedging staging, and a payload caller with --stdin still delivers its +# bytes through the worker. +# Stage litter older than the reap age does not survive a worker pass while +# fresh staging does. +set -u + +# shellcheck source=tests/lib.sh +. "$(dirname "${BASH_SOURCE[0]}")/lib.sh" + +ROOT=$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd -P) +# shellcheck source=bin/fm-timeout-lib.sh +. "$ROOT/bin/fm-timeout-lib.sh" + +TMP_ROOT=$(fm_test_tmproot fm-remote-transport-lanes) +mkdir -p "$TMP_ROOT" +TMP_ROOT=$(cd "$TMP_ROOT" && pwd -P) +REMOTE_ROOT="$TMP_ROOT/remote-root" +HOME_A="$TMP_ROOT/home-a" +HOME_B="$TMP_ROOT/home-b" +HOME_EDGE="$TMP_ROOT/home-a " +LOCAL_HOME="$TMP_ROOT/local-home" +ACCOUNT_HOME="$TMP_ROOT/account" +STATE_ROOT="$TMP_ROOT/remote-jobs" +FAKEBIN=$(fm_fakebin "$TMP_ROOT/fakebin") +mkdir -p "$REMOTE_ROOT/bin" "$HOME_A" "$HOME_B" "$HOME_EDGE" "$LOCAL_HOME/data" "$ACCOUNT_HOME" + +cleanup_lane_fixture() { + if [ -f "$STATE_ROOT/worker.pid" ]; then + fm_remote_job_stop_worker_tree "$(cat "$STATE_ROOT/worker.pid")" || true + fi + rm -rf -- "$TMP_ROOT" +} +trap cleanup_lane_fixture EXIT + +cp "$ROOT/bin/fm-remote-job-lib.sh" "$ROOT/bin/fm-remote-job-worker.sh" \ + "$ROOT/bin/fm-remote-entrypoint.sh" "$ROOT/bin/fm-remote-delta-read.sh" \ + "$ROOT/bin/fm-remote-secondmate-control.sh" "$ROOT/bin/fm-backend.sh" \ + "$ROOT/bin/fm-pending-reply-lib.sh" "$ROOT/bin/fm-task-inbox-lib.sh" \ + "$ROOT/bin/fm-wake-lib.sh" "$ROOT/bin/fm-marker-lib.sh" \ + "$ROOT/bin/fm-operational-input.sh" "$ROOT/bin/fm-tmux-lib.sh" \ + "$ROOT/bin/fm-composer-lib.sh" "$ROOT/bin/fm-cursor-lib.sh" \ + "$ROOT/bin/fm-classify-lib.sh" "$ROOT/bin/fm-timeout-lib.sh" \ + "$REMOTE_ROOT/bin/" +mkdir -p "$REMOTE_ROOT/bin/backends" +cp "$ROOT/bin/backends/herdr.sh" "$REMOTE_ROOT/bin/backends/herdr.sh" +printf 'fixture\n' > "$REMOTE_ROOT/AGENTS.md" +# Appends its tag to a shared log, then optionally sleeps: the log order is the +# observable execution order. +cat > "$REMOTE_ROOT/bin/fm-mark-job.sh" <<'SH' +#!/bin/bash +printf '%s\n' "$1" >> "$2" +sleep "${3:-0}" +SH +cat > "$REMOTE_ROOT/bin/fm-touch-job.sh" <<'SH' +#!/bin/bash +printf 'ran\n' > "$1" +SH +# Marks its start, sleeps, then marks completion: cancellation must leave the +# start marker without the completion marker. +cat > "$REMOTE_ROOT/bin/fm-two-phase-job.sh" <<'SH' +#!/bin/bash +printf 'started\n' > "$1" +sleep "$3" +printf 'finished\n' > "$2" +SH +cat > "$REMOTE_ROOT/bin/fm-stdin-probe.sh" <<'SH' +#!/bin/bash +while IFS= read -r line || [ -n "$line" ]; do printf 'stdin=%s\n' "$line"; done +SH +chmod +x "$REMOTE_ROOT/bin"/*.sh +git -C "$REMOTE_ROOT" init -q -b main +git -C "$REMOTE_ROOT" config user.email test@example.com +git -C "$REMOTE_ROOT" config user.name Test +git -C "$REMOTE_ROOT" add AGENTS.md bin +git -C "$REMOTE_ROOT" commit -qm 'lane transport fixture' + +# ios routes to home A, build routes to home B. +cat > "$LOCAL_HOME/data/secondmates.md" < "$FAKEBIN/fake-ssh" <<'SH' +#!/usr/bin/env bash +while [ "$#" -gt 0 ]; do + case "$1" in + -o) shift 2 ;; + --) shift; break ;; + *) exit 90 ;; + esac +done +shift 2 +exec "$FM_FAKE_REMOTE_ENTRYPOINT" "$@" +SH +chmod +x "$FAKEBIN/fake-ssh" + +export FM_REMOTE_JOB_STATE_ROOT="$STATE_ROOT" +export FM_REMOTE_JOB_PLATFORM_OVERRIDE=Linux +export FM_REMOTE_JOB_QUEUE_TIMEOUT=60 +export FM_REMOTE_JOB_TIMEOUT=30 +export FM_REMOTE_JOB_STAGE_REAP_SECONDS=1 +# shellcheck source=bin/fm-remote-job-lib.sh +. "$ROOT/bin/fm-remote-job-lib.sh" + +fm_remote_job_prepare_state "$ACCOUNT_HOME" || fail "$FM_REMOTE_JOB_ERROR" +rm -f -- "$STATE_ROOT/seq" +SEQ_PIDS=() +for i in $(seq 1 20); do + fm_remote_job_next_seq > "$TMP_ROOT/seq-$i" & + SEQ_PIDS+=("$!") +done +for pid in "${SEQ_PIDS[@]}"; do + wait "$pid" || fail "a concurrent sequence allocator failed" +done +SEQ_RESULTS=$(cat "$TMP_ROOT"/seq-* | sort -n) +SEQ_EXPECTED=$(seq 1 20) +[ "$SEQ_RESULTS" = "$SEQ_EXPECTED" ] \ + || fail "concurrent sequence claims were not unique and monotonic: $SEQ_RESULTS" +[ "$(find "$STATE_ROOT/.seq-claims" -mindepth 1 -maxdepth 1 -type d | wc -l | tr -d ' ')" = 20 ] \ + || fail "concurrent sequence allocations did not retain every durable claim" +mkdir "$STATE_ROOT/.seq-claims/999998" "$STATE_ROOT/.seq-claims/999999" +touch -t 200001010000 "$STATE_ROOT/.seq-claims/999998" +fm_remote_job_reap_stale "$ACCOUNT_HOME" || fail "sequence claim reaping failed" +assert_absent "$STATE_ROOT/.seq-claims/999998" "an expired sequence claim survived stale reaping" +assert_present "$STATE_ROOT/.seq-claims/999999" "a fresh sequence claim was reaped" +mkdir "$STATE_ROOT/.seq-claims/999997" +touch -t 200001010000 "$STATE_ROOT/.seq-claims/999997" +fm_remote_job_reap_stale "$ACCOUNT_HOME" || fail "rate-limited sequence claim reaping failed" +assert_present "$STATE_ROOT/.seq-claims/999997" "sequence claims were rescanned before the hourly interval" +touch -t 200001010000 "$STATE_ROOT/.seq-claims-reaped" +fm_remote_job_reap_stale "$ACCOUNT_HOME" || fail "expired sequence claim reaping failed" +assert_absent "$STATE_ROOT/.seq-claims/999997" "an expired sequence claim survived the next hourly scan" +rmdir "$STATE_ROOT/.seq-claims/999999" +pass "atomic sequence claims remain unique and reap only after expiry" + +fm_on() { + FM_HOME="$LOCAL_HOME" \ + FM_ROOT_OVERRIDE="$REMOTE_ROOT" \ + FM_SSH_BIN="$FAKEBIN/fake-ssh" \ + FM_FAKE_REMOTE_ENTRYPOINT="$REMOTE_ROOT/bin/fm-remote-entrypoint.sh" \ + "$ROOT/bin/fm-on.sh" "$@" +} + +job_state() { # + fm_remote_job_read_state "$STATE_ROOT/jobs/$1" 2>/dev/null || true +} + +wait_for_state() { # + local i=0 + while [ "$i" -lt 200 ]; do + [ "$(job_state "$1")" = "$2" ] && return 0 + i=$((i + 1)) + sleep 0.05 + done + return 1 +} + +HOME="$ACCOUNT_HOME" FM_ROOT_OVERRIDE="$REMOTE_ROOT" FM_REMOTE_JOB_STATE_ROOT="$STATE_ROOT" \ + FM_REMOTE_JOB_PLATFORM_OVERRIDE=Linux \ + "$REMOTE_ROOT/bin/fm-remote-job-worker.sh" > "$TMP_ROOT/worker.out" 2> "$TMP_ROOT/worker.err" & +for _ in $(seq 1 100); do + [ -f "$STATE_ROOT/worker.ready" ] && break + sleep 0.05 +done +assert_present "$STATE_ROOT/worker.ready" "the worker did not publish its readiness heartbeat" + +# T9: home B's job completes while home A runs a long job, and A's queued job +# stays strictly behind A's running job. +LOG_A="$TMP_ROOT/log-a" +LOG_B="$TMP_ROOT/log-b" +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh a1 "$LOG_A" 4 < /dev/null > /dev/null +A1=$FM_REMOTE_JOB_ID +wait_for_state "$A1" running || fail "home A's long job did not begin running" +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh a2 "$LOG_A" 0 < /dev/null > /dev/null +A2=$FM_REMOTE_JOB_ID +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_EDGE" fm-mark-job.sh b1 "$LOG_B" 0 < /dev/null > /dev/null +B1=$FM_REMOTE_JOB_ID +B_BEGAN=$(date +%s) +fm_remote_job_wait "$ACCOUNT_HOME" "$B1" || fail "$FM_REMOTE_JOB_ERROR" +B_ELAPSED=$(( $(date +%s) - B_BEGAN )) +[ "$FM_REMOTE_JOB_EXIT" -eq 0 ] || fail "home B's job behind home A's long job did not complete" +[ "$B_ELAPSED" -le 3 ] || fail "home B's job waited ${B_ELAPSED}s behind home A's long job" +[ "$(job_state "$A1")" = running ] || fail "home A's long job should still be running for the FIFO assertion" +[ "$(cat "$LOG_A")" = a1 ] || fail "home A's queued job ran beside its running job: $(cat "$LOG_A")" +fm_remote_job_reap "$ACCOUNT_HOME" "$B1" || fail "home B's job could not be reaped" +fm_remote_job_wait "$ACCOUNT_HOME" "$A1" || fail "$FM_REMOTE_JOB_ERROR" +fm_remote_job_wait "$ACCOUNT_HOME" "$A2" || fail "$FM_REMOTE_JOB_ERROR" +[ "$(printf '%s' "$(cat "$LOG_A")")" = "$(printf 'a1\na2')" ] \ + || fail "home A's jobs did not execute in stage order: $(cat "$LOG_A")" +fm_remote_job_reap "$ACCOUNT_HOME" "$A1" || fail "home A's first job could not be reaped" +fm_remote_job_reap "$ACCOUNT_HOME" "$A2" || fail "home A's second job could not be reaped" +pass "lanes run homes concurrently while each home stays FIFO" + +# T9 stage order: five jobs staged in rapid succession behind a busy lane must +# execute in staging-sequence order, not the queue directory's random-id order. +: > "$LOG_A" +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh hold "$LOG_A" 2 < /dev/null > /dev/null +HOLD=$FM_REMOTE_JOB_ID +wait_for_state "$HOLD" running || fail "the lane-holding job did not begin running" +RAPID_IDS=() +for tag in r1 r2 r3 r4 r5; do + fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh "$tag" "$LOG_A" 0 < /dev/null > /dev/null + RAPID_IDS+=("$FM_REMOTE_JOB_ID") +done +fm_remote_job_wait "$ACCOUNT_HOME" "$HOLD" || fail "$FM_REMOTE_JOB_ERROR" +fm_remote_job_reap "$ACCOUNT_HOME" "$HOLD" || true +for id in "${RAPID_IDS[@]}"; do + fm_remote_job_wait "$ACCOUNT_HOME" "$id" || fail "$FM_REMOTE_JOB_ERROR" + fm_remote_job_reap "$ACCOUNT_HOME" "$id" || true +done +[ "$(cat "$LOG_A")" = "$(printf 'hold\nr1\nr2\nr3\nr4\nr5')" ] \ + || fail "rapidly staged same-home jobs did not execute in stage order: $(tr '\n' ' ' < "$LOG_A")" +pass "same-home jobs staged in the same second execute in staging-sequence order" + +: > "$LOG_A" +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh publish-hold "$LOG_A" 3 < /dev/null > /dev/null +PUBLISH_HOLD=$FM_REMOTE_JOB_ID +wait_for_state "$PUBLISH_HOLD" running || fail "the publication-order lane holder did not begin running" +( + { + printf 'delayed payload\n' + sleep 5 + } | fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" \ + fm-mark-job.sh delayed "$LOG_A" 0 +) > "$TMP_ROOT/delayed-stage-id" & +DELAYED_STAGE_PID=$! +for _ in $(seq 1 200); do + ls "$STATE_ROOT/jobs"/.stage.* >/dev/null 2>&1 && break + sleep 0.02 +done +ls "$STATE_ROOT/jobs"/.stage.* >/dev/null 2>&1 \ + || fail "the delayed stdin stage did not begin capturing" +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh fast "$LOG_A" 0 < /dev/null > /dev/null +FAST_STAGE=$FM_REMOTE_JOB_ID +wait "$DELAYED_STAGE_PID" || fail "the delayed stdin stage failed to publish" +DELAYED_STAGE=$(cat "$TMP_ROOT/delayed-stage-id") +FAST_SEQ=$(fm_remote_job_read_number "$STATE_ROOT/jobs/$FAST_STAGE" seq) \ + || fail "the fast stage lost its sequence" +DELAYED_SEQ=$(fm_remote_job_read_number "$STATE_ROOT/jobs/$DELAYED_STAGE" seq) \ + || fail "the delayed stage lost its sequence" +[ "$FAST_SEQ" -lt "$DELAYED_SEQ" ] \ + || fail "sequence order did not follow publication order: fast=$FAST_SEQ delayed=$DELAYED_SEQ" +fm_remote_job_wait "$ACCOUNT_HOME" "$PUBLISH_HOLD" || fail "$FM_REMOTE_JOB_ERROR" +fm_remote_job_wait "$ACCOUNT_HOME" "$FAST_STAGE" || fail "$FM_REMOTE_JOB_ERROR" +fm_remote_job_wait "$ACCOUNT_HOME" "$DELAYED_STAGE" || fail "$FM_REMOTE_JOB_ERROR" +[ "$(cat "$LOG_A")" = "$(printf 'publish-hold\nfast\ndelayed')" ] \ + || fail "execution order diverged from publication sequence: $(tr '\n' ' ' < "$LOG_A")" +fm_remote_job_reap "$ACCOUNT_HOME" "$PUBLISH_HOLD" || true +fm_remote_job_reap "$ACCOUNT_HOME" "$FAST_STAGE" || true +fm_remote_job_reap "$ACCOUNT_HOME" "$DELAYED_STAGE" || true +pass "same-home sequence order follows completed staging publication" + +# T3a: a caller killed while its job is still queued cancels it; the worker +# never executes it. +fm_remote_job_stage "$ACCOUNT_HOME" "$REMOTE_ROOT" "$HOME_A" fm-mark-job.sh hold2 "$LOG_A" 4 < /dev/null > /dev/null +HOLD2=$FM_REMOTE_JOB_ID +wait_for_state "$HOLD2" running || fail "the cancellation fixture's lane holder did not begin running" +QUEUED_EFFECT="$TMP_ROOT/queued-cancel-effect" +fm_on ios fm-touch-job.sh "$QUEUED_EFFECT" > /dev/null 2>&1 & +QUEUED_CALLER=$! +QUEUED_JOB= +for _ in $(seq 1 200); do + for job in "$STATE_ROOT"/jobs/job-*; do + [ -d "$job" ] || continue + [ "${job##*/}" = "$HOLD2" ] && continue + [ "$(job_state "${job##*/}")" = queued ] && QUEUED_JOB=${job##*/} && break + done + [ -n "$QUEUED_JOB" ] && break + sleep 0.05 +done +[ -n "$QUEUED_JOB" ] || fail "the doomed caller's job never appeared in the queue" +kill -TERM "$QUEUED_CALLER" 2>/dev/null || true +wait "$QUEUED_CALLER" 2>/dev/null || true +for _ in $(seq 1 200); do + [ ! -d "$STATE_ROOT/jobs/$QUEUED_JOB" ] && break + sleep 0.05 +done +[ ! -d "$STATE_ROOT/jobs/$QUEUED_JOB" ] \ + || fail "the cancelled queued job's record survived (state: $(job_state "$QUEUED_JOB"))" +fm_remote_job_wait "$ACCOUNT_HOME" "$HOLD2" || fail "$FM_REMOTE_JOB_ERROR" +fm_remote_job_reap "$ACCOUNT_HOME" "$HOLD2" || true +sleep 1 +assert_absent "$QUEUED_EFFECT" "the worker executed a queued job whose caller was killed" +pass "a caller killed mid-wait cancels its queued job before execution" + +# T3b: a caller killed while its job is running terminates the job's process +# group instead of letting it run to completion for nobody. +RUN_START="$TMP_ROOT/running-cancel-start" +RUN_FINISH="$TMP_ROOT/running-cancel-finish" +fm_on build fm-two-phase-job.sh "$RUN_START" "$RUN_FINISH" 8 > /dev/null 2>&1 & +RUNNING_CALLER=$! +for _ in $(seq 1 200); do + [ -f "$RUN_START" ] && break + sleep 0.05 +done +assert_present "$RUN_START" "the running-cancellation fixture never started" +kill -TERM "$RUNNING_CALLER" 2>/dev/null || true +wait "$RUNNING_CALLER" 2>/dev/null || true +CANCEL_BEGAN=$(date +%s) +for _ in $(seq 1 200); do + ls "$STATE_ROOT"/jobs/job-* >/dev/null 2>&1 || break + sleep 0.05 +done +CANCEL_ELAPSED=$(( $(date +%s) - CANCEL_BEGAN )) +ls "$STATE_ROOT"/jobs/job-* >/dev/null 2>&1 \ + && fail "the cancelled running job's record survived" +[ "$CANCEL_ELAPSED" -le 6 ] || fail "running-job cancellation took ${CANCEL_ELAPSED}s" +sleep 2 +assert_absent "$RUN_FINISH" "a cancelled running job's process group ran to completion" +pass "a caller killed mid-wait stops its running job's process group" + +# T3c: a caller whose parent exits WITHOUT delivering any signal - the shape a +# dead ssh channel leaves behind - still cancels through the entrypoint's +# parent-liveness probe. +ORPHAN_START="$TMP_ROOT/orphan-cancel-start" +ORPHAN_FINISH="$TMP_ROOT/orphan-cancel-finish" +# shellcheck disable=SC2016 # Expansion is deliberately deferred to the child shell. +env FM_HOME="$LOCAL_HOME" FM_ROOT_OVERRIDE="$REMOTE_ROOT" \ + FM_SSH_BIN="$FAKEBIN/fake-ssh" \ + FM_FAKE_REMOTE_ENTRYPOINT="$REMOTE_ROOT/bin/fm-remote-entrypoint.sh" \ + FM_REMOTE_JOB_STATE_ROOT="$STATE_ROOT" FM_REMOTE_JOB_PLATFORM_OVERRIDE=Linux \ + bash -c ' + "$1/bin/fm-on.sh" build fm-two-phase-job.sh "$2" "$3" 12 >/dev/null 2>&1 & + while [ ! -f "$2" ]; do sleep 0.1; done + ' _ "$ROOT" "$ORPHAN_START" "$ORPHAN_FINISH" +assert_present "$ORPHAN_START" "the orphan-cancellation fixture never started" +ORPHAN_BEGAN=$(date +%s) +for _ in $(seq 1 300); do + ls "$STATE_ROOT"/jobs/job-* >/dev/null 2>&1 || break + sleep 0.05 +done +ORPHAN_ELAPSED=$(( $(date +%s) - ORPHAN_BEGAN )) +ls "$STATE_ROOT"/jobs/job-* >/dev/null 2>&1 \ + && fail "the orphaned caller's job record survived its disconnect" +[ "$ORPHAN_ELAPSED" -le 10 ] || fail "orphan-disconnect cancellation took ${ORPHAN_ELAPSED}s" +sleep 2 +assert_absent "$ORPHAN_FINISH" "a job abandoned by a signal-less disconnect ran to completion" +pass "a signal-less caller disconnect cancels the abandoned job through the parent probe" + +# T3: after the cancellations, a burst of short bounded commands meets its own +# budget - no convoy behind abandoned work. +BURST_BEGAN=$(date +%s) +for tag in c1 c2 c3; do + rc=0 + fm_run_timed 15 env FM_HOME="$LOCAL_HOME" FM_ROOT_OVERRIDE="$REMOTE_ROOT" \ + FM_SSH_BIN="$FAKEBIN/fake-ssh" \ + FM_FAKE_REMOTE_ENTRYPOINT="$REMOTE_ROOT/bin/fm-remote-entrypoint.sh" \ + FM_REMOTE_JOB_STATE_ROOT="$STATE_ROOT" FM_REMOTE_JOB_PLATFORM_OVERRIDE=Linux \ + "$ROOT/bin/fm-on.sh" ios fm-touch-job.sh "$TMP_ROOT/burst-$tag" >/dev/null 2>&1 || rc=$? + [ "$rc" -eq 0 ] || fail "post-cancellation burst command $tag failed with $rc" + assert_present "$TMP_ROOT/burst-$tag" "post-cancellation burst command $tag did not run" +done +BURST_ELAPSED=$(( $(date +%s) - BURST_BEGAN )) +[ "$BURST_ELAPSED" -le 12 ] || fail "the post-cancellation burst convoyed for ${BURST_ELAPSED}s" +pass "bounded reads after a cancellation meet their own budget with no convoy" + +# T6: a non-payload call with an OPEN stdin pipe completes instead of wedging +# staging on a stdin capture that never reaches EOF. +printf 'rsm\n' > "$HOME_A/.fm-secondmate-home" +printf '# fixture secondmate home\n' > "$HOME_A/AGENTS.md" +mkdir -p "$HOME_A/state" "$HOME_A/bin" +rc=0 +fm_run_timed 20 env FM_HOME="$LOCAL_HOME" FM_ROOT_OVERRIDE="$REMOTE_ROOT" \ + FM_SSH_BIN="$FAKEBIN/fake-ssh" \ + FM_FAKE_REMOTE_ENTRYPOINT="$REMOTE_ROOT/bin/fm-remote-entrypoint.sh" \ + FM_REMOTE_JOB_STATE_ROOT="$STATE_ROOT" FM_REMOTE_JOB_PLATFORM_OVERRIDE=Linux \ + "$ROOT/bin/fm-on.sh" ios fm-remote-secondmate-control.sh state rsm \ + < <(sleep 30) > "$TMP_ROOT/state-out" 2> "$TMP_ROOT/state-err" || rc=$? +[ "$rc" -ne 124 ] || fail "a control-state call with an open stdin pipe wedged staging" +assert_grep 'missing' "$TMP_ROOT/state-out" \ + "the control-state call did not complete through the worker: $(cat "$TMP_ROOT/state-err")" +pass "an open caller stdin no longer wedges a non-payload remote command" + +# A live explicit stdin stage can exceed the litter age while waiting for EOF; +# the stale sweep must retain it until its owning entrypoint publishes the job. +rc=0 +{ + printf 'slow payload one\n' + sleep 3 + printf 'slow payload two\n' +} | fm_on --stdin ios fm-stdin-probe.sh > "$TMP_ROOT/slow-payload-out" 2> "$TMP_ROOT/slow-payload-err" || rc=$? +expect_code 0 "$rc" "a live slow stdin stage must survive stale reaping: $(cat "$TMP_ROOT/slow-payload-err")" +assert_grep 'stdin=slow payload one' "$TMP_ROOT/slow-payload-out" "the slow stdin stage lost its first bytes" +assert_grep 'stdin=slow payload two' "$TMP_ROOT/slow-payload-out" "the slow stdin stage was reaped before EOF" +pass "a live explicit-stdin stage survives the staging-litter age bound" + +# T6: a payload caller with --stdin still delivers its bytes. +printf 'payload byte one\npayload byte two\n' > "$TMP_ROOT/payload" +fm_on --stdin ios fm-stdin-probe.sh < "$TMP_ROOT/payload" > "$TMP_ROOT/payload-out" 2>/dev/null \ + || fail "the --stdin payload call failed" +assert_grep 'stdin=payload byte one' "$TMP_ROOT/payload-out" "--stdin did not deliver the payload" +assert_grep 'stdin=payload byte two' "$TMP_ROOT/payload-out" "--stdin lost part of the payload" +pass "--stdin still delivers a payload caller's bytes" + +# Stage litter: an abandoned .stage.* older than the reap age does not survive +# a worker pass, while fresh staging is left alone. +OLD_STAGE="$STATE_ROOT/jobs/.stage.abandoned" +FRESH_STAGE="$STATE_ROOT/jobs/.stage.fresh" +mkdir -p "$OLD_STAGE" "$FRESH_STAGE" +touch -t 200001010000 "$OLD_STAGE" +for _ in $(seq 1 100); do + [ ! -d "$OLD_STAGE" ] && break + sleep 0.05 +done +[ ! -d "$OLD_STAGE" ] || fail "stage litter older than the reap age survived the worker pass" +assert_present "$FRESH_STAGE" "the worker reaped fresh staging that is still in use" +rmdir "$FRESH_STAGE" +pass "abandoned stage litter is reaped by age while fresh staging survives" + +echo "ALL TESTS PASSED" diff --git a/tests/fm-send-remote-delivery.test.sh b/tests/fm-send-remote-delivery.test.sh index 8a686dc9cb0..ee3736b014c 100755 --- a/tests/fm-send-remote-delivery.test.sh +++ b/tests/fm-send-remote-delivery.test.sh @@ -105,6 +105,12 @@ count=$(cat "$FM_SSH_COUNT" 2>/dev/null || echo 0) count=$((count + 1)) printf '%s\n' "$count" > "$FM_SSH_COUNT" printf '%s\n' "$*" >> "$FM_SSH_LOG" +if [ -n "${FM_FAKE_SSH_HANG:-}" ]; then + # A busy remote lane: the transport attempt never returns on its own. The + # real sleep, because the stubbed one on PATH returns immediately. + /bin/sleep "$FM_FAKE_SSH_HANG" + exit 255 +fi if [ "${FM_FAKE_SSH_AFTER_AMBIGUOUS_RC:-0}" -ne 0 ] && [ "$count" -gt 1 ]; then exit "$FM_FAKE_SSH_AFTER_AMBIGUOUS_RC" fi @@ -603,6 +609,95 @@ test_remote_transport_loss_preserves_expectation() { pass "fm-send remote: ssh 255 fails with resend-safe guidance and preserves the expectation" } +test_remote_send_budget_bounds_busy_lane() { + local dir fb ssh_log home rhome rc err began elapsed count pend delivery corr ssh_before + dir="$TMP_ROOT/remote-budget"; mkdir -p "$dir" + fb=$(make_stubs "$dir"); ssh_log="$dir/ssh.log"; : > "$ssh_log" + rhome=$(setup_remote_secondmate_home remote-budget) + home=$(setup_remote_parent_home remote-budget "$rhome") + delivery=aaaabbbbccccdddd + + rc=0 + send_env "$fb" "$home" "$ssh_log" FM_SEND_REMOTE_BUDGET=invalid \ + "$SEND" rsm --key Enter >"$dir/key-invalid.out" 2>"$dir/key-invalid.err" || rc=$? + [ "$rc" -ne 0 ] || fail "an invalid remote key budget must fail" + assert_contains "$(cat "$dir/key-invalid.err")" "must be a positive integer" \ + "an invalid remote key budget must explain its validation failure" + [ ! -f "$ssh_log.count" ] || fail "an invalid remote key budget reached the transport" + + began=$(date +%s) + rc=0 + send_env "$fb" "$home" "$ssh_log" FM_FAKE_SSH_HANG=60 FM_SEND_REMOTE_BUDGET=2 \ + "$SEND" rsm --key Enter >"$dir/key.out" 2>"$dir/key.err" || rc=$? + elapsed=$(( $(date +%s) - began )) + expect_code 1 "$rc" "a bounded remote key must preserve the existing failure contract" + [ "$elapsed" -le 15 ] || fail "the bounded remote key waited ${elapsed}s behind the busy lane" + assert_contains "$(cat "$dir/key.err")" "completion may be unknown" \ + "a bounded remote key failure must preserve its existing diagnostic" + [ "$(cat "$ssh_log.count")" = 1 ] \ + || fail "a bounded remote key must make exactly one transport attempt" + printf '0\n' > "$ssh_log.count" + + # T5: a fire-and-forget send to a mate behind a busy lane returns its + # unconfirmed result within its own budget instead of waiting the lane out. + began=$(date +%s) + rc=0 + send_env "$fb" "$home" "$ssh_log" FM_FAKE_SSH_HANG=60 FM_SEND_REMOTE_BUDGET=2 \ + "$SEND" rsm --fire-and-forget "$delivery" "reconcile your own books" \ + >"$dir/out" 2>"$dir/err" || rc=$? + elapsed=$(( $(date +%s) - began )) + err=$(cat "$dir/err") + expect_code 3 "$rc" "a budget-bounded fire-and-forget send must report unconfirmed: $err" + [ "$elapsed" -le 15 ] || fail "the bounded send waited ${elapsed}s behind the busy lane" + assert_contains "$err" "delivery-id=$delivery" \ + "the bounded unconfirmed result must name the reusable delivery id" + [ "$(cat "$ssh_log.count")" = 1 ] \ + || fail "a budget hit must not retry into the same busy lane, got $(cat "$ssh_log.count") attempts" + + # A retry with the same delivery id against the recovered lane dedups onto + # the same remote record. + send_env "$fb" "$home" "$ssh_log" \ + "$SEND" rsm --fire-and-forget "$delivery" "reconcile your own books" \ + >"$dir/retry.out" 2>"$dir/retry.err" \ + || fail "the same-delivery-id retry after the budget hit failed" + count=$(remote_inbox_records "$rhome" | grep -c . || true) + [ "$count" = 1 ] || fail "the same-delivery-id retry did not dedup onto one record, found $count" + + # A reply-bearing send names the budget and prints the correlation-reusing + # resend command, with the expectation preserved as delivery-unknown. + rc=0 + send_env "$fb" "$home" "$ssh_log" FM_FAKE_SSH_HANG=60 FM_SEND_REMOTE_BUDGET=2 \ + "$SEND" rsm "please rename the metric" >"$dir/reply.out" 2>"$dir/reply.err" || rc=$? + err=$(cat "$dir/reply.err") + [ "$rc" -ne 0 ] || fail "a budget-bounded reply-bearing send must not claim confirmed delivery" + assert_contains "$err" "within its 2s budget" \ + "the budget-bounded failure must name the budget that bounded it" + assert_contains "$err" "Only the correlation-reusing resend below is idempotent" \ + "the budget-bounded failure must print the supported safe resend boundary" + pend=$(pending_record "$home") + [ -n "$pend" ] || fail "a budget-bounded reply-bearing send must preserve its expectation" + [ "$(grep '^phase=' "$pend" | tail -1 | cut -d= -f2-)" = delivery_unknown ] \ + || fail "the preserved expectation must record unknown delivery: $(cat "$pend")" + + # Invalid transport configuration fails before a correlation-reusing resend + # mutates the preserved expectation or reaches the transport. + corr=$(fm_pending_reply_get "$pend" corr_id) + cp "$pend" "$dir/pending-before-invalid-budget" + ssh_before=$(cat "$ssh_log.count") + rc=0 + send_env "$fb" "$home" "$ssh_log" FM_SEND_REMOTE_BUDGET=invalid \ + FM_PENDING_REPLY_EXISTING_CORR="$corr" \ + "$SEND" rsm "please rename the metric" >"$dir/invalid.out" 2>"$dir/invalid.err" || rc=$? + [ "$rc" -ne 0 ] || fail "an invalid remote budget must fail the resend" + assert_contains "$(cat "$dir/invalid.err")" "must be a positive integer" \ + "an invalid remote budget must explain its validation failure" + [ "$(cat "$ssh_log.count")" = "$ssh_before" ] \ + || fail "an invalid remote budget reached the remote transport" + cmp -s "$dir/pending-before-invalid-budget" "$pend" \ + || fail "an invalid remote budget mutated the reusable pending expectation: $(cat "$pend")" + pass "fm-send remote: the remote leg is budget-bounded and stays idempotent across the bound" +} + test_local_secondmate_pending_keeps_expectation_armed() { local dir fb log home rc rec corr dir="$TMP_ROOT/local-pending-expectation"; mkdir -p "$dir" @@ -693,6 +788,7 @@ test_remote_slash_rides_inbox test_remote_real_failure_still_fails test_remote_exit3_no_longer_delivered test_remote_transport_loss_preserves_expectation +test_remote_send_budget_bounds_busy_lane test_local_pending_reports_delivered_unconfirmed test_local_pending_does_not_close_resolve_key test_local_secondmate_pending_keeps_expectation_armed