diff --git a/.agents/skills/operational-home-layout/SKILL.md b/.agents/skills/operational-home-layout/SKILL.md index f51340058bb..de91784fb17 100644 --- a/.agents/skills/operational-home-layout/SKILL.md +++ b/.agents/skills/operational-home-layout/SKILL.md @@ -110,7 +110,7 @@ state/ runtime records and signals; gitignored .afk-contract the away or quiet posture record; bin/fm-afk-contract.sh owns its mode, schema, entry, archive, and lock contract; its sibling .afk-contract.lock serializes actions authorized by the live record afk-contracts/ archived away and quiet records; bin/fm-afk-contract.sh owns their archive contract .afk durable away/quiet-mode daemon flag on the harnesses that still launch the daemon (never on Pi); present = sub-supervisor may inject escalations, first line `away` (default, set by /afk, cleared on user return) or `quiet` (set by /quiet, cleared only on explicit /quiet off) per the single owner fm_afk_mode() in bin/fm-wake-lib.sh - .lock-session trusted Claude session-lock sidecar; written only by bin/fm-lock.sh; never touch + .lock-session trusted Claude session or Codex thread session-lock sidecar; written only by bin/fm-lock.sh; never touch .watch.lock .wake-queue.lock watcher singleton and queue serialization locks .claude-autoarm.lock .claude-autoarm-epoch .claude-autoarm-failure-notified .claude-autoarm-failure-alarmed .turnend-claude-blocks .turnend-claude-blocks.lock Claude Stop auto-arm single-flight, epoch, failure-episode, attended-alarm, guard-budget, and budget-lock records; never touch .cursor-park-owner .cursor-park-owner.lock .turnend-cursor-blocks Cursor stop-hook owner record, publication and commit lock, and bounded repair-nag budget; never touch diff --git a/bin/fm-lock.sh b/bin/fm-lock.sh index 5c954599f99..d6b9da912b2 100755 --- a/bin/fm-lock.sh +++ b/bin/fm-lock.sh @@ -18,13 +18,33 @@ # and left byte-identical when it already names that id. A same-session # confirmation never rewrites line 1 while the recorded pid is alive, because # bin/fm-startup-network.sh compares that pid across its deferred sweeps; a dead -# recorded pid is reclaimed and rewritten to this session's anchor. +# recorded pid is reclaimed and rewritten to this session's anchor. An anchor +# that is the shared managed Codex daemon is refused without a trusted thread +# id, since no sidecar could then tell its threads apart. # # Usage: fm-lock.sh acquire; exit 1 unless ownership is verified # fm-lock.sh status print holder and liveness; always exits 0. # A held lock is not proof the holder is consuming # wakes. Machine-readable lock fields live on # fm-inbox.sh ready, from the same inspect helper. +# fm-lock.sh take-over --expect-pid PID --expect-session codex:ID|none \ +# --attest-owner-ended PID/codex:ID|none +# Explicitly replace an idle shared Codex daemon +# lock after matching its exact pid and sidecar and +# attesting that the previous thread has ended. +# Refuses while a watcher is live or either lock +# or watcher beacon was recently active. Those +# guards do not prove the thread ended: taking +# over an active thread can split one home between +# two primaries. The operator must verify closure. +# Backs up the prior lock before replacement. +# A Codex Desktop primary live during the upgrade +# to recorded Codex thread ids has no sidecar and +# sees its own lock as foreign. Recover it: stop +# that thread's watcher (fm-watch-arm.sh --stop), +# wait out the quiet window, then run take-over +# --expect-pid PID --expect-session none with +# --attest-owner-ended PID/none. set -u SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" @@ -49,13 +69,82 @@ if [ "${1:-}" = "status" ]; then case "$FM_LOCK_INSPECT_STATE" in free) echo "lock: free" ;; unreadable) echo "lock: unreadable" ;; - held) echo "lock: held by live harness pid $FM_LOCK_INSPECT_PID" ;; + held) + if fm_session_lock_shared_codex_pid "$FM_LOCK_INSPECT_PID"; then + recorded=$(fm_session_lock_recorded_session_id "$STATE" 2>/dev/null || true) + [ -n "$recorded" ] || recorded=none + case "$recorded" in + none) + echo "lock: held by live harness pid $FM_LOCK_INSPECT_PID (managed Codex daemon; session none)" + echo "recovery after verified idle: bin/fm-lock.sh take-over --expect-pid $FM_LOCK_INSPECT_PID --expect-session none" + echo "verify the old thread ended, then add --attest-owner-ended $FM_LOCK_INSPECT_PID/none; an active-thread takeover can split ownership" + ;; + codex:*) + codex_id=${recorded#codex:} + case "$codex_id" in + ''|*[!A-Za-z0-9_-]*) echo "lock: held by live harness pid $FM_LOCK_INSPECT_PID (managed Codex daemon; unrecognized sidecar)" ;; + *) + echo "lock: held by live harness pid $FM_LOCK_INSPECT_PID (managed Codex daemon; session $recorded)" + echo "recovery after verified idle: bin/fm-lock.sh take-over --expect-pid $FM_LOCK_INSPECT_PID --expect-session $recorded" + echo "verify the old thread ended, then add --attest-owner-ended $FM_LOCK_INSPECT_PID/$recorded; an active-thread takeover can split ownership" + ;; + esac + ;; + *) echo "lock: held by live harness pid $FM_LOCK_INSPECT_PID (managed Codex daemon; unrecognized sidecar)" ;; + esac + else + echo "lock: held by live harness pid $FM_LOCK_INSPECT_PID" + fi + ;; *) echo "lock: stale (pid $FM_LOCK_INSPECT_PID dead or not a harness)" ;; esac exit 0 fi +TAKEOVER=0 +EXPECT_PID= +EXPECT_SESSION= +ATTEST_OWNER= +TAKEOVER_BACKUP= +if [ "${1:-}" = take-over ]; then + [ "$#" -eq 7 ] && [ "$2" = --expect-pid ] && [ "$4" = --expect-session ] \ + && [ "$6" = --attest-owner-ended ] || { + echo "usage: fm-lock.sh take-over --expect-pid PID --expect-session codex:ID|none --attest-owner-ended PID/codex:ID|none" >&2 + exit 2 + } + TAKEOVER=1 + EXPECT_PID=$3 + EXPECT_SESSION=$5 + ATTEST_OWNER=$7 + case "$EXPECT_PID" in ''|*[!0-9]*) echo "error: expected pid must be numeric" >&2; exit 2 ;; esac + case "$EXPECT_SESSION" in + none) ;; + codex:*) + codex_id=${EXPECT_SESSION#codex:} + case "$codex_id" in ''|*[!A-Za-z0-9_-]*) echo "error: expected session must be codex:ID or none" >&2; exit 2 ;; esac + ;; + *) echo "error: expected session must be codex:ID or none" >&2; exit 2 ;; + esac + [ "$ATTEST_OWNER" = "$EXPECT_PID/$EXPECT_SESSION" ] || { + echo "error: attest-owner-ended must name the exact previous pid/session; verify that thread has ended" >&2 + exit 2 + } +elif [ "$#" -ne 0 ]; then + echo "usage: fm-lock.sh [status|take-over --expect-pid PID --expect-session codex:ID|none --attest-owner-ended PID/codex:ID|none]" >&2 + exit 2 +fi + me=$(fm_session_lock_anchor_pid) || { echo "error: cannot locate harness process in ancestry" >&2; exit 1; } +if ! fm_session_lock_trusted_session_id >/dev/null; then + if [ "$TAKEOVER" -eq 1 ]; then + echo "error: take-over requires a verified Claude session or Codex thread identity" >&2 + exit 1 + fi + if fm_session_lock_shared_codex_pid "$me"; then + echo "error: refusing to record the shared managed Codex daemon pid $me without a trusted thread id; export a valid CODEX_THREAD_ID ([A-Za-z0-9_-], at most 128 chars) or run the primary under a per-session harness anchor" >&2 + exit 1 + fi +fi probe=$(mktemp "$STATE/.lock-write.XXXXXX" 2>/dev/null) || { echo "error: cannot write session lock; operate read-only until resolved" >&2 exit 1 @@ -166,7 +255,7 @@ confirm_own_lock() { # waited=1 fi recorded=$(cat "$LOCK" 2>/dev/null || true) - if [ "$recorded" = "$me" ] || fm_session_lock_owned_by_self "$STATE"; then + if fm_session_lock_owned_by_self "$STATE"; then publish_lock_session_or_die commit_lock_session release_claim_lock @@ -179,6 +268,58 @@ confirm_own_lock() { # return 1 } +takeover_preflight() { + local old recorded watcher_pid quiet=${FM_LOCK_TAKEOVER_QUIET_SECONDS:-300} + case "$quiet" in ''|*[!0-9]*) die_takeover "FM_LOCK_TAKEOVER_QUIET_SECONDS must be numeric" ;; esac + [ "$quiet" -ge 60 ] && [ "$quiet" -le 3600 ] || die_takeover "quiet interval must be 60-3600 seconds" + [ -f "$LOCK" ] && [ ! -L "$LOCK" ] || die_takeover "session lock is missing or unsafe" + old=$(cat "$LOCK" 2>/dev/null) || die_takeover "session lock is unreadable" + [ "$old" = "$EXPECT_PID" ] || die_takeover "session lock pid changed" + if [ -e "$LOCK_SESSION" ] || [ -L "$LOCK_SESSION" ]; then + [ -f "$LOCK_SESSION" ] && [ ! -L "$LOCK_SESSION" ] \ + || die_takeover "session identity is unsafe" + recorded=$(fm_session_lock_recorded_session_id "$STATE") \ + || die_takeover "session identity is unreadable" + else + recorded=none + fi + [ "$recorded" = "$EXPECT_SESSION" ] || die_takeover "session identity changed" + fm_session_lock_shared_codex_pid "$old" || die_takeover "recorded owner is not a live managed Codex daemon" + fm_harness_pid_alive "$old" || die_takeover "recorded owner cannot be verified" + if [ -e "$STATE/.watch.lock" ] || [ -L "$STATE/.watch.lock" ]; then + [ -d "$STATE/.watch.lock" ] && [ ! -L "$STATE/.watch.lock" ] \ + || die_takeover "watcher lock is unsafe" + [ -f "$STATE/.watch.lock/pid" ] && [ ! -L "$STATE/.watch.lock/pid" ] \ + || die_takeover "watcher identity is unreadable" + watcher_pid=$(cat "$STATE/.watch.lock/pid" 2>/dev/null) \ + || die_takeover "watcher identity is unreadable" + case "$watcher_pid" in ''|*[!0-9]*) die_takeover "watcher identity is malformed" ;; esac + if fm_pid_alive "$watcher_pid"; then + die_takeover "a watcher is still live; stop or wait for its own handoff" + fi + fi + [ ! -L "$STATE/.last-watcher-beat" ] || die_takeover "watcher beacon is unsafe" + [ "$(fm_path_age "$LOCK")" -ge "$quiet" ] || die_takeover "session lock is still recent" + [ "$(fm_path_age "$STATE/.last-watcher-beat")" -ge "$quiet" ] || die_takeover "watcher beacon is still recent" +} + +die_takeover() { echo "error: take-over refused: $1" >&2; exit 1; } + +backup_takeover_lock() { + local backup + backup=$(umask 077; mktemp -d "$STATE/.lock-takeover.XXXXXX") \ + || die_takeover "cannot create a prior-lock backup" + if ! { cp -P "$LOCK" "$backup/lock" \ + && { [ "$EXPECT_SESSION" = none ] \ + || cp -P "$LOCK_SESSION" "$backup/lock-session"; } \ + && printf '%s\n' "$EXPECT_PID/$EXPECT_SESSION" > "$backup/expected-owner"; }; then + rm -f "$backup/lock" "$backup/lock-session" "$backup/expected-owner" + rmdir "$backup" 2>/dev/null || true + die_takeover "cannot preserve the prior lock" + fi + TAKEOVER_BACKUP=$backup +} + refuse_live_owner() { # local recorded if recorded=$(fm_session_lock_recorded_session_id "$STATE"); then @@ -189,9 +330,9 @@ refuse_live_owner() { # exit 1 } -if [ -f "$LOCK" ] && [ ! -L "$LOCK" ]; then +if [ "$TAKEOVER" -eq 0 ] && [ -f "$LOCK" ] && [ ! -L "$LOCK" ]; then old=$(cat "$LOCK" 2>/dev/null || true) - if [ "$old" = "$me" ] || fm_session_lock_owned_by_self "$STATE"; then + if fm_session_lock_owned_by_self "$STATE"; then confirm_own_lock "$old" old=$(cat "$LOCK" 2>/dev/null || true) fi @@ -210,7 +351,12 @@ if ! fm_lock_try_acquire "$CLAIM_LOCK"; then fi CLAIM_LOCK_HELD=1 -if [ -e "$LOCK" ] || [ -L "$LOCK" ]; then +if [ "$TAKEOVER" -eq 1 ]; then + takeover_preflight + backup_takeover_lock +fi + +if [ "$TAKEOVER" -eq 0 ] && { [ -e "$LOCK" ] || [ -L "$LOCK" ]; }; then if [ ! -f "$LOCK" ] || [ -L "$LOCK" ]; then echo "error: session lock is not a regular file; operate read-only until resolved" >&2 exit 1 @@ -219,12 +365,8 @@ if [ -e "$LOCK" ] || [ -L "$LOCK" ]; then echo "error: session lock is unreadable; operate read-only until resolved" >&2 exit 1 } - if [ "$old" != "$me" ] && fm_harness_pid_alive "$old"; then - fm_session_lock_owned_by_self "$STATE" && confirm_own_lock "$old" - old=$(cat "$LOCK" 2>/dev/null || true) - if [ "$old" != "$me" ] && fm_harness_pid_alive "$old"; then - refuse_live_owner "$old" - fi + if ! fm_session_lock_owned_by_self "$STATE" && fm_harness_pid_alive "$old"; then + refuse_live_owner "$old" fi fi # The sidecar goes first: a fresh pid beside a previous session's id would let @@ -272,4 +414,8 @@ if [ ! -f "$LOCK" ] || [ -L "$LOCK" ] || [ "$written" != "$me" ]; then fi commit_lock_session release_claim_lock -echo "lock acquired: harness pid $me" +if [ "$TAKEOVER" -eq 1 ]; then + echo "lock taken over: prior managed Codex daemon pid $EXPECT_PID; harness pid $me; prior lock backup $TAKEOVER_BACKUP" +else + echo "lock acquired: harness pid $me" +fi diff --git a/bin/fm-procevent-remote-reply.sh b/bin/fm-procevent-remote-reply.sh index 7d545c0508b..b2389d4d23d 100755 --- a/bin/fm-procevent-remote-reply.sh +++ b/bin/fm-procevent-remote-reply.sh @@ -10,18 +10,20 @@ # fm-procevent-remote-reply.sh self-announcing # fm-procevent-remote-reply.sh source-id # fm-procevent-remote-reply.sh relisten +# fm-procevent-remote-reply.sh lag-check +# fm-procevent-remote-reply.sh rebase --expect-offset # fm-procevent-remote-reply.sh retire # # `arm` registers one blocking, non-destructive delta source for the remote # home's state/parent-replies.status log. The process-event runner owns blocking, -# capture, publication, and one machine-wide source owner. Each captured delta is -# terminal for that exact registration; `handle` validates and idempotently -# ingests it, acknowledges the captured generation, then registers the next -# cursor-anchored source. `relisten` tells that runner to poll again in the same -# process, still holding the claim, after an empty window and after that re-arm. +# capture, publication, and one machine-wide source owner. Each captured delta +# is applied through `handle`, which validates and idempotently ingests it, +# acknowledges the captured generation, then registers the next cursor-anchored +# source. A delta is not terminal for the runner: `relisten` polls again in the +# same process, still holding the claim, after an empty window and after re-arm. # A window the remote job worker preempted is reported to the runner as an empty # window, so it relistens too (see JOB_PREEMPTED below). -# A continuity break is escalated and not re-armed, so the registration is dropped +# Only a continuity break is terminal: it is escalated and not re-armed, so the registration is dropped # and the runner stops. The runner does not refresh the owner lease. # # `autohandle` is the runner's own entry into that same `handle`: it takes the @@ -40,6 +42,13 @@ # completely quiet. Only a capture autohandle could NOT # fully apply is published as a `check` wake for the manual handler, and # running `handle` on that wake is idempotent. +# A shortened log remains a continuity break. `rebase` is the explicit recovery +# after the terminal result was handled: it requires the exact current cursor, +# an acknowledged truncation result, and the last ingested delta result. It +# re-reads the remote log from that delta's verified start, compares every +# retained byte with the ingested payload, rewinds only to that proven +# complete-line boundary, and re-arms the listener. It refuses changed bytes, +# missing receipts, an active registration, and a non-shortened source. # # This channel is a status-stream MIRROR, not a correlated-reply channel. A local # secondmate appends its whole status stream straight into the parent's @@ -95,6 +104,8 @@ DOCUMENT_LOCAL_FAILURE=2 . "$SCRIPT_DIR/fm-secondmate-registry-lib.sh" # shellcheck source=bin/fm-pending-reply-lib.sh . "$SCRIPT_DIR/fm-pending-reply-lib.sh" +# shellcheck source=bin/fm-timeout-lib.sh +. "$SCRIPT_DIR/fm-timeout-lib.sh" die() { printf 'error: %s\n' "$1" >&2; exit 1; } usage() { sed -n '2,66p' "$0" | sed 's/^# \{0,1\}//'; exit 2; } @@ -161,7 +172,10 @@ write_cursor() { # printf 'prefix_sha256=%s\n' "$hash" } > "$tmp" || { rm -f -- "$tmp"; return 1; } chmod 600 "$tmp" || { rm -f -- "$tmp"; return 1; } - mv -f -- "$tmp" "$path" + mv -f -- "$tmp" "$path" || return 1 + # Cursor progress ends the old lag episode before another network probe is + # due. A watcher also compares the receipt against this committed cursor. + rm -f -- "$CURSOR_DIR/$id.lag" "$CURSOR_DIR/$id.lag-ready" || true } ingest_receipt_matches() { # @@ -268,20 +282,290 @@ WINDOW_CLOSED_EMPTY=75 JOB_PREEMPTED=76 cmd_source() { - local id=${1:-} started rc=0 + local id=${1:-} started rc=0 attempt marker key reason validate_id "$id" - read_cursor "$id" - started=$(fm_pending_reply_now) - "$SCRIPT_DIR/fm-on.sh" "$id" fm-remote-delta-read.sh \ - "$REMOTE_LOG" "$CURSOR_OFFSET" "$CURSOR_HASH" "$WAIT_SECONDS" < /dev/null || rc=$? - if [ "$rc" -eq "$WINDOW_CLOSED_EMPTY" ]; then - fm_pending_reply_note_remote_channel_caught_up "$STATE" "$id" "$started" || true - elif [ "$rc" -eq "$JOB_PREEMPTED" ]; then - rc=$WINDOW_CLOSED_EMPTY + mkdir -p "$CURSOR_DIR" || return 1 + [ -d "$CURSOR_DIR" ] && [ ! -L "$CURSOR_DIR" ] || return 1 + marker="$CURSOR_DIR/$id.source-failed" + [ ! -e "$marker" ] || { [ -f "$marker" ] && [ ! -L "$marker" ]; } || return 1 + SOURCE_ATTEMPT=$(umask 077; mktemp "$CURSOR_DIR/.source.XXXXXX") || return 1 + trap 'rm -f -- "$SOURCE_ATTEMPT"' EXIT + # The runner retains its claim during brief transport/read failures. Three + # consecutive failures open one durable failure episode and exit the runner, + # leaving recovery to reconcile's launch floor. + for attempt in 1 2 3; do + read_cursor "$id" + started=$(fm_pending_reply_now) + rc=0 + "$SCRIPT_DIR/fm-on.sh" "$id" fm-remote-delta-read.sh \ + "$REMOTE_LOG" "$CURSOR_OFFSET" "$CURSOR_HASH" "$WAIT_SECONDS" < /dev/null > "$SOURCE_ATTEMPT" || rc=$? + case "$rc" in + 0|"$WINDOW_CLOSED_EMPTY"|"$JOB_PREEMPTED") + cat -- "$SOURCE_ATTEMPT" || return 1 + rm -f -- "$marker" + if [ "$rc" -eq "$WINDOW_CLOSED_EMPTY" ]; then + fm_pending_reply_note_remote_channel_caught_up "$STATE" "$id" "$started" || true + elif [ "$rc" -eq "$JOB_PREEMPTED" ]; then + rc=$WINDOW_CLOSED_EMPTY + fi + return "$rc" + ;; + esac + [ "$attempt" -eq 3 ] || sleep "$attempt" + done + key="remote-reply-source-failed-$id" + reason="check: remote reply listener $id failed three consecutive reads (last exit $rc); inspect the remote route and re-ensure remote-reply-$id" + if [ ! -e "$marker" ]; then + if fm_wake_append check "$key" "$reason"; then + (umask 077; printf '%s\n' "$rc" > "$marker") || true + fi fi return "$rc" } +# The process-event claim can disappear after a failed remote read while its +# registration remains. A watcher normally re-ensures it, but a live watcher +# alone does not prove the cursor is advancing. Probe the remote append-only +# log's size on a bounded cadence and, once per episode where it stays ahead of +# the committed cursor, re-ensure the listener and announce the episode. A moving cursor resets the episode; a failed probe +# proves no lag and leaves the existing observation untouched. +cmd_lag_check() { # + local id=${1:-} threshold=${FM_REMOTE_REPLY_LAG_SECONDS:-120} + local cadence=${FM_REMOTE_REPLY_LAG_PROBE_SECONDS:-} now size marker probe + local prior_offset prior_since prior_alerted since key reason tmp ready + validate_id "$id" + remote_route_exists "$id" + if [ -z "$cadence" ]; then + case "$WAIT_SECONDS" in ''|*[!0-9]*) die "remote reply wait window must be whole seconds" ;; esac + [ "$WAIT_SECONDS" -le 300 ] || die "remote reply wait window exceeds its safety bound" + cadence=$((10#$WAIT_SECONDS + 35)) + fi + case "$threshold:$cadence" in *[!0-9:]*|:*|*:) die "lag thresholds must be whole seconds" ;; esac + [ "$threshold" -ge 1 ] && [ "$threshold" -le 3600 ] || die "lag threshold must be 1-3600 seconds" + [ "$cadence" -ge 1 ] && [ "$cadence" -le 360 ] || die "lag probe cadence must be 1-360 seconds" + mkdir -p "$CURSOR_DIR" || return 1 + [ -d "$CURSOR_DIR" ] && [ ! -L "$CURSOR_DIR" ] || return 1 + marker="$CURSOR_DIR/$id.lag" + probe="$CURSOR_DIR/$id.lag-probe" + ready="$CURSOR_DIR/$id.lag-ready" + [ ! -L "$marker" ] && [ ! -L "$probe" ] && [ ! -L "$ready" ] || return 1 + [ "$(fm_path_age "$probe")" -ge "$cadence" ] || return 0 + touch "$probe" || return 1 + read_cursor "$id" + size=$(fm_run_timed 15 "$SCRIPT_DIR/fm-on.sh" "$id" fm-remote-delta-read.sh size "$REMOTE_LOG" 2>/dev/null) || return 0 + case "$size" in ''|*[!0-9]*) return 0 ;; esac + [ "${#size}" -le 18 ] || return 0 + if [ "$size" -le "$CURSOR_OFFSET" ]; then + [ ! -L "$marker" ] || return 1 + rm -f -- "$marker" "$ready" + return 0 + fi + now=$(date +%s) + prior_offset= + prior_since= + prior_alerted=0 + if [ -e "$marker" ] || [ -L "$marker" ]; then + [ -f "$marker" ] && [ ! -L "$marker" ] || return 1 + read -r prior_offset prior_since prior_alerted < "$marker" || true + fi + case "$prior_since" in ''|*[!0-9]*) prior_since= ;; esac + if [ "$prior_offset" != "$CURSOR_OFFSET" ] || [ -z "$prior_since" ] || [ "$prior_since" -gt "$now" ]; then + prior_since=$now + prior_alerted=0 + rm -f -- "$ready" || return 1 + fi + since=$prior_since + if [ "$prior_alerted" = 1 ]; then return 0; fi + if [ $((now - since)) -lt "$threshold" ]; then + tmp=$(umask 077; mktemp "$CURSOR_DIR/.lag.XXXXXX") || return 1 + printf '%s %s 0\n' "$CURSOR_OFFSET" "$since" > "$tmp" && mv -f -- "$tmp" "$marker" || { rm -f -- "$tmp"; return 1; } + return 0 + fi + key="remote-reply-lag-$id-$since" + reason="check: remote reply channel stalled: mate=$id remote_bytes=$size cursor=$CURSOR_OFFSET for $((now - since))s; inspect the process-event listener and watcher, then re-ensure remote-reply-$id" + tmp=$(umask 077; mktemp "$CURSOR_DIR/.lag.XXXXXX") || return 1 + if ! fm_wake_append check "$key" "$reason"; then + rm -f -- "$tmp" + return 1 + fi + # The watcher probes in a background worker. Publish its delivery receipt + # before marking the episode alerted so a failed write retries next probe. + printf '%s\n%s %s\n' "$reason" "$CURSOR_OFFSET" "$since" > "$tmp" \ + && mv -f -- "$tmp" "$ready" || { rm -f -- "$tmp"; return 1; } + tmp=$(umask 077; mktemp "$CURSOR_DIR/.lag.XXXXXX") || return 1 + printf '%s %s 1\n' "$CURSOR_OFFSET" "$since" > "$tmp" && mv -f -- "$tmp" "$marker" || { rm -f -- "$tmp"; return 1; } + "$SCRIPT_DIR/fm-procevent.sh" ensure-listening "$(source_id "$id")" >/dev/null 2>&1 || true + printf '%s\n' "$reason" +} + +# Pick the acknowledged truncation and the last ingested delta that justify +# changing this exact cursor. Captured result bodies, including their raw +# payload bytes, remain durable after handling; their ingestion receipts bind +# those bytes to the committed cursor. A digest alone cannot prove a trimmed +# prefix, so missing capture evidence is a refusal. +rebase_find_evidence() { # + local id=$1 sid result base seq class path from to hash reason + sid=$(source_id "$id") + REBASE_DELTA= + REBASE_FROM= + REBASE_FROM_HASH= + REBASE_BREAK=0 + for result in "$STATE/procevent-inbox/$sid".*.result; do + [ -e "$result" ] || continue + [ -f "$result" ] && [ ! -L "$result" ] || die "unsafe remote reply result: $result" + base=${result%.result} + [ -f "$base.handled" ] && [ ! -L "$base.handled" ] || continue + seq=${base##*.} + case "$seq" in ''|*[!0-9]*) continue ;; esac + class=$(classify_result "$result") + path=$(result_field "$result" path 2>/dev/null || true) + [ "$path" = "$REMOTE_LOG" ] || continue + from=$(result_field "$result" from_offset 2>/dev/null || true) + hash=$(result_field "$result" from_prefix_sha256 2>/dev/null || true) + if [ "$class" = continuity-broken ]; then + reason=$(result_field "$result" reason 2>/dev/null || true) + [ "$reason" = truncated ] && [ "$from" = "$CURSOR_OFFSET" ] \ + && [ "$hash" = "$CURSOR_HASH" ] && REBASE_BREAK=1 + continue + fi + [ "$class" = delta ] || continue + to=$(result_field "$result" to_offset 2>/dev/null || true) + hash=$(result_field "$result" to_prefix_sha256 2>/dev/null || true) + [ "$to" = "$CURSOR_OFFSET" ] && [ "$hash" = "$CURSOR_HASH" ] || continue + ingest_receipt_matches "$id" "$seq" "$result" || continue + from=$(result_field "$result" from_offset) || die "last ingested delta has no start offset" + case "$from" in ''|*[!0-9]*) die "last ingested delta has an invalid start offset" ;; esac + if [ -z "$REBASE_FROM" ] || [ "$from" -gt "$REBASE_FROM" ]; then + REBASE_DELTA=$result + REBASE_FROM=$from + REBASE_FROM_HASH=$(result_field "$result" from_prefix_sha256) \ + || die "last ingested delta has no start hash" + fi + done + [ "$REBASE_BREAK" -eq 1 ] || die "no acknowledged truncation matches the current cursor" + [ -n "$REBASE_DELTA" ] || die "last ingested delta is unavailable; cannot prove the retained tail" +} + +cmd_rebase() { # --expect-offset + local id=${1:-} expected=${3:-} sid source cursor status_file remote_size + local probe old_payload retained_payload old_bytes retained_bytes retained_hash + local new_to new_hash verified_size complete_bytes schema blank rc=0 backup rollback tmp arm_out line + [ "$#" -eq 3 ] && [ "$2" = --expect-offset ] \ + || die "usage: fm-procevent-remote-reply.sh rebase --expect-offset " + validate_id "$id" + case "$expected" in ''|*[!0-9]*) die "expected offset must be numeric" ;; esac + sid=$(source_id "$id") + source="$STATE/procevent/$sid.source" + cursor=$(cursor_path "$id") + status_file="$STATE/$id.status" + ( + REBASE_LIFECYCLE_LOCK=$(secondmate_reply_lifecycle_lock_path "$STATE" "$id") + fm_lock_acquire_wait "$REBASE_LIFECYCLE_LOCK" || die "cannot lock remote reply lifecycle" + trap 'fm_lock_release "$REBASE_LIFECYCLE_LOCK"' EXIT + remote_route_exists "$id" + [ -d "$CURSOR_DIR" ] && [ ! -L "$CURSOR_DIR" ] \ + || die "remote reply cursor directory is unavailable or unsafe" + [ ! -e "$source" ] && [ ! -L "$source" ] \ + || die "reply source is still registered; rebase requires a terminal continuity break" + [ -f "$cursor" ] && [ ! -L "$cursor" ] || die "reply cursor is unavailable or unsafe" + [ -f "$status_file" ] && [ ! -L "$status_file" ] \ + || die "reply status is unavailable or unsafe" + read_cursor "$id" + [ "$CURSOR_OFFSET" = "$expected" ] || die "reply cursor changed from the expected offset" + [ "${#CURSOR_OFFSET}" -le 18 ] || die "reply cursor exceeds the rebase size bound" + retirement_capture_scan "$id" || die "cannot inspect captured reply results" + [ "$RETIREMENT_PENDING" -eq 0 ] \ + || die "rebase refused with unhandled captured reply results" + rebase_find_evidence "$id" + remote_size=$(fm_run_timed 15 "$SCRIPT_DIR/fm-on.sh" "$id" \ + fm-remote-delta-read.sh size "$REMOTE_LOG" 2>/dev/null) \ + || die "cannot read remote reply size for rebase" + case "$remote_size" in ''|*[!0-9]*) die "remote reply size is invalid" ;; esac + [ "${#remote_size}" -le 18 ] && [ "$remote_size" -lt "$CURSOR_OFFSET" ] \ + && [ "$remote_size" -ge "$REBASE_FROM" ] \ + || die "remote log is not a verifiable suffix trim of the ingested delta" + tmp=$(umask 077; mktemp -d "$CURSOR_DIR/.rebase.XXXXXX") \ + || die "cannot stage rebase proof" + REBASE_TMP=$tmp + trap 'rm -rf -- "$REBASE_TMP"; fm_lock_release "$REBASE_LIFECYCLE_LOCK"' EXIT + old_bytes=$(result_field "$REBASE_DELTA" payload_bytes) \ + || die "ingested payload size is ambiguous" + case "$old_bytes" in ''|*[!0-9]*) die "ingested payload size is invalid" ;; esac + retained_bytes=$((remote_size - REBASE_FROM)) + [ "$retained_bytes" -le "$old_bytes" ] \ + || die "remote tail extends beyond the last ingested delta" + blank=$(LC_ALL=C awk '$0 == "" { print NR; exit }' "$REBASE_DELTA") + case "$blank" in ''|*[!0-9]*) die "ingested delta has no payload boundary" ;; esac + old_payload="$tmp/ingested-payload" + tail -n "+$((blank + 1))" "$REBASE_DELTA" > "$old_payload" \ + || die "cannot read ingested tail" + [ "$(LC_ALL=C wc -c < "$old_payload" | tr -d ' ')" = "$old_bytes" ] \ + || die "ingested payload size changed" + retained_payload="$tmp/retained-payload" + head -c "$retained_bytes" "$old_payload" > "$retained_payload" \ + || die "cannot stage retained ingested tail" + [ "$(LC_ALL=C wc -c < "$retained_payload" | tr -d ' ')" = "$retained_bytes" ] \ + || die "retained ingested tail is incomplete" + retained_hash=$(sha256_file "$retained_payload") \ + || die "cannot hash retained ingested tail" + probe="$tmp/remote-result" + fm_run_timed 30 "$SCRIPT_DIR/fm-on.sh" "$id" fm-remote-delta-read.sh \ + verify-rebase "$REMOTE_LOG" "$REBASE_FROM" "$REBASE_FROM_HASH" \ + "$retained_bytes" "$retained_hash" \ + > "$probe" 2> "$tmp/remote-error" || rc=$? + [ "$rc" -eq 0 ] \ + || die "retained remote tail differs or could not be verified (exit $rc): $(head -n 1 "$tmp/remote-error")" + schema=$(result_field "$probe" schema) || die "remote rebase proof has no schema" + verified_size=$(result_field "$probe" remote_bytes) \ + || die "remote rebase proof has no source size" + new_to=$(result_field "$probe" offset) || die "remote rebase proof has no cursor" + new_hash=$(result_field "$probe" prefix_sha256) \ + || die "remote rebase proof has no cursor hash" + complete_bytes=$(LC_ALL=C od -An -v -tu1 "$retained_payload" | awk ' + { for (i = 1; i <= NF; i++) { bytes++; if ($i == 10) complete=bytes } } + END { print complete + 0 } + ') + [ "$schema" = fm-remote-rebase-verification.v1 ] \ + && [ "$verified_size" = "$remote_size" ] \ + && [ "$new_to" = "$((REBASE_FROM + complete_bytes))" ] \ + && [ "$new_to" -lt "$CURSOR_OFFSET" ] \ + || die "remote rebase proof does not match the retained ingested tail" + case "$new_hash" in *[!A-Fa-f0-9]*|'') die "remote rebase proof has an invalid cursor hash" ;; esac + [ "${#new_hash}" -eq 64 ] || die "remote rebase proof has an invalid cursor hash" + backup=$(umask 077; mktemp "$CURSOR_DIR/$id.rebase-prior.XXXXXX") \ + || die "cannot preserve prior cursor" + cp -P "$cursor" "$backup" || die "cannot preserve prior cursor" + rm -f -- "$(fm_pending_reply_remote_channel_watermark_path "$STATE" "$id")" \ + || die "cannot clear stale remote-channel caught-up evidence" + write_cursor "$id" "$new_to" "$new_hash" || die "cannot commit verified rebase cursor" + arm_out="$tmp/arm.out" + if ! cmd_arm_locked "$id" > "$arm_out"; then + if [ ! -e "$source" ] && [ ! -L "$source" ]; then + rollback=$(umask 077; mktemp "$CURSOR_DIR/.cursor.rollback.XXXXXX") \ + || die "rebase registration failed and prior cursor could not be restored: $backup" + cp -P "$backup" "$rollback" && mv -f -- "$rollback" "$cursor" \ + || die "rebase registration failed and prior cursor could not be restored: $backup" + die "rebase registration failed; prior cursor restored from $backup" + fi + die "registration returned failure after a source appeared; rebased cursor retained and prior cursor backed up at $backup" + fi + "$SCRIPT_DIR/fm-procevent.sh" ensure-listening "$sid" >/dev/null \ + || die "rebased cursor and registered source, but listener is not confirmed; run bin/fm-procevent.sh ensure-listening $sid" + line="resolved [key=remote-reply-continuity-$id] [at=$(date +%s)]: verified retained remote reply tail and re-armed at offset $new_to" + append_status_once "$status_file" "$line" \ + || die "rebase succeeded but could not publish its resolution" + fm_lock_release "$REBASE_LIFECYCLE_LOCK" || die "cannot release remote reply lifecycle lock" + trap 'rm -rf -- "$REBASE_TMP"' EXIT + # The test seam holds this command after the launch while a second break is + # captured, pinning the durable order of resolution and renewed blockage. + if [ "${FM_TEST_SEAM:-0}" = 1 ] && [ -n "${FM_TEST_REBASE_AFTER_ENSURE_HOOK:-}" ]; then + "$FM_TEST_REBASE_AFTER_ENSURE_HOOK" "$sid" \ + || die "rebase test hook failed" + fi + printf 'rebased: %s offset=%s prior-cursor=%s\n' "$id" "$new_to" "$backup" + ) +} + safe_doc_path() { case "$1" in data/*.md) ;; @@ -544,7 +828,7 @@ cmd_ingest() { die "result does not continue the current cursor for $id" fi if [ "$class" = continuity-broken ]; then - line="blocked [key=remote-reply-continuity-$id]: remote reply continuity broke for $id ($reason)" + line="blocked [key=remote-reply-continuity-$id]: remote reply continuity broke for $id ($reason; cursor=$from remote_bytes=$to generation=${seq:-manual})" append_rc=0 if status_event_recorded "$status_file" "$line"; then append_rc=1 @@ -744,6 +1028,8 @@ cmd_retire_finalize_locked() { fi rm -f -- "$(cursor_path "$id")" rm -f -- "$CURSOR_DIR/$id".*.ingested + rm -f -- "$CURSOR_DIR/$id.lag" "$CURSOR_DIR/$id.lag-ready" \ + "$CURSOR_DIR/$id.lag-probe" "$CURSOR_DIR/$id.source-failed" rm -f -- "$(fm_pending_reply_remote_channel_watermark_path "$STATE" "$id")" } @@ -780,10 +1066,12 @@ case "${1:-}" in autohandle) shift; [ "$#" -eq 3 ] || usage; cmd_autohandle "$@" ;; ingest) shift; [ "$#" -eq 2 ] || usage; cmd_ingest "$@" ;; classify) shift; [ "$#" -eq 1 ] || usage; classify_result "$1" ;; - terminal) shift; [ "$#" -eq 1 ] || usage; [ -s "$1" ] ;; + terminal) shift; [ "$#" -eq 1 ] || usage; [ "$(classify_result "$1")" = continuity-broken ] ;; self-announcing) shift; [ "$#" -eq 0 ] || usage; exit 0 ;; source-id) shift; [ "$#" -eq 1 ] || usage; source_id "$1" ;; relisten) shift; [ "$#" -eq 0 ] || usage; exit 0 ;; + lag-check) shift; [ "$#" -eq 1 ] || usage; cmd_lag_check "$@" ;; + rebase) shift; cmd_rebase "$@" ;; retire) shift; [ "$#" -ge 1 ] && [ "$#" -le 2 ] || usage; cmd_retire "$@" ;; retire-quiesce-locked) shift; [ "$#" -ge 1 ] && [ "$#" -le 2 ] || usage; require_parent_lifecycle_lock "$1"; cmd_retire_quiesce_locked "$@" ;; retire-finalize-locked) shift; [ "$#" -ge 1 ] && [ "$#" -le 2 ] || usage; require_parent_lifecycle_lock "$1"; cmd_retire_finalize_locked "$@" ;; diff --git a/bin/fm-remote-delta-read.sh b/bin/fm-remote-delta-read.sh index 84aef13d051..8422ec38ba3 100755 --- a/bin/fm-remote-delta-read.sh +++ b/bin/fm-remote-delta-read.sh @@ -3,12 +3,16 @@ # # Usage: # fm-remote-delta-read.sh [wait-seconds] +# fm-remote-delta-read.sh size +# fm-remote-delta-read.sh verify-rebase # -# The reader validates continuity by hashing the exact prefix represented by the -# caller's cursor. It then blocks until at least one complete appended line is -# available, returns at most 65536 payload bytes, and never truncates or consumes -# the source. A shortened or changed prefix returns a structured continuity-break -# result instead of silently rebasing the cursor. +# Normal delta mode validates continuity by hashing the exact prefix represented +# by the caller's cursor. It then blocks until at least one complete appended +# line is available, returns at most 65536 payload bytes, and never truncates or +# consumes the source. A shortened or changed prefix returns a structured +# continuity-break result instead of silently rebasing the cursor. +# verify-rebase checks every retained byte of one bounded snapshot against the +# already ingested tail and returns its last complete-line cursor. # # The log is sampled every FM_REMOTE_DELTA_POLL_SECONDS (default 0.5 seconds). # A complete line is visible on the next sample, and the window deadline can @@ -143,6 +147,85 @@ emit_break() { # printf 'reason=%s\n\n' "$1" } +# Verify one snapshot against the last ingested delta, including an incomplete +# final line. Return only the last complete-line cursor. A size change between +# the caller's size probe and this snapshot is a refusal, not a partial proof. +if [ "${1:-}" = verify-rebase ]; then + [ "$#" -eq 6 ] || usage + REL=$2 + OFFSET=$3 + PREFIX=$4 + RETAINED=$5 + EXPECTED_TAIL_HASH=$6 + for NUMBER in "$OFFSET" "$RETAINED"; do + case "$NUMBER" in ''|*[!0-9]*) die "rebase offsets must be nonnegative integers" ;; esac + done + [ "${#OFFSET}" -le 18 ] && [ "${#RETAINED}" -le 7 ] \ + || die "rebase snapshot exceeds its size bound" + [ "$RETAINED" -le 1048576 ] || die "rebase tail exceeds its size bound" + for HASH in "$PREFIX" "$EXPECTED_TAIL_HASH"; do + case "$HASH" in *[!A-Fa-f0-9]*|'') die "rebase hash must be hexadecimal" ;; esac + [ "${#HASH}" -eq 64 ] || die "rebase hash has the wrong length" + done + PREFIX=$(printf '%s' "$PREFIX" | tr 'A-F' 'a-f') + EXPECTED_TAIL_HASH=$(printf '%s' "$EXPECTED_TAIL_HASH" | tr 'A-F' 'a-f') + LOG=$(resolve_log "$REL") + [ -e "$LOG" ] || die "rebase source log is missing" + TMP=$(mktemp -d "${TMPDIR:-/tmp}/fm-remote-rebase.XXXXXX") \ + || die "cannot stage remote rebase proof" + trap 'rm -rf -- "$TMP"' EXIT + MAX_BYTES=$RETAINED + snapshot_log "$LOG" "$TMP/source" "$TMP/size" \ + || die "rebase source could not be captured safely" + IFS= read -r SIZE < "$TMP/size" + [ "$SIZE" -eq $((OFFSET + RETAINED)) ] \ + || die "rebase source size changed before verification" + copy_prefix "$TMP/source" "$OFFSET" "$TMP/prefix" + [ "$(sha256_file "$TMP/prefix")" = "$PREFIX" ] \ + || die "rebase source prefix differs from the ingested cursor" + tail -c "+$((OFFSET + 1))" "$TMP/source" > "$TMP/tail" + [ "$(LC_ALL=C wc -c < "$TMP/tail" | tr -d ' ')" = "$RETAINED" ] \ + || die "rebase source tail size changed" + [ "$(sha256_file "$TMP/tail")" = "$EXPECTED_TAIL_HASH" ] \ + || die "rebase source tail differs from the ingested bytes" + COMPLETE_BYTES=$(LC_ALL=C od -An -v -tu1 "$TMP/tail" | awk ' + { for (i = 1; i <= NF; i++) { bytes++; if ($i == 10) complete=bytes } } + END { print complete + 0 } + ') + TO=$((OFFSET + COMPLETE_BYTES)) + copy_prefix "$TMP/source" "$TO" "$TMP/to-prefix" + printf 'schema=fm-remote-rebase-verification.v1\n' + printf 'remote_bytes=%s\n' "$SIZE" + printf 'offset=%s\n' "$TO" + printf 'prefix_sha256=%s\n' "$(sha256_file "$TMP/to-prefix")" + exit 0 +fi + +if [ "${1:-}" = size ]; then + [ "$#" -eq 2 ] || usage + LOG=$(resolve_log "$2") + if [ ! -e "$LOG" ]; then + printf '0\n' + exit 0 + fi + LOG_PARENT=$(dirname "$LOG") + LOG_BASE=$(basename "$LOG") + ( + CDPATH='' cd -- "$LOG_PARENT" 2>/dev/null || exit 1 + [ "$(pwd -P)" = "$LOG_PARENT" ] || exit 1 + perl -MFcntl=:DEFAULT -e ' + use strict; + use warnings; + my $path = shift; + sysopen(my $in, $path, O_RDONLY | O_NOFOLLOW) or exit 1; + my @stat = stat $in or exit 1; + -f _ or exit 1; + print "$stat[7]\n"; + ' "$LOG_BASE" + ) || die "log size could not be read safely: $2" + exit 0 +fi + [ "$#" -ge 3 ] && [ "$#" -le 4 ] || usage REL=$1 OFFSET=$2 diff --git a/bin/fm-session-lock-lib.sh b/bin/fm-session-lock-lib.sh index 58fc088fd59..fdca76d8cca 100644 --- a/bin/fm-session-lock-lib.sh +++ b/bin/fm-session-lock-lib.sh @@ -8,9 +8,9 @@ # Stop hook fires inside the lock-owning primary session before it may arm or # rewake. Two signals decide ownership, either one sufficient: the recorded pid # is a member of this process's contiguous harness ancestry, or the trusted -# Claude session id below matches the id recorded beside a live lock. Neither -# signal ever fails open: no id, no sidecar, an untrusted id, or a different -# recorded id leaves the ancestry verdict exactly as it was. +# Claude or Codex thread id below matches the id recorded beside a live lock. +# A managed Codex app-server is shared across threads, so its ancestry alone +# never proves ownership. Other harnesses retain the ancestry verdict. # This file is sourced by scripts and has no side effects on source. # Cursor process identity is NOT expressible as a command-name pattern and is @@ -172,6 +172,19 @@ fm_harness_pid_alive() { fm_harness_process_matches "$comm" "$args" } +# A Codex Desktop tool shell can descend from the shared managed app-server. +# That process survives individual threads and is never sufficient ownership +# evidence by itself. Other Codex processes retain the ordinary PID rule. +fm_session_lock_shared_codex_pid() { # + local pid=$1 comm args + comm=$(ps -o comm= -p "$pid" 2>/dev/null) || return 1 + args=$(ps -o args= -p "$pid" 2>/dev/null) || return 1 + case "${comm##*/}:$args" in + codex:*app-server*--managed-daemon*) return 0 ;; + esac + return 1 +} + # --- trusted same-session identity ------------------------------------------- # Claude Code hands every hook and tool shell CLAUDE_CODE_SESSION_ID (the # session's conversation id) and CLAUDE_PID (the pid of the process running the @@ -195,25 +208,40 @@ fm_harness_pid_alive() { # non-goal. Two genuinely different live sessions sharing one id is not a # supported state (Claude refuses to resume a running session under its id). -# Print the Claude session id this process may own with, or return 1. $1 is the +# Print the verified session or thread id this process may own with, or return 1. $1 is the # ancestry list an earlier walk already produced, so a caller that walked once # need not walk again. fm_session_lock_trusted_session_id() { # [] - local id=${CLAUDE_CODE_SESSION_ID:-} claude_pid=${CLAUDE_PID:-} pids=${1:-} pid comm args - [ -n "$id" ] || return 1 - case "$id" in *$'\n'*|*$'\r'*) return 1 ;; esac - case "$claude_pid" in ''|*[!0-9]*) return 1 ;; esac + local id=${CLAUDE_CODE_SESSION_ID:-} claude_pid=${CLAUDE_PID:-} codex_id=${CODEX_THREAD_ID:-} + local pids=${1:-} pid comm args if [ -z "$pids" ]; then pids=$(fm_harness_ancestry_pids) || return 1 fi + if [ -n "$id" ] && [ -n "$claude_pid" ]; then + case "$id" in *$'\n'*|*$'\r'*) id= ;; esac + case "$claude_pid" in *[!0-9]*) id= ;; esac + fi + if [ -n "$id" ] && [ -n "$claude_pid" ]; then + while IFS= read -r pid; do + [ "$pid" = "$claude_pid" ] || continue + comm=$(ps -o comm= -p "$pid" 2>/dev/null) || return 1 + args=$(ps -o args= -p "$pid" 2>/dev/null) + fm_harness_process_matches "$comm" "$args" || return 1 + [ "$FM_HARNESS_IS_CLAUDE" -eq 1 ] || return 1 + printf '%s\n' "$id" + return 0 + done </dev/null) || return 1 args=$(ps -o args= -p "$pid" 2>/dev/null) - fm_harness_process_matches "$comm" "$args" || return 1 - [ "$FM_HARNESS_IS_CLAUDE" -eq 1 ] || return 1 - printf '%s\n' "$id" - return 0 + case "${comm##*/}:$args" in + codex:*) printf 'codex:%s\n' "$codex_id"; return 0 ;; + esac done < [] # the sidecar still names that session. Every other session records the # outermost pid of its contiguous run, exactly as before. fm_session_lock_anchor_pid() { - local pids + local pids trusted pids=$(fm_harness_ancestry_pids) || return 1 - if fm_session_lock_trusted_session_id "$pids" >/dev/null; then - printf '%s\n' "$CLAUDE_PID" - return 0 - fi + trusted=$(fm_session_lock_trusted_session_id "$pids") || trusted= + case "$trusted" in + ''|codex:*) ;; + *) printf '%s\n' "$CLAUDE_PID"; return 0 ;; + esac _fm_harness_outermost_pid "$pids" } @@ -280,6 +309,11 @@ fm_session_lock_owned_by_self() { ''|*[!0-9]*) return 1 ;; esac pids=$(fm_harness_ancestry_pids) || return 1 + if fm_session_lock_shared_codex_pid "$lock_pid"; then + fm_session_lock_same_session "$state" "$pids" || return 1 + fm_harness_pid_alive "$lock_pid" + return $? + fi while IFS= read -r pid; do [ "$pid" = "$lock_pid" ] && return 0 done </dev/null || true) @@ -305,13 +339,8 @@ fm_session_lock_foreign_owner_live() { ''|*[!0-9]*) return 1 ;; esac fm_harness_pid_alive "$lock_pid" || return 1 - pids=$(fm_harness_ancestry_pids) || return 1 - while IFS= read -r pid; do - [ "$pid" = "$lock_pid" ] && return 1 - done </dev/null || return 1 + fm_session_lock_owned_by_self "$state" && return 1 # shellcheck disable=SC2034 # Output global, read by the sourcing guard caller. FM_SESSION_LOCK_FOREIGN_OWNER_PID=$lock_pid return 0 diff --git a/bin/fm-watch.sh b/bin/fm-watch.sh index 6e51f76770a..b671a0ec90b 100755 --- a/bin/fm-watch.sh +++ b/bin/fm-watch.sh @@ -1135,6 +1135,72 @@ secondmate_liveness_tick() { [ "$failed" -eq 0 ] } +# The remote reply mirror has its own progress evidence: the remote log size +# and the committed local cursor. The adapter owns the episode marker, its one +# ensure-listening repair, and wake publication. Remote size reads can take 15s +# each, so run them outside the watcher's signal and inactive-outcome path. +remote_reply_lag_tick() { + local meta id kind host + for meta in "$STATE"/*.meta; do + [ -e "$meta" ] || continue + kind=$(fm_meta_get "$meta" kind) + [ "$kind" = secondmate ] || continue + host=$(fm_meta_get "$meta" remote_host) + [ -n "$host" ] || continue + id=${meta##*/} + id=${id%.meta} + case "$id" in ''|*[!A-Za-z0-9._-]*) continue ;; esac + FM_HOME="$FM_HOME" FM_STATE_OVERRIDE="$STATE" \ + "$SCRIPT_DIR/fm-procevent-remote-reply.sh" lag-check "$id" >/dev/null 2>&1 || true + done +} + +remote_reply_lag_start() { + ( + trap - EXIT HUP INT TERM + mkdir -p "$STATE/remote-replies" || exit 1 + [ -d "$STATE/remote-replies" ] && [ ! -L "$STATE/remote-replies" ] || exit 1 + fm_lock_try_acquire "$STATE/remote-replies/.lag-check.lock" || exit 0 + trap 'fm_lock_release "$STATE/remote-replies/.lag-check.lock" || true' EXIT + remote_reply_lag_tick + ) /dev/null 2>&1 & +} + +remote_reply_lag_after_output() { + [ "$1" -eq 0 ] || return 0 + rm -f -- "$REMOTE_REPLY_LAG_READY" +} + +remote_reply_lag_surface() { + local ready reason id marker cursor marker_offset marker_since marker_alerted + local ready_offset ready_since cursor_offset + for ready in "$STATE/remote-replies"/*.lag-ready; do + [ -e "$ready" ] || continue + [ -f "$ready" ] && [ ! -L "$ready" ] || continue + id=${ready##*/} + id=${id%.lag-ready} + case "$id" in ''|*[!A-Za-z0-9._-]*) continue ;; esac + marker="$STATE/remote-replies/$id.lag" + [ -f "$marker" ] && [ ! -L "$marker" ] || continue + read -r marker_offset marker_since marker_alerted < "$marker" || continue + [ "$marker_alerted" = 1 ] || continue + read -r ready_offset ready_since < <(sed -n '2p' "$ready") || continue + [ "$ready_offset" = "$marker_offset" ] && [ "$ready_since" = "$marker_since" ] || continue + cursor="$STATE/remote-replies/$id.cursor" + cursor_offset=0 + if [ -e "$cursor" ] || [ -L "$cursor" ]; then + [ -f "$cursor" ] && [ ! -L "$cursor" ] || continue + cursor_offset=$(sed -n 's/^offset=//p' "$cursor") + fi + [ "$cursor_offset" = "$ready_offset" ] || continue + IFS= read -r reason < "$ready" || continue + case "$reason" in 'check: remote reply channel stalled: mate='*) ;; *) continue ;; esac + REMOTE_REPLY_LAG_READY=$ready + FM_WAKE_POST_OUTPUT_ACTION=remote_reply_lag_after_output + wake "$reason" + done +} + # Consecutive wedge-escalation count for a window past FM_WEDGE_DEMAND_INSPECT_COUNT # (default 3): a pane that keeps re-wedging on the SAME stale hash - each # escalation gets absorbed again as "still validating" one poll later, since the @@ -2665,6 +2731,11 @@ while :; do # alive. Supervision scripts warn when this goes stale with tasks in flight. touch "$STATE/.last-watcher-beat" + # A completed background probe is already queued; surface it before other + # potentially slow per-cycle work, then start the next bounded probe. + remote_reply_lag_surface + remote_reply_lag_start + # Opt-in fleet activity ledger (docs/fleet-ledger.md): pick up newly appended # status lines before this cycle can exit on a wake. Off costs one file test. [ ! -e "$CONFIG/fleet-ledger" ] || FM_HOME=$FM_HOME FM_STATE_OVERRIDE=$STATE FM_CONFIG_OVERRIDE=$CONFIG "$SCRIPT_DIR/fm-fleet-ledger.sh" capture || true diff --git a/docs/remote-secondmates.md b/docs/remote-secondmates.md index e13eebacf9f..c551fe5eacc 100644 --- a/docs/remote-secondmates.md +++ b/docs/remote-secondmates.md @@ -568,6 +568,9 @@ So a mirrored reply reaches the primary status channel without depending on the A mirrored line that carries a correlation token settles its pending-reply record and closes that request's own open escalation decision. A remote reply reaches the primary only through this asynchronous mirror. +The listener retains its process-event owner across transient read failures and retries three times; when all three fail it queues one actionable failure for a continuing episode and exits, leaving recovery to the reconcile launch floor. +The primary watcher schedules a detached probe of the remote log size against the committed local cursor on a bounded cadence; when the remote log stays ahead beyond the configured threshold the adapter runs `fm-procevent.sh ensure-listening` for that source once in the episode and queues one stalled-channel wake. +Cursor progress or catchup resets that lag episode. Because of that, the primary treats a missing correlated report as a missed report only once the mirror has been read through the end of the remote log after that turn ended. A remote mate that did answer is therefore never asked to repost while its answer is still in flight. A genuinely missing answer still gets exactly one repost once the mirror is known to be current. @@ -576,8 +579,10 @@ The [process-to-event operating contract](configuration.md#process-to-event-sour ### Source log continuity -The source log is never truncated or consumed. -A shortened or changed prefix stops the relay and surfaces a continuity failure instead of silently resetting the cursor. +The relay never truncates or consumes the source log. +A source shortened or changed by another writer stops the relay and surfaces a continuity failure instead of silently resetting the cursor. +For a shortened source, `bin/fm-procevent-remote-reply.sh rebase` requires an expected cursor and compares every retained remote byte with its acknowledged ingested bytes before rewinding to a complete-line boundary and re-arming; its script header owns the exact command and refusal contract. +Changed bytes or missing ingestion evidence remain blocked for investigation. ### SSH exit 255 and unavailable homes diff --git a/docs/scripts.md b/docs/scripts.md index 8d656ec203e..3914ee2a307 100644 --- a/docs/scripts.md +++ b/docs/scripts.md @@ -48,7 +48,7 @@ The shared no-mistakes gate lifecycle boundary is summarized in [architecture.md | `fm-ensure-agents-md.sh` | Manually initialize project agent-memory files (see the helper's header and help) | | `fm-guard.sh` | Warn on primary-checkout tangles, main-session pending wakes, and unhealthy supervision | | `fm-primary-scope-lib.sh` | Shared marker-or-plain-checkout primary-home predicate for tracked hooks | -| `fm-session-lock-lib.sh` | Shared session-lock ownership from harness ancestry or a trusted Claude session id for fm-lock.sh and the Claude Stop auto-arm, plus the read-only lock inspection behind `fm-lock.sh status` and `fm-inbox.sh ready` | +| `fm-session-lock-lib.sh` | Shared session-lock ownership from harness ancestry or a trusted Claude session or Codex thread id, with shared Codex daemons requiring the thread id; also owns the read-only lock inspection behind `fm-lock.sh status` and `fm-inbox.sh ready` | | `fm-claude-stop-autoarm.sh` | Claude Stop `asyncRewake` hook owning tokenless watcher continuity with single-flight exit-2 rewake (docs/watcher-continuity.md) | | `fm-turnend-guard.sh` | Shared primary turn-end guard predicate so no turn ends blind (docs/turnend-guard.md) | | `fm-turnend-guard-grok.sh` | Grok Stop-hook adapter for the primary turn-end guard | diff --git a/docs/sessionstart-nudge.md b/docs/sessionstart-nudge.md index 2b111302985..6c15f065e97 100644 --- a/docs/sessionstart-nudge.md +++ b/docs/sessionstart-nudge.md @@ -96,11 +96,13 @@ So `clear` or `compact` cannot skip startup sweeps after a truncated run. `bin/fm-lock.sh` treats a lock as this session's own when it is owned through either of these: -- The shared ancestry verdict. -- A trusted same-session Claude id. +- The shared ancestry verdict, except when the owner is a managed Codex app-server shared across threads. +- A trusted same-session Claude id or Codex thread id recorded beside that shared daemon's lock. So a proven `clear` or `compact` re-emit re-verifies ownership and proceeds. A lock another live session took meanwhile still produces the ordinary read-only digest. +An idle managed Codex daemon remains live after its owning Desktop thread ends; `bin/fm-lock.sh take-over` provides a guarded handoff when its exact recorded pid and thread identity are known, its watcher is gone, and its lock and beacon have been quiet long enough. +The operator must also verify that the old thread has ended and explicitly attest to that exact pid and session; quiet lock evidence alone cannot rule out an active thread, and an incorrect takeover can split ownership. ### Nudge wrapper on a run-tier harness diff --git a/docs/turnend-guard.md b/docs/turnend-guard.md index 23af69a0eb1..0bc56651051 100644 --- a/docs/turnend-guard.md +++ b/docs/turnend-guard.md @@ -101,8 +101,8 @@ When an active home instead has a live session lock held by a verified harness t Ownership is the shared `fm_session_lock_owned_by_self` verdict in `bin/fm-session-lock-lib.sh`. The current session owns the lock when either of these holds: -- The recorded pid is a member of the current session's contiguous harness ancestry. -- The trusted Claude session id recorded beside the lock in `state/.lock-session` matches this hook's own environment while the recorded pid is still a live harness. +- The recorded pid is a member of the current session's contiguous harness ancestry, unless it is a managed Codex app-server shared across threads. +- The trusted Claude session id or Codex thread id recorded beside the lock in `state/.lock-session` matches this hook's own environment while the recorded pid is still a live harness. That second signal keeps a background Claude session owning its own lock after the transient helper chain between its hooks and its recorded owner is recycled. The library's header owns the trust gate (`CLAUDE_PID` must be a Claude-shaped member of the current run). diff --git a/docs/watcher-continuity.md b/docs/watcher-continuity.md index f14717596b4..acd3d3dd3b0 100644 --- a/docs/watcher-continuity.md +++ b/docs/watcher-continuity.md @@ -98,8 +98,8 @@ The hook handles the session lock as follows: Whether the session owns that lock is the shared `fm_session_lock_owned_by_self` verdict in `bin/fm-session-lock-lib.sh`. That verdict accepts either of two cases: -- A recorded pid inside the current harness ancestry. -- A live lock recorded under this same trusted Claude session id. +- A recorded pid inside the current harness ancestry, unless it is a managed Codex app-server shared across threads. +- A live lock recorded under this same trusted Claude session id or Codex thread id. With that verdict, a background session keeps arming after its transient helper chain is recycled. [`turnend-guard.md`](turnend-guard.md#guard-predicates) owns the Claude guard's behavior when that live owner is genuinely another session. @@ -203,6 +203,18 @@ It enters its poll loop immediately and keeps scanning signals, stale panes, and - Codex retains its bounded foreground checkpoint protocol. - Grok retains its tracked background-task notification protocol. +A Codex Desktop primary that was live during the upgrade to sidecar-recorded Codex thread identity holds a shared-daemon lock with no `state/.lock-session`, so its own thread now reads that lock as foreign. +The lock and watcher beacon being quiet does not prove the old Codex thread ended. +The operator must verify that the old thread is closed; taking over an active thread can split one home between two primaries. +There is no automatic adoption; recover it explicitly: + +1. Stop that thread's watcher: `bin/fm-watch-arm.sh --stop`. +2. Wait out the quiet window (`FM_LOCK_TAKEOVER_QUIET_SECONDS`, default 300 seconds) so neither `state/.lock` nor `state/.last-watcher-beat` was touched within it. +3. From the Codex thread taking the lock, run `bin/fm-lock.sh take-over --expect-pid PID --expect-session none --attest-owner-ended PID/none`, with `PID` from `bin/fm-lock.sh status`. + +For a sidecar-recorded old thread, use `--expect-session codex:ID --attest-owner-ended PID/codex:ID` instead. +The command preserves a copy of the prior lock and sidecar in `state/.lock-takeover.*` before replacing them. + No adapter starts a replacement with a fire-and-forget shell `&` from a model command. The Claude hook's detached handling successor is launched by the hook itself, which waits for the successor's status line before it exits. diff --git a/tests/fm-remote-delta-read.test.sh b/tests/fm-remote-delta-read.test.sh index 49c37fe920a..a9920beadde 100755 --- a/tests/fm-remote-delta-read.test.sh +++ b/tests/fm-remote-delta-read.test.sh @@ -34,6 +34,16 @@ run_reader() { # [rel] "$READER" "${4:-$DELTA_LOG_REL}" "$1" "$2" "$3" } +printf 'sized\n' > "$DELTA_HOME/$DELTA_LOG_REL" +[ "$(FM_HOME="$DELTA_HOME" "$READER" size "$DELTA_LOG_REL")" = 6 ] \ + || fail 'the size probe did not report exact remote bytes' +FM_HOME="$DELTA_HOME" "$READER" size '../outside' >/dev/null 2>&1 \ + && fail 'the size probe accepted traversal' +ln -s replies.status "$DELTA_HOME/state/size-link.status" +FM_HOME="$DELTA_HOME" "$READER" size state/size-link.status >/dev/null 2>&1 \ + && fail 'the size probe accepted a symlink' +pass 'the remote size probe reports bytes and refuses unsafe paths' + # A growing log returns the complete appended lines with exact boundaries. : > "$DELTA_HOME/$DELTA_LOG_REL" run_reader 0 "$EMPTY_SHA" 4 > "$TMP_ROOT/growth.out" & diff --git a/tests/fm-remote-reply.test.sh b/tests/fm-remote-reply.test.sh index dfdc212c1b9..b9bb5fa2f9b 100755 --- a/tests/fm-remote-reply.test.sh +++ b/tests/fm-remote-reply.test.sh @@ -53,6 +53,13 @@ if [ -n "${FM_REMOTE_REPLY_POLL_LOG:-}" ]; then printf 'x\n' >> "$FM_REMOTE_REPLY_POLL_LOG" fi [ "${FM_REMOTE_REPLY_FAIL_READ:-}" != 1 ] || exit 255 +[ ! -e "${FM_REMOTE_REPLY_FAIL_FLAG:-/nonexistent}" ] || exit 255 +if [ -e "${FM_REMOTE_REPLY_FAIL_DELTA_FLAG:-/nonexistent}" ]; then + if ! printf '%s' "${@: -1}" | base64 -d 2>/dev/null | tr '\0' '\n' | grep -qx size; then + printf 'x\n' >> "$FM_REMOTE_REPLY_FAIL_DELTA_FLAG" + exit 255 + fi +fi host=$1 entry=$2 shift 2 @@ -173,6 +180,9 @@ if [ -z "$RESULT" ]; then fail "the remote reply delta was not durably captured" fi assert_grep 'done [corr=0123456789abcdef]' "$RESULT" "captured delta lost the correlated status line" +if remote_env "$ADAPTER" terminal "$RESULT"; then + fail "an ordinary delta was classified terminal and tried to retire its re-armed listener" +fi # One remote note, one announcement: the adapter declares self-announcing, so a # fully autohandled capture publishes NO check wake - the mirrored status bytes # are the single announcement, observed here through the same signature-vs-seen @@ -823,18 +833,38 @@ fi stop_reply_listener || fail "the continuity listener did not stop" pass "a remote reply listener stays owned across empty waits and a delta" -# A failed transport is not an empty wait: do not launch a second read under -# the same owner, even when the launch floor is short. +# A failed transport retries briefly under the same owner, then publishes one +# durable failure and exits so reconcile's launch floor owns recovery. : > "$TMP_ROOT/failed-polls" -FM_REMOTE_REPLY_FAIL_READ=1 FM_REMOTE_REPLY_POLL_LOG="$TMP_ROOT/failed-polls" \ +touch "$TMP_ROOT/fail-remote-read" +FM_REMOTE_REPLY_FAIL_FLAG="$TMP_ROOT/fail-remote-read" \ + FM_REMOTE_REPLY_POLL_LOG="$TMP_ROOT/failed-polls" \ FM_PROCEVENT_LAUNCH_FLOOR_SECONDS=1 \ remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" >/dev/null 2>&1 & failed_reader=$! wait "$failed_reader" || fail "failed reader did not leave the runner" sleep 2 -[ "$(wc -l < "$TMP_ROOT/failed-polls" | tr -d ' ')" -eq 1 ] \ - || fail "failed reader relaunched within the launch floor" -pass "a failed remote read exits instead of relistening" +[ "$(wc -l < "$TMP_ROOT/failed-polls" | tr -d ' ')" -eq 3 ] \ + || fail "the listener did not stop after its three-read retry budget" +assert_grep 'remote reply listener ios failed three consecutive reads' "$PARENT/state/.wake-queue" \ + "three failed reads did not publish a durable failure" +[ -f "$PARENT/state/remote-replies/ios.source-failed" ] \ + || fail "the failed-read episode left no durable marker" +printf 'working: recovered after transport failure\n' >> "$REMOTE/state/parent-replies.status" +rm -f "$TMP_ROOT/fail-remote-read" +remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" >/dev/null 2>&1 & +for _ in $(seq 1 100); do + grep -q 'recovered after transport failure' "$PARENT/state/ios.status" && break + sleep 0.1 +done +assert_grep 'recovered after transport failure' "$PARENT/state/ios.status" \ + "a relaunched listener did not ingest after transport recovery" +assert_absent "$PARENT/state/remote-replies/ios.source-failed" \ + "a recovered read did not close the failed-read episode" +[ "$(grep -c 'remote reply listener ios failed three consecutive reads' "$PARENT/state/.wake-queue")" -eq 1 ] \ + || fail "one failed-read episode published duplicate wakes" +stop_reply_listener || fail "the recovered listener did not stop" +pass "a failed remote read retries a bounded number of times, wakes once, and exits" # Make local ingestion persistently fail after the delta has been captured. # Its durable generation must remain the only copy until reconciliation. @@ -988,9 +1018,106 @@ assert_grep "offset=$replay_offset" "$PARENT/state/remote-replies/ios.cursor" \ "the recapture did not rebuild the lost cursor" pass "a cursor-loss whole-log recapture is acknowledged quietly with no duplicate wake" +# The old log can lose a few bytes from an already ingested final line. Its +# terminal continuity break is correct, but recovery must prove that the +# retained complete tail matches the acknowledged raw delta before rewinding. +stop_reply_listener || fail "the reply listener did not stop before the rebase fixture" +printf 'note: rebase anchor stays intact\n' >> "$REMOTE/state/parent-replies.status" +rebase_anchor_offset=$(LC_ALL=C wc -c < "$REMOTE/state/parent-replies.status" | tr -d ' ') +printf 'done [corr=0123456789abcdef]: rebase tail is complete\n' \ + >> "$REMOTE/state/parent-replies.status" +rebase_old_offset=$(LC_ALL=C wc -c < "$REMOTE/state/parent-replies.status" | tr -d ' ') +GEN=$((GEN + 1)) +await_reply_result "$PARENT/state/procevent-inbox/$SID.$GEN.result" \ + || fail "the two-line rebase tail was not ingested" +rebase_delta_gen=$GEN +assert_grep "offset=$rebase_old_offset" "$PARENT/state/remote-replies/ios.cursor" \ + "the rebase fixture did not commit the full tail" +stop_reply_listener || fail "the reply listener did not stop before its suffix was shortened" +cp "$REMOTE/state/parent-replies.status" "$TMP_ROOT/rebase-full-source" +perl -e 'my $p=shift; my $n=-s $p; truncate($p,$n-5) or die $!' \ + "$REMOTE/state/parent-replies.status" || fail "could not shorten the rebase fixture by five bytes" +cp "$REMOTE/state/parent-replies.status" "$TMP_ROOT/rebase-short-source" +GEN=$((GEN + 1)) +remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" > "$TMP_ROOT/rebase-break.out" 2>&1 & +RUNNER=$! +wait "$RUNNER" || fail "the shortened reply log did not produce a terminal result" +rebase_break="$PARENT/state/procevent-inbox/$SID.$GEN.result" +assert_present "$rebase_break" "the shortened reply log produced no captured break" +assert_present "${rebase_break%.result}.handled" "the truncation result was not acknowledged" +[ "$(remote_env "$ADAPTER" classify "$rebase_break")" = continuity-broken ] \ + || fail "a five-byte shrink did not break continuity" +assert_grep 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status" \ + "the five-byte shrink did not publish a continuity block" +assert_absent "$PARENT/state/procevent/$SID.source" \ + "the shortened source stayed registered before a guarded rebase" +if remote_env "$ADAPTER" rebase ios --expect-offset "$((rebase_old_offset - 1))" \ + > "$TMP_ROOT/rebase-wrong-offset.out" 2>&1; then + fail "rebase accepted a stale expected cursor offset" +fi +mv "$PARENT/state/remote-replies/ios.$rebase_delta_gen.ingested" \ + "$TMP_ROOT/rebase-ingest-receipt" +if remote_env "$ADAPTER" rebase ios --expect-offset "$rebase_old_offset" \ + > "$TMP_ROOT/rebase-no-proof.out" 2>&1; then + fail "rebase accepted a tail without its ingestion receipt" +fi +assert_grep 'last ingested delta is unavailable' "$TMP_ROOT/rebase-no-proof.out" \ + "missing durable ingestion proof did not explain the rebase refusal" +mv "$TMP_ROOT/rebase-ingest-receipt" \ + "$PARENT/state/remote-replies/ios.$rebase_delta_gen.ingested" +perl -0pi -e 's/rebase anchor/rebase Anchor/' "$REMOTE/state/parent-replies.status" +if remote_env "$ADAPTER" rebase ios --expect-offset "$rebase_old_offset" \ + > "$TMP_ROOT/rebase-changed-tail.out" 2>&1; then + fail "rebase accepted changed bytes in the retained complete tail" +fi +assert_grep 'retained remote tail differs' "$TMP_ROOT/rebase-changed-tail.out" \ + "changed retained bytes did not explain the rebase refusal" +cp "$TMP_ROOT/rebase-short-source" "$REMOTE/state/parent-replies.status" +perl -e 'my ($p,$n)=@ARGV; truncate($p,$n+5) or die $!' \ + "$REMOTE/state/parent-replies.status" "$replay_offset" \ + || fail "could not leave only an incomplete retained line" +perl -0pi -e 's/note:/Note:/' "$REMOTE/state/parent-replies.status" +if remote_env "$ADAPTER" rebase ios --expect-offset "$rebase_old_offset" \ + > "$TMP_ROOT/rebase-changed-fragment.out" 2>&1; then + fail "rebase accepted changed bytes in an incomplete retained line" +fi +assert_grep 'retained remote tail differs' "$TMP_ROOT/rebase-changed-fragment.out" \ + "changed incomplete bytes bypassed the rebase proof" +assert_grep "offset=$rebase_old_offset" "$PARENT/state/remote-replies/ios.cursor" \ + "a refused rebase changed the committed cursor" +assert_absent "$PARENT/state/procevent/$SID.source" \ + "a refused rebase registered a source" +cp "$TMP_ROOT/rebase-short-source" "$REMOTE/state/parent-replies.status" +remote_env "$ADAPTER" rebase ios --expect-offset "$rebase_old_offset" \ + > "$TMP_ROOT/rebase-success.out" \ + || fail "a verified five-byte suffix trim could not rebase: $(cat "$TMP_ROOT/rebase-success.out")" +assert_grep "rebased: ios offset=$rebase_anchor_offset" "$TMP_ROOT/rebase-success.out" \ + "rebase did not stop at the last verified complete line" +rebase_backup=$(sed -n 's/^rebased: ios offset=[0-9]* prior-cursor=//p' "$TMP_ROOT/rebase-success.out") +assert_grep "offset=$rebase_old_offset" "$rebase_backup" \ + "rebase did not preserve the prior cursor for inspection" +assert_grep "offset=$rebase_anchor_offset" "$PARENT/state/remote-replies/ios.cursor" \ + "rebase did not commit its verified complete-line cursor" +assert_grep 'resolved [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status" \ + "rebase did not publish a continuity resolution" +assert_present "$PARENT/state/procevent/$SID.source" \ + "rebase did not register the reply source" +cp "$TMP_ROOT/rebase-full-source" "$REMOTE/state/parent-replies.status" +printf 'done [corr=abcdef0123456789]: reply after verified rebase\n' \ + >> "$REMOTE/state/parent-replies.status" +GEN=$((GEN + 1)) +await_reply_result "$PARENT/state/procevent-inbox/$SID.$GEN.result" \ + || fail "the rebased listener did not ingest the restored tail and next reply" +assert_grep 'reply after verified rebase' "$PARENT/state/ios.status" \ + "the next reply after rebase did not reach the parent" +[ "$(grep -cF 'rebase tail is complete' "$PARENT/state/ios.status")" -eq 1 ] \ + || fail "rebase replay duplicated an already ingested tail line" +pass "a five-byte suffix trim rebases only after exact tail proof and resumes without duplicate mirroring" + # The adapter re-armed at the committed cursor. Truncation is detected from the # next blocking source and escalated once; it is never silently treated as a new # log or re-armed past the break. +continuity_before=$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status") stop_reply_listener || fail "the reply listener did not stop before the continuity break" printf 'failed [corr=fedcba9876543210]: source was replaced\n' > "$REMOTE/state/parent-replies.status" GEN=$((GEN + 1)) @@ -1001,6 +1128,8 @@ RESULT_TWELVE=$(find "$PARENT/state/procevent-inbox" -name "$SID.$GEN.result" -p [ -n "$RESULT_TWELVE" ] || fail "continuity break produced no durable result" [ "$(remote_env "$ADAPTER" classify "$RESULT_TWELVE")" = continuity-broken ] \ || fail "truncated source was not classified as a continuity break" +remote_env "$ADAPTER" terminal "$RESULT_TWELVE" \ + || fail "a continuity break did not classify as terminal" set +e remote_env "$ADAPTER" handle ios "$GEN" "$RESULT_TWELVE" > "$TMP_ROOT/handle-nine.out" 2>&1 handle_rc=$? @@ -1009,9 +1138,9 @@ set -e assert_grep 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status" "continuity break did not escalate" assert_absent "$PARENT/state/procevent/$SID.source" "continuity break was re-armed without an operator rebase" remote_env "$ADAPTER" ingest ios "$RESULT_TWELVE" >/dev/null 2>&1 || true -[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 1 ] \ - || fail "continuity replay duplicated the escalation" -status_line_at_epoch "$(grep -F 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" >/dev/null \ +[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq "$((continuity_before + 1))" ] \ + || fail "a new continuity episode was lost or replay duplicated its escalation" +status_line_at_epoch "$(grep -F 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status" | tail -1)" >/dev/null \ || fail "new continuity escalation has unknown emission time" if [ "${FM_TEST_EVIDENCE:-0}" = 1 ]; then printf '\nNew continuity escalation after ingest retry:\n' @@ -1029,10 +1158,164 @@ assert_absent "$PARENT/state/procevent/$SID.source" \ "refused retirement left the reply source running past its pending-result check" remote_env "$ADAPTER" handle ios "$GEN" "$RESULT_TWELVE" >/dev/null 2>&1 || [ "$?" -eq 3 ] \ || fail "pending continuity result could not be acknowledged after retirement refusal" +printf '0 1700000000 1\n' > "$PARENT/state/remote-replies/ios.lag" +printf 'check: remote reply channel stalled: mate=ios\n0 1700000000\n' \ + > "$PARENT/state/remote-replies/ios.lag-ready" remote_env "$ADAPTER" retire ios >/dev/null assert_absent "$PARENT/state/remote-replies/ios.cursor" "adapter retirement left its cursor" assert_absent "$PARENT/state/remote-replies/ios.caught-up" \ "adapter retirement left a caught-up watermark a later route could inherit" +assert_absent "$PARENT/state/remote-replies/ios.source-failed" \ + "adapter retirement left a failed-read episode a later route could inherit" +assert_absent "$PARENT/state/remote-replies/ios.lag-ready" \ + "retirement left a pending lag receipt for the watcher to announce" pass "remote reply retirement quiesces and refuses unhandled captured results" +# A watcher compares the remote log size with the committed cursor. One lag +# episode re-ensures the listener and wakes only after the bound. While delta +# reads keep failing the cursor stays behind, and a later probe in the same +# episode neither repairs nor wakes again. Cursor progress clears its marker. +lag_failures() { wc -l < "$TMP_ROOT/fail-delta-read" | tr -d ' '; } +export FM_REMOTE_REPLY_FAIL_DELTA_FLAG="$TMP_ROOT/fail-delta-read" +: > "$FM_REMOTE_REPLY_FAIL_DELTA_FLAG" +remote_env "$ADAPTER" arm ios >/dev/null +FM_REMOTE_REPLY_LAG_SECONDS=1 FM_REMOTE_REPLY_LAG_PROBE_SECONDS=1 \ + remote_env "$ADAPTER" lag-check ios > "$TMP_ROOT/lag-first.out" +[ ! -s "$TMP_ROOT/lag-first.out" ] || fail "lag woke before its bound" +[ "$(lag_failures)" -eq 0 ] || fail "lag repaired its listener before its bound" +sleep 1.1 +FM_REMOTE_REPLY_LAG_SECONDS=1 FM_REMOTE_REPLY_LAG_PROBE_SECONDS=1 \ + remote_env "$ADAPTER" lag-check ios > "$TMP_ROOT/lag-second.out" +assert_grep 'remote reply channel stalled: mate=ios' "$TMP_ROOT/lag-second.out" \ + "an aged remote log ahead of its cursor did not wake" +assert_grep 'remote reply channel stalled: mate=ios' \ + "$PARENT/state/remote-replies/ios.lag-ready" \ + "the background watcher has no durable receipt to surface after a lag probe" +for _ in $(seq 1 100); do + [ "$(lag_failures)" -ge 3 ] && [ "$(reply_owner)" = none ] && break + sleep 0.1 +done +[ "$(lag_failures)" -eq 3 ] || fail "a stalled-channel episode did not re-ensure its listener" +[ "$(reply_owner)" = none ] || fail "the re-ensured listener outlived its failed-read budget" +sleep 1.1 +FM_REMOTE_REPLY_LAG_SECONDS=1 FM_REMOTE_REPLY_LAG_PROBE_SECONDS=1 \ + remote_env "$ADAPTER" lag-check ios > "$TMP_ROOT/lag-third.out" +[ ! -s "$TMP_ROOT/lag-third.out" ] || fail "one lag episode woke repeatedly" +sleep 0.5 +[ "$(lag_failures)" -eq 3 ] || fail "one lag episode re-ensured its listener repeatedly" +read -r _ _ lag_alerted < "$PARENT/state/remote-replies/ios.lag" +[ "$lag_alerted" = 1 ] || fail "the still-behind lag episode lost its alerted state" +[ "$(grep -c 'remote reply channel stalled: mate=ios' "$PARENT/state/.wake-queue")" -eq 1 ] \ + || fail "one lag episode queued duplicate wakes" +rm -f "$FM_REMOTE_REPLY_FAIL_DELTA_FLAG" +unset FM_REMOTE_REPLY_FAIL_DELTA_FLAG +remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" > "$TMP_ROOT/lag-catchup.out" 2>&1 & +for _ in $(seq 1 100); do + grep -q 'source was replaced' "$PARENT/state/ios.status" && break + sleep 0.1 +done +assert_grep 'source was replaced' "$PARENT/state/ios.status" \ + "lagged reply was not ingested after the listener restarted" +for _ in $(seq 1 100); do + [ ! -e "$PARENT/state/remote-replies/ios.lag-ready" ] && break + sleep 0.1 +done +assert_absent "$PARENT/state/remote-replies/ios.lag-ready" \ + "committing cursor progress did not clear its obsolete lag receipt" +stop_reply_listener || fail "lag catchup listener did not stop" +for _ in $(seq 1 100); do + [ "$(reply_owner)" = none ] && break + sleep 0.1 +done +sleep 1.1 +FM_REMOTE_REPLY_LAG_SECONDS=1 FM_REMOTE_REPLY_LAG_PROBE_SECONDS=1 \ + remote_env "$ADAPTER" lag-check ios > "$TMP_ROOT/lag-caught-up.out" +[ ! -s "$TMP_ROOT/lag-caught-up.out" ] || fail "a caught-up channel still woke" +assert_absent "$PARENT/state/remote-replies/ios.lag" "catchup did not clear the lag episode" +assert_absent "$PARENT/state/remote-replies/ios.lag-ready" \ + "catchup left a stale receipt for the watcher to announce" +pass "an aged remote reply lag re-ensures its listener and wakes once per episode, and catchup resets it" + +# The live failure mode is a detached reconcile launch that never proves a +# claim, followed by an attached start that can read and apply the same reply. +# Keep both paths in this relay fixture instead of treating a launch-failed +# wake as proof that the remote route itself is broken. +REAL_PERL=$(command -v perl) +export FM_TEST_REAL_PERL="$REAL_PERL" +cat > "$FAKEBIN/perl" <<'SH' +#!/usr/bin/env bash +if [ "${3:-}" = detach ]; then exit 125; fi +exec "$FM_TEST_REAL_PERL" "$@" +SH +chmod +x "$FAKEBIN/perl" +detached_rc=0 +FM_TEST_REAL_PERL="$REAL_PERL" FM_PROCEVENT_LAUNCH_CONFIRM_SECONDS=1 \ + remote_env "$ROOT/bin/fm-procevent.sh" reconcile > "$TMP_ROOT/detached-fail.out" 2>&1 \ + || detached_rc=$? +[ "$detached_rc" -ne 0 ] || fail "an unconfirmed detached launch was reported as healthy" +assert_grep 'failed=1' "$TMP_ROOT/detached-fail.out" \ + "the detached path did not report its unconfirmed launch" +[ "$(reply_owner)" = none ] || fail "a failed detached launch invented an owner" +assert_grep 'remote-reply-ios is registered but its launch did not prove' "$PARENT/state/.wake-queue" \ + "the detached failure had no durable actionable wake" +rm -f "$FAKEBIN/perl" +unset FM_TEST_REAL_PERL +printf 'working: attached launch after detached failure\n' >> "$REMOTE/state/parent-replies.status" +remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" > "$TMP_ROOT/attached-recovery.out" 2>&1 & +for _ in $(seq 1 100); do + grep -q 'attached launch after detached failure' "$PARENT/state/ios.status" && break + sleep 0.1 +done +assert_grep 'attached launch after detached failure' "$PARENT/state/ios.status" \ + "an attached launch did not ingest after detached confirmation failed" +stop_reply_listener || fail "attached recovery listener did not stop" +pass "a failed detached launch surfaces durably and an attached launch ingests the backlog" + +# Re-arm can expose another break immediately. Force it after ensure-listening +# returns, before the rebase command exits, so the resolution must already +# precede the new block in the durable status stream. +GEN=$(find "$PARENT/state/procevent-inbox" -maxdepth 1 -name "$SID.*.result" -print \ + | awk -F. '{ print $(NF - 1) }' | sort -n | tail -1) +GEN=${GEN:-0} +printf 'note: rebase ordering anchor\n' >> "$REMOTE/state/parent-replies.status" +printf 'done [corr=abcdef0123456789]: rebase ordering tail\n' \ + >> "$REMOTE/state/parent-replies.status" +order_old_offset=$(LC_ALL=C wc -c < "$REMOTE/state/parent-replies.status" | tr -d ' ') +GEN=$((GEN + 1)) +await_reply_result "$PARENT/state/procevent-inbox/$SID.$GEN.result" \ + || fail "the ordering fixture was not ingested" +stop_reply_listener || fail "the ordering fixture listener did not stop" +perl -e 'my $p=shift; my $n=-s $p; truncate($p,$n-5) or die $!' \ + "$REMOTE/state/parent-replies.status" || fail "could not shorten the ordering fixture" +GEN=$((GEN + 1)) +remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" > "$TMP_ROOT/order-break.out" 2>&1 & +RUNNER=$! +wait "$RUNNER" || fail "the ordering fixture did not publish its first break" +assert_present "$PARENT/state/procevent-inbox/$SID.$GEN.handled" \ + "the ordering fixture break was not acknowledged" +order_next_gen=$((GEN + 1)) +cat > "$TMP_ROOT/rebase-order-hook" <<'SH' +#!/usr/bin/env bash +printf 'failed: a second continuity break after rebase\n' > "$FM_TEST_REBASE_REMOTE_LOG" +for _ in $(seq 1 800); do + [ -f "$FM_TEST_REBASE_NEXT_HANDLED" ] && exit 0 + sleep 0.05 +done +exit 1 +SH +chmod +x "$TMP_ROOT/rebase-order-hook" +FM_TEST_REBASE_AFTER_ENSURE_HOOK="$TMP_ROOT/rebase-order-hook" \ + FM_TEST_REBASE_REMOTE_LOG="$REMOTE/state/parent-replies.status" \ + FM_TEST_REBASE_NEXT_HANDLED="$PARENT/state/procevent-inbox/$SID.$order_next_gen.handled" \ + remote_env "$ADAPTER" rebase ios --expect-offset "$order_old_offset" \ + > "$TMP_ROOT/rebase-order.out" 2>&1 \ + || fail "the second break could not be reproduced during re-arm: $(cat "$TMP_ROOT/rebase-order.out")" +order_resolved_line=$(grep -nF 'resolved [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status" | tail -1 | cut -d: -f1) +order_blocked_line=$(grep -nF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status" | tail -1 | cut -d: -f1) +[ "$order_blocked_line" -gt "$order_resolved_line" ] \ + || fail "a rebase resolution hid the immediately subsequent continuity break" +assert_absent "$PARENT/state/procevent/$SID.source" \ + "the second continuity break left its source registered" +pass "rebase resolution precedes a break detected during listener re-arm" + echo "ALL TESTS PASSED" diff --git a/tests/fm-session-lock-ancestry.test.sh b/tests/fm-session-lock-ancestry.test.sh index 6076057bb17..14d82c7bd35 100755 --- a/tests/fm-session-lock-ancestry.test.sh +++ b/tests/fm-session-lock-ancestry.test.sh @@ -38,13 +38,14 @@ NAMED_CLAUDE="$FAKEBIN/claude" # liveness questions are decided by the process table alone (FM_TEST_KILL_RC=1 # makes every pid dead). The suite itself may run inside a Claude session whose # CLAUDE_CODE_SESSION_ID and CLAUDE_PID would leak into the expression, so both -# are scrubbed and only FM_TEST_SESSION_ID and FM_TEST_CLAUDE_PID reach it. +# are scrubbed and only fixture identities reach it. lib_eval() { # local fakebin=$1 expr=$2 local -a session_env=() [ -z "${FM_TEST_SESSION_ID:-}" ] || session_env+=("CLAUDE_CODE_SESSION_ID=$FM_TEST_SESSION_ID") [ -z "${FM_TEST_CLAUDE_PID:-}" ] || session_env+=("CLAUDE_PID=$FM_TEST_CLAUDE_PID") - env -u CLAUDE_CODE_SESSION_ID -u CLAUDE_PID ${session_env[@]+"${session_env[@]}"} \ + [ -z "${FM_TEST_CODEX_THREAD_ID:-}" ] || session_env+=("CODEX_THREAD_ID=$FM_TEST_CODEX_THREAD_ID") + env -u CLAUDE_CODE_SESSION_ID -u CLAUDE_PID -u CODEX_THREAD_ID ${session_env[@]+"${session_env[@]}"} \ PATH="$fakebin:$PATH" bash -c " . \"\$0\" kill() { return \${FM_TEST_KILL_RC:-0}; } @@ -91,6 +92,9 @@ SH FM_TEST_CLAUDE_SHAPE="$shape" lib_eval "$fakebin" "fm_session_lock_owned_by_self '$dir/state'" \ || fail "$shape: the session holding the lock did not recognize itself as the owner" done + printf 'codex:stale-thread\n' > "$dir/state/.lock-session" + lib_eval "$fakebin" "fm_session_lock_owned_by_self '$dir/state'" \ + || fail "a stale Codex sidecar masked direct ownership of a Claude anchor" pass "session-lock: a version-named Claude Code session is identified from its install path and argv[0]" } @@ -427,6 +431,219 @@ test_anchor_pid_is_the_model_loop_process_only_for_a_trusted_id() { pass "session-lock: a trusted id anchors the lock on the model-loop process, anything else on the outermost pid" } +test_managed_codex_daemon_is_not_a_thread_identity() { + local dir fakebin state got + dir="$TMP_ROOT/codex-shared-daemon" + fakebin=$(fm_fakebin "$dir") + state="$dir/state" + mkdir -p "$state" + cat > "$fakebin/ps" <<'SH' +#!/usr/bin/env bash +field= pid= +while [ "$#" -gt 0 ]; do + case "$1" in + -o) field=$2; shift 2 ;; + -p) pid=$2; shift 2 ;; + *) shift ;; + esac +done +case "$pid:$field" in + 900:comm=) printf '%s\n' codex ;; + 900:args=) printf '%s\n' 'codex app-server --managed-daemon' ;; + 900:ppid=) printf '%s\n' 1 ;; + *:comm=) printf '%s\n' bash ;; + *:args=) printf '%s\n' 'bash /repo/bin/fm-lock.sh' ;; + *:ppid=) printf '%s\n' 900 ;; +esac +SH + chmod +x "$fakebin/ps" + printf '900\n' > "$state/.lock" + printf 'codex:thread-one\n' > "$state/.lock-session" + got=$(FM_TEST_CODEX_THREAD_ID=thread-one lib_eval "$fakebin" 'fm_session_lock_trusted_session_id') \ + || fail "the Codex thread identity was not trusted under its daemon" + [ "$got" = codex:thread-one ] || fail "the Codex sidecar identity was '$got'" + FM_TEST_CODEX_THREAD_ID=thread-one owned "$fakebin" "$state" \ + || fail "the owning Codex thread did not recognize its lock" + if FM_TEST_CODEX_THREAD_ID=thread-two owned "$fakebin" "$state"; then + fail "another Codex thread claimed the shared daemon's lock" + fi + if owned "$fakebin" "$state"; then + fail "a Codex shell without a thread id claimed the shared daemon's lock" + fi + rm -f "$state/.lock-session" + if FM_TEST_CODEX_THREAD_ID=thread-one owned "$fakebin" "$state"; then + fail "a shared daemon lock without a sidecar was claimed by ancestry" + fi + FM_TEST_CODEX_THREAD_ID=thread-two foreign_owner "$fakebin" "$state" >/dev/null \ + || fail "the other Codex thread did not see a live foreign owner" + pass "session-lock: shared Codex daemon ancestry cannot impersonate another thread" +} + +test_guarded_codex_takeover_requires_quiet_exact_owner() { + local dir fakebin state owner watcher out backup + dir="$TMP_ROOT/codex-takeover" + fakebin=$(fm_fakebin "$dir") + state="$dir/state" + mkdir -p "$state" + sleep 120 & + owner=$! + cat > "$fakebin/ps" <<'SH' +#!/usr/bin/env bash +field= pid= +while [ "$#" -gt 0 ]; do + case "$1" in + -o) field=$2; shift 2 ;; + -p) pid=$2; shift 2 ;; + *) shift ;; + esac +done +if [ "$pid" = "$FM_TEST_OWNER_PID" ]; then + case "$field" in + comm=) printf '%s\n' codex ;; + args=) printf '%s\n' 'codex app-server --managed-daemon' ;; + ppid=) printf '%s\n' 1 ;; + esac +else + case "$field" in + comm=) printf '%s\n' bash ;; + args=) printf '%s\n' 'bash /repo/bin/fm-lock.sh' ;; + ppid=) printf '%s\n' "$FM_TEST_OWNER_PID" ;; + esac +fi +SH + chmod +x "$fakebin/ps" + printf '%s\n' "$owner" > "$state/.lock" + printf 'codex:thread-old\n' > "$state/.lock-session" + codex_lock() { + env -u CLAUDE_CODE_SESSION_ID -u CLAUDE_PID \ + FM_HOME="$dir" FM_TEST_OWNER_PID="$owner" CODEX_THREAD_ID=thread-new \ + PATH="$fakebin:$PATH" "$ROOT/bin/fm-lock.sh" "$@" + } + out=$(codex_lock status) + assert_contains "$out" "take-over --expect-pid $owner --expect-session codex:thread-old" \ + "status did not show a guarded command with its exact current owner" + assert_contains "$out" "--attest-owner-ended $owner/codex:thread-old" \ + "status did not require an explicit old-thread attestation" + if codex_lock >/dev/null 2>&1; then + fail "an ordinary Codex lock acquisition stole the shared daemon's lock" + fi + if codex_lock take-over --expect-pid "$owner" --expect-session codex:thread-old \ + --attest-owner-ended "$owner/codex:thread-old" >/dev/null 2>&1; then + fail "take-over accepted a recent owner lock" + fi + perl -e 'utime(time-600,time-600,$ARGV[0]) or die $!' "$state/.lock" + sleep 120 & + watcher=$! + mkdir -p "$state/.watch.lock" + printf '%s\n' "$watcher" > "$state/.watch.lock/pid" + if codex_lock take-over --expect-pid "$owner" --expect-session codex:thread-old \ + --attest-owner-ended "$owner/codex:thread-old" >/dev/null 2>&1; then + fail "take-over accepted a live watcher" + fi + kill "$watcher" 2>/dev/null || true + wait "$watcher" 2>/dev/null || true + rm -f "$state/.watch.lock/pid" + rmdir "$state/.watch.lock" + if codex_lock take-over --expect-pid "$owner" --expect-session codex:wrong \ + --attest-owner-ended "$owner/codex:wrong" >/dev/null 2>&1; then + fail "take-over accepted a stale expected session" + fi + if codex_lock take-over --expect-pid "$owner" --expect-session codex:thread-old >/dev/null 2>&1; then + fail "take-over accepted an owner without an old-thread attestation" + fi + if codex_lock take-over --expect-pid "$owner" --expect-session codex:thread-old \ + --attest-owner-ended "$owner/codex:wrong" >/dev/null 2>&1; then + fail "take-over accepted an attestation naming a different thread" + fi + out=$(codex_lock take-over --expect-pid "$owner" --expect-session codex:thread-old \ + --attest-owner-ended "$owner/codex:thread-old") \ + || fail "guarded take-over refused an idle exact owner: $out" + assert_contains "$out" 'lock taken over' "take-over did not report its verified handoff" + backup=${out##*prior lock backup } + [ -f "$backup/lock" ] && [ "$(cat "$backup/lock")" = "$owner" ] \ + || fail "take-over did not preserve the prior daemon lock" + [ "$(cat "$backup/lock-session")" = codex:thread-old ] \ + || fail "take-over did not preserve the prior thread sidecar" + [ "$(cat "$state/.lock-session")" = codex:thread-new ] \ + || fail "take-over did not publish the new thread identity" + codex_lock >/dev/null || fail "the new Codex thread did not retain its lock" + rm -f "$state/.lock-session" + perl -e 'utime(time-600,time-600,$ARGV[0]) or die $!' "$state/.lock" + if codex_lock >/dev/null 2>&1; then + fail "a managed Codex daemon with no sidecar was claimed by ancestry" + fi + out=$(codex_lock status) + assert_contains "$out" "take-over --expect-pid $owner --expect-session none" \ + "status did not offer guarded recovery for a missing sidecar" + if codex_lock take-over --expect-pid "$owner" --expect-session none >/dev/null 2>&1; then + fail "missing-sidecar recovery accepted a takeover without attestation" + fi + codex_lock take-over --expect-pid "$owner" --expect-session none \ + --attest-owner-ended "$owner/none" >/dev/null \ + || fail "guarded take-over could not repair a missing sidecar" + [ "$(cat "$state/.lock-session")" = codex:thread-new ] \ + || fail "missing-sidecar recovery did not publish a trusted thread" + kill "$owner" 2>/dev/null || true + wait "$owner" 2>/dev/null || true + pass "session-lock: exact old-thread attestation, quiet guards, and prior-lock backup govern managed-daemon takeover" +} + +test_codex_daemon_without_thread_id_never_acquires() { + local dir fakebin state owner out + dir="$TMP_ROOT/codex-no-thread-id" + fakebin=$(fm_fakebin "$dir") + state="$dir/state" + mkdir -p "$state" + sleep 120 & + owner=$! + cat > "$fakebin/ps" <<'SH' +#!/usr/bin/env bash +field= pid= +while [ "$#" -gt 0 ]; do + case "$1" in + -o) field=$2; shift 2 ;; + -p) pid=$2; shift 2 ;; + *) shift ;; + esac +done +if [ "$pid" = "$FM_TEST_OWNER_PID" ]; then + case "$field" in + comm=) printf '%s\n' codex ;; + args=) printf '%s\n' 'codex app-server --managed-daemon' ;; + ppid=) printf '%s\n' 1 ;; + esac +else + case "$field" in + comm=) printf '%s\n' bash ;; + args=) printf '%s\n' 'bash /repo/bin/fm-lock.sh' ;; + ppid=) printf '%s\n' "$FM_TEST_OWNER_PID" ;; + esac +fi +SH + chmod +x "$fakebin/ps" + codex_lock_as() { # + local -a id_env=() + [ -z "$1" ] || id_env=("CODEX_THREAD_ID=$1") + env -u CLAUDE_CODE_SESSION_ID -u CLAUDE_PID -u CODEX_THREAD_ID ${id_env[@]+"${id_env[@]}"} \ + FM_HOME="$dir" FM_TEST_OWNER_PID="$owner" PATH="$fakebin:$PATH" "$ROOT/bin/fm-lock.sh" + } + for bad_id in '' 'bad id!'; do + if out=$(codex_lock_as "$bad_id" 2>&1); then + fail "a Codex thread without a usable thread id ('$bad_id') acquired the shared daemon lock" + fi + assert_contains "$out" "CODEX_THREAD_ID" "the refusal did not tell the operator to provide a trusted thread id" + [ ! -e "$state/.lock" ] || fail "a refused Codex acquisition still wrote the lock" + [ ! -e "$state/.lock-session" ] || fail "a refused Codex acquisition still wrote the sidecar" + done + codex_lock_as thread-one >/dev/null || fail "a Codex thread with a trusted id could not acquire a free lock" + [ "$(cat "$state/.lock")" = "$owner" ] || fail "the trusted Codex thread did not record the daemon pid" + [ "$(cat "$state/.lock-session")" = codex:thread-one ] || fail "the trusted Codex thread did not publish its sidecar" + codex_lock_as thread-one >/dev/null || fail "the trusted Codex thread did not retain its lock" + kill "$owner" 2>/dev/null || true + wait "$owner" 2>/dev/null || true + pass "session-lock: a Codex thread without a trusted id never records the shared daemon pid" +} + # --- end-to-end layer: the real Stop auto-arm in real process trees ---------- install_autoarm_scripts() { @@ -1106,6 +1323,9 @@ test_harness_beyond_a_gap_never_owns_the_lock test_competing_version_named_session_is_seen_as_live test_same_session_id_owns_a_recycled_background_chain test_anchor_pid_is_the_model_loop_process_only_for_a_trusted_id +test_managed_codex_daemon_is_not_a_thread_identity +test_guarded_codex_takeover_requires_quiet_exact_owner +test_codex_daemon_without_thread_id_never_acquires test_e2e_version_named_session_claims_the_home test_e2e_daemon_parented_session_claims_the_home test_e2e_daemon_parented_version_named_session_keeps_its_lock diff --git a/tests/fm-watch-triage.test.sh b/tests/fm-watch-triage.test.sh index c2a2d60924d..50e795a6994 100755 --- a/tests/fm-watch-triage.test.sh +++ b/tests/fm-watch-triage.test.sh @@ -5803,6 +5803,56 @@ procevent_watch_bg() { # FM_POLL=0.2 FM_SIGNAL_GRACE=1 FM_CHECK_INTERVAL=999999 FM_HEARTBEAT=999999 "$WATCH" > "$out" & } +test_remote_reply_probe_keeps_signal_scan_live() { + local dir state probe_pid ready + dir=$(make_case remote-lag-nonblocking); state="$dir/state" + ( + FM_HOME=$dir + FM_STATE_OVERRIDE=$state + # shellcheck source=/dev/null + . "$WATCH" + remote_reply_lag_tick() { + : > "$state/probe-started" + while [ ! -e "$state/probe-release" ]; do sleep 0.1; done + } + remote_reply_lag_start + probe_pid=$! + trap ': > "$state/probe-release"; wait "$probe_pid" 2>/dev/null || true' EXIT + for _ in $(seq 1 100); do + [ -e "$state/probe-started" ] && break + sleep 0.05 + done + [ -e "$state/probe-started" ] || fail "the detached remote probe never started" + printf 'blocked [at=1700000000]: actionable\n' > "$state/crew.status" + scan_signals > "$dir/signals" + assert_grep "$state/crew.status" "$dir/signals" \ + "a pending remote size read blocked the ordinary status scan" + : > "$state/probe-release" + wait "$probe_pid" || fail "the detached remote probe did not finish" + + mkdir -p "$state/remote-replies" + ready="$state/remote-replies/slow.lag-ready" + printf '0 1700000000 1\n' > "$state/remote-replies/slow.lag" + printf 'check: remote reply channel stalled: mate=slow\n0 1700000000\n' > "$ready" + wake() { + printf '%s\n' "$1" > "$dir/lag-reason" + "$FM_WAKE_POST_OUTPUT_ACTION" 0 + } + remote_reply_lag_surface + assert_grep 'check: remote reply channel stalled: mate=slow' "$dir/lag-reason" \ + "a completed remote probe was not surfaced" + [ ! -e "$ready" ] || fail "a delivered lag receipt was not retired" + printf '0 1700000001 1\n' > "$state/remote-replies/slow.lag" + printf 'check: remote reply channel stalled: mate=slow\n0 1700000001\n' > "$ready" + printf 'schema=fm-remote-reply-cursor.v1\noffset=9\n' > "$state/remote-replies/slow.cursor" + rm -f "$dir/lag-reason" + remote_reply_lag_surface + [ ! -e "$dir/lag-reason" ] || fail "a progressed cursor surfaced an obsolete lag receipt" + trap - EXIT + ) + pass "remote reply size probes cannot block ordinary signal scans and completed probes surface" +} + test_procevent_captured_result_surfaces_proactively() { local dir state out drain_out pid beacon_age dir=$(make_case procevent-delivery); state="$dir/state" @@ -6729,6 +6779,7 @@ test_timer_repair_drops_a_finished_write_deferral_chain test_terminal_first_sight_drops_a_finished_write_deferral_chain test_triage_log_size_cap_accepts_spaced_wc_counts test_procevent_captured_result_surfaces_proactively +test_remote_reply_probe_keeps_signal_scan_live test_procevent_unacknowledged_result_redrains_until_handled test_procevent_marker_keys_are_injective test_procevent_headlines_classify_queue_keys