From 0973dbfebe3e7c8cf908cd670eb8bb8550e7905c Mon Sep 17 00:00:00 2001 From: dnth Date: Fri, 4 Sep 2026 21:24:47 +0800 Subject: [PATCH] fix(watcher): bound the wake consumption wait so an unacknowledged steer cannot park the successor chain The shared Pi/OMP watcher core (bin/fm-primary-watch-core.ts) waited forever for a delivered wake to be acknowledged through before_agent_start. OMP only emits that event when the wake starts a new turn; a wake queued as a steer into an already-running turn never does, so owner.restoring stayed true and every later actionable close was enqueued without a successor arm or a follow-up. A secondmate whose crew ends turns every 30-50s lost its watcher within a minute of every arm. sendWake now races the acknowledgement against a new FM_WATCH_WAKE_CONSUME_TIMEOUT_MS bound (default 15000ms); on timeout it drops the wake token and returns generationIsLive(owner), so the pending loop never re-sends the same steer and the successor chain continues. The restoring gate and the close handler are unchanged. Regression tests in tests/fm-omp-primary.test.sh and tests/fm-pi-watch-extension.test.sh drive two consecutive actionable closes with no before_agent_start and prove a third arm starts, exactly one wake is delivered per close, and the durable wake-queue row is untouched; both fail against the unbounded core. Claude-Session: https://claude.ai/code/session_018sbgQjjCZMmpvcfvxyan1Q --- bin/fm-primary-watch-core.ts | 24 +++++- docs/configuration.md | 1 + docs/verification/supervision.md | 14 ++++ docs/watcher-continuity.md | 2 + tests/fm-omp-primary.test.sh | 121 ++++++++++++++++++++++++++++ tests/fm-pi-watch-extension.test.sh | 81 +++++++++++++++++++ 6 files changed, 242 insertions(+), 1 deletion(-) diff --git a/bin/fm-primary-watch-core.ts b/bin/fm-primary-watch-core.ts index 87b988b914a..052f14e7c1f 100644 --- a/bin/fm-primary-watch-core.ts +++ b/bin/fm-primary-watch-core.ts @@ -297,6 +297,11 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar const retryBaseMs = positiveInteger("FM_WATCH_REARM_RETRY_BASE_MS", 250); const retryMaxMs = positiveInteger("FM_WATCH_REARM_RETRY_MAX_MS", 4000); const retryLimit = positiveInteger("FM_WATCH_REARM_RETRY_LIMIT", 5); + // A delivered wake is acknowledged when the runtime starts a turn whose prompt + // is that wake. A wake the runtime queues into an already-running turn never + // starts one, so that acknowledgement can never arrive; bound the wait so one + // unacknowledged delivery cannot hold the successor chain forever. + const wakeConsumeTimeoutMs = positiveInteger("FM_WATCH_WAKE_CONSUME_TIMEOUT_MS", 15000); const armReadyTimeoutMs = positiveInteger( armReadyTimeoutEnv, process.platform === "win32" ? 35000 : 12000, @@ -576,7 +581,24 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar owner.wakeAcknowledgements.set(token, { content, settle: settleConsumption }); try { await sendFollowUp(content); - return await consumption; + let consumeTimer: ReturnType | undefined; + const consumeTimeout = new Promise<"timeout">((resolveTimeout) => { + consumeTimer = setTimeout(() => resolveTimeout("timeout"), wakeConsumeTimeoutMs); + consumeTimer.unref(); + }); + try { + const consumed = await Promise.race([consumption, consumeTimeout]); + if (consumed !== "timeout") return consumed; + } finally { + clearTimeout(consumeTimer); + } + // The runtime accepted the wake without starting a turn for it: the message + // is already queued into the running conversation, so it counts as + // delivered. Returning true here (never false) keeps the pending loop from + // re-sending the same wake. + owner.wakeAcknowledgements.delete(token); + settleConsumption(true); + return generationIsLive(owner); } catch (error) { owner.wakeAcknowledgements.delete(token); settleConsumption(false); diff --git a/docs/configuration.md b/docs/configuration.md index bcc15c78ceb..ba5a1e1bd87 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -686,6 +686,7 @@ FM_WATCH_ARM_RETIRE_TIMEOUT_MS=1000 # milliseconds Pi/OpenCode wait for an unr FM_WATCH_REARM_RETRY_BASE_MS=250 # Pi/OpenCode adapter base delay for continuity restoration retries FM_WATCH_REARM_RETRY_MAX_MS=4000 # Pi/OpenCode adapter cap for exponential continuity retry delay FM_WATCH_REARM_RETRY_LIMIT=5 # Pi/OpenCode adapter launch-failure retries before surfacing restoration failure +FM_WATCH_WAKE_CONSUME_TIMEOUT_MS=15000 # milliseconds the shared Pi/OMP watcher core waits for a delivered wake to be acknowledged by a new turn before treating a wake queued into an already-running turn as consumed and continuing the successor chain FM_WATCH_CYCLE_LOG_MAX_BYTES=262144 # size cap for the arm-owned watcher lifecycle ledger FM_WATCH_CYCLE_LOG_KEEP_LINES=1000 # newest complete lifecycle rows considered when the ledger is capped FM_WATCHER_STALE_GRACE=300 # defaults to FM_GUARD_GRACE; seconds a live watcher lock may have a stale beacon before re-arm errors diff --git a/docs/verification/supervision.md b/docs/verification/supervision.md index 24165234d17..2d9dc531e63 100644 --- a/docs/verification/supervision.md +++ b/docs/verification/supervision.md @@ -330,6 +330,20 @@ Stale prior-generation tool callbacks could not mutate the active child, repeate Plain Pi and pi-signed share the same tracked `.pi/extensions/fm-primary-pi-watch.ts` path, so both inherit the generation owner stated once in `bin/fm-primary-watch-core.ts`. For OMP, the 18.1.5 SDK source at `src/session/agent-session.ts` shows `newSession()` emitting `session_switch` with `reason: "new"`, `switchSession()` emitting `reason: "resume"`, `fork()` emitting `reason: "fork"`, and `reload()` delegating to `switchSession()`; none of those paths emits `session_shutdown`. The portable OMP regression invokes only those `session_switch` shapes, including a second `resume` event for reload, and proves each replacement arms a new real child while carrying one unconsumed actionable close exactly once. +A 2026-09-04 regression pins the bounded acknowledgement wait after a secondmate home lost its watcher chain behind a wake steered into its own recovery turn: + +```sh +bin/fm-test-run.sh tests/fm-omp-primary.test.sh +bin/fm-test-run.sh tests/fm-pi-watch-extension.test.sh +``` + +```text +ok - OMP unacknowledged wake delivery keeps the successor chain and delivers once per close +ok - Pi unacknowledged wake delivery keeps the successor chain and delivers once per close +``` + +Both cases drive two consecutive actionable closes whose deliveries never receive `before_agent_start` under `FM_WATCH_WAKE_CONSUME_TIMEOUT_MS=300`, and prove that a third real arm child starts, exactly one wake is delivered per close, and the durable queue row is untouched. +Against the unbounded core the same cases time out waiting for the third arm. A 2026-09-03 differential check against installed OMP 18.1.5 ran the opt-in CLI guard on both the watcher-continuity port and pristine `origin/main`: ```sh diff --git a/docs/watcher-continuity.md b/docs/watcher-continuity.md index d9808466469..9334422ab13 100644 --- a/docs/watcher-continuity.md +++ b/docs/watcher-continuity.md @@ -27,6 +27,7 @@ While supervision is still needed and away mode remains inactive, an actionable After an actionable Pi, OMP, or OpenCode child close, the adapter starts and verifies one singleton successor before it delivers the original wake. OMP delivers that follow-up through its host API as a hidden custom `steer` with `triggerTurn`, which starts an idle handling turn without touching an editable TUI draft. +A delivered follow-up is acknowledged when the runtime starts a turn whose prompt is that wake; a wake the runtime queues into an already-running turn never starts one, so the core waits at most `FM_WATCH_WAKE_CONSUME_TIMEOUT_MS` for that acknowledgement, then treats the wake as consumed and continues the successor chain instead of parking every later actionable close behind it. It confirms the handling handoff against that successor before scheduling the follow-up, retries once against the current generation and successor, and treats a failed confirmation as a restoration failure: it classifies the error, retires a successor that is no longer alive, and surfaces exactly one typed message. A failed confirmation is never swallowed. It waits at most one readiness timeout per attempt, then sends TERM and waits a bounded retirement confirmation before the next lock-verified exponential retry. @@ -98,6 +99,7 @@ Only the watcher process touches `state/.last-watcher-beat`; no helper process c `tests/fm-pi-watch-extension.test.sh` checks Pi's first-cycle-or-explicit-repair tool metadata and ownership-based redundant-call no-ops, then simulates actionable and empty child closes against the actual Pi and OpenCode close handlers, blocks prompt delivery to prove the successor launches first, verifies single-flight behavior, changes the session lock before close to prove ownership is rechecked, and hangs each successor arm to prove bounded fallback delivery includes the typed restoration failure. The same suite covers Pi same-process replacement for `/new`, `/resume`, `/fork`, and reload, same-instance shutdown-plus-start, automatic re-arm before any model turn, an in-flight actionable close consumed exactly once by the replacement, stale prior-generation callbacks, repeated transitions with exactly one live cycle, disappearance of the shutting-down refusal after a valid replacement activates, and terminal quit still refusing late rearm. `tests/fm-omp-primary.test.sh` drives OMP `/new`, `/resume`, `/fork`, and reload through real watcher child processes and the adapter's public `session_switch` event, proving automatic re-arm plus the same in-flight actionable handoff without using `session_shutdown` as a replacement signal. +Both suites also prove the bounded acknowledgement wait: two consecutive actionable closes whose deliveries never receive `before_agent_start` still start a third arm and deliver exactly one wake per close. It also covers a fresh factory bind without a prior shutdown: the superseded binding retires its arm child, cannot arm or reclaim ownership through retained session callbacks, the live binding keeps exactly one arm child, and repeated binds never accumulate more than one process-exit fallback. `tests/fm-omp-primary.test.sh` covers OMP's binding of the same core to its native session, watcher, and shutdown surfaces, pins the input-preserving custom-steer delivery contract, pins the recovery handling handshake OMP performs before delivering that steer, and pins the single typed wake a refused handshake produces. The opt-in `tests/fm-omp-primary-live-e2e.test.sh` proves a real OMP watcher wake reaches the session while its exact pending draft remains intact. diff --git a/tests/fm-omp-primary.test.sh b/tests/fm-omp-primary.test.sh index a824f712197..5dafa540cf7 100755 --- a/tests/fm-omp-primary.test.sh +++ b/tests/fm-omp-primary.test.sh @@ -1081,6 +1081,126 @@ JS pass "OMP /new /resume /fork and reload session_switch paths auto-arm and carry an in-flight actionable close" } +# A wake that OMP queues into an already-running turn never starts a turn, so +# before_agent_start never acknowledges it. The shared core must bound that wait: +# the next actionable close still has to start its successor and deliver its own +# wake exactly once instead of parking the whole chain behind the first +# unacknowledged delivery. +test_native_omp_unacknowledged_wake_keeps_successor_chain() { + local fixture out status=0 + fixture="$TMP_ROOT/native-unacknowledged-wake" + mkdir -p "$fixture/.omp/extensions/lib" "$fixture/bin" "$fixture/state" "$fixture/config" + cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" + cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$fixture/.omp/extensions/lib/fm-task-inbox-doorbell.ts" + cp "$ROOT/bin/fm-primary-watch-core.ts" "$fixture/bin/fm-primary-watch-core.ts" + cp "$ROOT/bin/fm-pi-compatible-runtimes" "$fixture/bin/fm-pi-compatible-runtimes" + : > "$fixture/AGENTS.md" + git init -q -b main "$fixture" + cat > "$fixture/bin/fm-gate-refuse-lib.sh" <<'SH' +fm_is_gate_agent() { return 1; } +SH + cat > "$fixture/bin/fm-primary-scope-lib.sh" <<'SH' +fm_primary_scope_matches() { return 0; } +SH + cat > "$fixture/bin/fm-operational-input.sh" <<'SH' +#!/usr/bin/env bash +printf 'encoded:%s:%s' "$2" "$(cat)" +SH + cat > "$fixture/bin/fm-sessionstart-nudge.sh" <<'SH' +#!/usr/bin/env bash +exit 0 +SH + cat > "$fixture/bin/fm-watch-arm.sh" <<'SH' +#!/usr/bin/env bash +state=${FM_STATE_OVERRIDE:?} +count=$(cat "$state/watch-count" 2>/dev/null || printf 0) +count=$((count + 1)) +printf '%s\n' "$count" > "$state/watch-count" +printf 'arm=%s predecessor=%s\n' "$$" "${FM_WATCH_PREDECESSOR_ARM_PID:-none}" >> "$state/arm-log" +printf 'watcher: started pid=%s (beacon fresh)\n' "$$" +trap 'exit 0' TERM INT +if [ "$count" -le 2 ]; then + while [ ! -e "$state/watch-trigger-$count" ]; do sleep 0.02; done + printf 'signal: omp unacknowledged wake %s\n' "$count" + exit 0 +fi +while [ ! -e "$state/watch-stop" ]; do sleep 0.02; done +SH + cat > "$fixture/bin/fm-turnend-guard.sh" <<'SH' +#!/usr/bin/env bash +exit 0 +SH + for script in fm-subagent-pretool-check.sh fm-cd-pretool-check.sh fm-arm-pretool-check.sh; do + cat > "$fixture/bin/$script" <<'SH' +#!/usr/bin/env bash +exit 0 +SH + done + chmod +x "$fixture/bin/"*.sh + + out=$(EXTENSION="$fixture/.omp/extensions/fm-primary-omp.ts" FM_HOME="$fixture" \ + FM_ROOT_OVERRIDE="$fixture" FM_STATE_OVERRIDE="$fixture/state" FM_CONFIG_OVERRIDE="$fixture/config" \ + FM_WATCH_WAKE_CONSUME_TIMEOUT_MS=300 node --input-type=module 2>&1 <<'JS' +import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { pathToFileURL } from "node:url"; + +const handlers = new Map(); +const steers = []; +const api = { + zod: { object: () => ({}) }, + on(name, handler) { handlers.set(name, handler); }, + registerCommand() {}, + registerTool() {}, + // The runtime queues every wake as a steer into a running turn: delivery + // resolves, no turn starts, and before_agent_start is never invoked for it. + sendMessage(message) { steers.push(String(message?.content ?? "")); }, +}; +const state = process.env.FM_STATE_OVERRIDE; +const bound = Number(process.env.FM_WATCH_WAKE_CONSUME_TIMEOUT_MS); +const count = () => existsSync(`${state}/watch-count`) ? Number(readFileSync(`${state}/watch-count`, "utf8").trim()) : 0; +const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); +async function waitFor(pred, label) { + for (let i = 0; i < 500; i += 1) { + if (pred()) return; + await sleep(10); + } + throw new Error(`timeout waiting for ${label}`); +} +const queue = `${state}/.wake-queue`; +const queueRow = "1\t1\tsignal\tcrew.turn-ended\tsignal: durable row\n"; +writeFileSync(queue, queueRow); +writeFileSync(`${state}/.lock`, `${process.pid}\n`); +process.argv[1] = process.env.EXTENSION; +const extension = await import(`${pathToFileURL(process.env.EXTENSION).href}?unacknowledged=${Date.now()}`); +extension.default(api); +const context = { sessionManager: { getSessionFile: () => undefined } }; +await handlers.get("session_start")({ type: "session_start" }, context); +await waitFor(() => count() === 1, "initial automatic OMP arm"); +writeFileSync(`${state}/watch-trigger-1`, "trigger\n"); +await waitFor(() => steers.length === 1 && count() === 2, "first OMP delivery and its successor"); +writeFileSync(`${state}/watch-trigger-2`, "trigger\n"); +await waitFor(() => count() === 3, "third OMP arm after the unacknowledged delivery"); +await waitFor(() => steers.length === 2, "second OMP delivery"); +await sleep(bound * 3); +if (count() !== 3) throw new Error(`expected exactly three arms, got ${count()}`); +if (steers.length !== 2) throw new Error(`expected exactly one delivery per close, got ${steers.length}: ${steers.join(" | ")}`); +if (!steers[0].includes("signal: omp unacknowledged wake 1") || !steers[1].includes("signal: omp unacknowledged wake 2")) { + throw new Error(`deliveries did not match their closes: ${steers.join(" | ")}`); +} +const armRows = readFileSync(`${state}/arm-log`, "utf8").trim().split("\n"); +if (!/predecessor=[0-9]+$/.test(armRows[2])) throw new Error(`third arm lost its predecessor identity: ${armRows.join(" | ")}`); +if (readFileSync(queue, "utf8") !== queueRow) throw new Error("the durable wake queue row was altered by the bounded acknowledgement"); +writeFileSync(`${state}/watch-stop`, "stop\n"); +await handlers.get("session_shutdown")({ type: "session_shutdown" }, context); +console.log("omp-unacknowledged-wake-ok"); +JS + ) || status=$? + expect_code 0 "$status" "OMP unacknowledged wake continuity" + assert_contains "$out" omp-unacknowledged-wake-ok "OMP unacknowledged-wake continuity test did not complete" + pass "OMP unacknowledged wake delivery keeps the successor chain and delivers once per close" +} + test_resolve_path_uses_node_when_readlink_f_is_unavailable test_exact_bun_omp_primary_identity test_standalone_omp_primary_identity @@ -1093,3 +1213,4 @@ test_native_primary_extension_contract test_native_omp_confirms_recovery_handling_delivery test_native_omp_refused_handling_delivery_is_typed_once test_native_omp_session_switch_carries_inflight_actionable_close +test_native_omp_unacknowledged_wake_keeps_successor_chain diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index bc13f37d927..c9cd66881e3 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -1257,6 +1257,86 @@ EOF pass "Pi session replacement auto-arms and carries its in-flight actionable close" } +# The Pi binding shares the core, so a delivery the runtime never acknowledges +# through before_agent_start must not park its successor chain either. +test_pi_unacknowledged_wake_keeps_successor_chain() { + local repo home plugin log stop out status + repo="$TMP_ROOT/pi-unacknowledged-wake-root" + home="$TMP_ROOT/pi-unacknowledged-wake-home" + log="$TMP_ROOT/pi-unacknowledged-wake.log" + stop="$TMP_ROOT/pi-unacknowledged-wake.stop" + mkdir -p "$repo/bin" "$home/state" "$home/config" + install_pi_watch_extension_fixture "$repo" + plugin="$repo/.pi/extensions/fm-primary-pi-watch.ts" + cat > "$repo/bin/fm-watch-arm.sh" <<'SH' +#!/usr/bin/env bash +printf 'arm=%s predecessor=%s\n' "$$" "${FM_WATCH_PREDECESSOR_ARM_PID:-none}" >> "${FM_ARM_LOG:?}" +count=$(grep -c '^arm=' "$FM_ARM_LOG") +printf 'watcher: started pid=%s (beacon fresh)\n' "$$" +trap 'exit 0' TERM INT +if [ "$count" -le 2 ]; then + while [ ! -e "$FM_ARM_LOG.trigger-$count" ]; do sleep 0.02; done + printf 'signal: pi unacknowledged wake %s\n' "$count" + exit 0 +fi +while [ ! -e "$FM_STOP_FILE" ]; do sleep 0.02; done +SH + chmod +x "$repo/bin/fm-watch-arm.sh" + out=$(PLUGIN="$plugin" FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_ARM_LOG="$log" FM_STOP_FILE="$stop" \ + FM_WATCH_WAKE_CONSUME_TIMEOUT_MS=300 node --input-type=module 2>&1 <<'EOF' +import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { pathToFileURL } from "node:url"; + +const handlers = new Map(); +const prompts = []; +const pi = { + on(event, handler) { handlers.set(event, handler); }, + registerCommand() {}, + registerTool() {}, + // Delivery resolves, but the runtime never starts a turn for it, so the + // before_agent_start acknowledgement never arrives. + sendUserMessage: async (message) => { prompts.push(String(message)); }, + events: { on() {}, emit() {} }, +}; +const bound = Number(process.env.FM_WATCH_WAKE_CONSUME_TIMEOUT_MS); +const armRows = () => existsSync(process.env.FM_ARM_LOG) + ? readFileSync(process.env.FM_ARM_LOG, "utf8").trim().split("\n").filter((row) => row.startsWith("arm=")) + : []; +const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); +async function waitFor(pred, label) { + for (let i = 0; i < 500; i += 1) { + if (pred()) return; + await sleep(10); + } + throw new Error(`timeout waiting for ${label}`); +} +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +const mod = await import(pathToFileURL(process.env.PLUGIN).href); +mod.default(pi); +await handlers.get("session_start")?.({ type: "session_start", reason: "startup" }, {}); +await waitFor(() => armRows().length === 1, "initial automatic Pi arm"); +writeFileSync(`${process.env.FM_ARM_LOG}.trigger-1`, "trigger\n"); +await waitFor(() => prompts.length === 1 && armRows().length === 2, "first Pi delivery and its successor"); +writeFileSync(`${process.env.FM_ARM_LOG}.trigger-2`, "trigger\n"); +await waitFor(() => armRows().length === 3, "third Pi arm after the unacknowledged delivery"); +await waitFor(() => prompts.length === 2, "second Pi delivery"); +await sleep(bound * 3); +if (armRows().length !== 3) throw new Error(`expected exactly three arms, got ${armRows().join(" | ")}`); +if (prompts.length !== 2) throw new Error(`expected exactly one delivery per close, got ${prompts.length}: ${prompts.join(" | ")}`); +if (!prompts[0].includes("signal: pi unacknowledged wake 1") || !prompts[1].includes("signal: pi unacknowledged wake 2")) { + throw new Error(`deliveries did not match their closes: ${prompts.join(" | ")}`); +} +if (!/predecessor=[0-9]+$/.test(armRows()[2])) throw new Error(`third arm lost its predecessor identity: ${armRows().join(" | ")}`); +writeFileSync(process.env.FM_STOP_FILE, "stop\n"); +await handlers.get("session_shutdown")?.({ type: "session_shutdown", reason: "quit" }, {}); +EOF + ) + status=$? + expect_code 0 "$status" "Pi unacknowledged wake must not park the successor chain" + [ -z "$out" ] || fail "Pi unacknowledged-wake test printed output: $out" + pass "Pi unacknowledged wake delivery keeps the successor chain and delivers once per close" +} + test_pi_rebind_without_shutdown_supersedes_prior_binding() { local repo home plugin child_pid_file arm_log out status repo="$TMP_ROOT/pi-rebind-root" @@ -2511,6 +2591,7 @@ test_pi_actionable_close_rechecks_session_lock test_pi_arm_distinguishes_session_lock_ownership test_pi_session_transition_generation_owner test_pi_session_replacement_carries_inflight_actionable_close +test_pi_unacknowledged_wake_keeps_successor_chain test_pi_rebind_without_shutdown_supersedes_prior_binding test_pi_process_exit_cleanup_listener_lifecycle test_pi_process_exit_cleanup_stops_arm_child