From 5ff81ffb47b595762352f2414f94d2c5d7fcabd6 Mon Sep 17 00:00:00 2001 From: Sungin Kim Date: Sat, 12 Sep 2026 00:11:11 +0000 Subject: [PATCH 1/4] fix(pi): restore watcher successors while an earlier wake is still being handled A hung branch settlement was holding the restore lock, so a later worker completion could leave monitoring stopped. Split restore from serialized delivery so later actionable closes still start a successor. Co-authored-by: Cursor --- .pi/extensions/fm-primary-pi-watch.ts | 52 +++++-- bin/fm-test-run.sh | 2 + docs/verification/supervision.md | 24 +++ docs/watcher-continuity.md | 1 + tests/fm-pi-hung-delivery-herdr-e2e.test.sh | 164 ++++++++++++++++++++ tests/fm-pi-watch-extension.test.sh | 157 ++++++++++++++++++- 6 files changed, 386 insertions(+), 14 deletions(-) create mode 100755 tests/fm-pi-hung-delivery-herdr-e2e.test.sh diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index 6c4c020cd87..d90f6d23b35 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -30,6 +30,13 @@ // consumes at the user message_start carrying the exact wake text; either // event finishes the pending record, and a still-unconsumed record rides the // replacement handoff. +// +// Restore versus delivery (stated once here): +// restoring is true only while a successor arm is started and verified. +// delivering is true while the serialized pending-wake pump is in flight. +// A later actionable close restores a successor even when the previous wake's +// branch settlement has not resolved. Failure closes during delivering still +// defer their bounded retry until that delivery settles. import { spawn, spawnSync, type ChildProcess } from "node:child_process"; import { createHash } from "node:crypto"; import { existsSync, mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs"; @@ -100,6 +107,7 @@ type SessionGeneration = { cleanupTimer: ReturnType | null; retryFailures: number; restoring: boolean; + delivering: boolean; seq: number; awayStandby: boolean; awayPoll: ReturnType | null; @@ -435,6 +443,7 @@ function createGeneration(): SessionGeneration { cleanupTimer: null, retryFailures: 0, restoring: false, + delivering: false, seq: 0, awayStandby: false, awayPoll: null, @@ -723,7 +732,7 @@ export default function (pi: ExtensionAPI) { if (retiring) retiring.kill("SIGTERM"); return; } - if (!owner.awayStandby || owner.child || owner.retryTimer || owner.restoring) return; + if (!owner.awayStandby || owner.child || owner.retryTimer || owner.restoring || owner.delivering) return; owner.awayStandby = false; if (!startArm(owner).ok) owner.awayStandby = true; else void processPendingActionables(owner); @@ -781,13 +790,33 @@ export default function (pi: ExtensionAPI) { owner.cleanupTimer = timer; } + async function restoreContinuity(owner: SessionGeneration, predecessorArmPid: string): Promise<{ + failure: string; + recovery?: { generation: string; watcherPid: string }; + }> { + if (owner.restoring) return { failure: "" }; + owner.restoring = true; + try { + return await restoreAfterActionableClose(owner, predecessorArmPid); + } finally { + owner.restoring = false; + } + } + async function processPendingActionables(owner: SessionGeneration): Promise { - if (!generationIsLive(owner) || owner.restoring || owner.pendingActionables.length === 0) return; + if (!generationIsLive(owner) || owner.pendingActionables.length === 0) return; if (awayModeActive()) { owner.awayStandby = true; return; } - owner.restoring = true; + if (owner.delivering) { + const newest = owner.pendingActionables[owner.pendingActionables.length - 1]; + if (newest && !newest.delivered && !owner.unconsumedWakes.has(newest.token)) { + void restoreContinuity(owner, newest.predecessorArmPid); + } + return; + } + owner.delivering = true; const attemptedCleanup = new Set(); try { while (generationIsLive(owner) && !awayModeActive() && owner.pendingActionables.length > 0) { @@ -832,7 +861,7 @@ export default function (pi: ExtensionAPI) { // A new restoration supersedes whatever became of the previous // successor; only a failure during this delivery is retried after it. owner.deferredClose = null; - const restoration = await restoreAfterActionableClose(owner, pending.predecessorArmPid); + const restoration = await restoreContinuity(owner, pending.predecessorArmPid); if (!generationIsLive(owner)) { settleClaim("failed"); releaseClaim(); @@ -875,8 +904,8 @@ export default function (pi: ExtensionAPI) { const detail = error instanceof Error ? error.message : String(error); surfaceFailure(owner, `watcher: FAILED - Pi extension could not deliver an actionable wake\n${detail}`); } finally { + owner.delivering = false; if (generationIsLive(owner)) { - owner.restoring = false; if (owner.pendingActionables.some((pending) => pending.delivered)) schedulePendingCleanup(owner); // No bare arm is launched here. A generation without a child at this // point has either delivered a typed restoration failure after its @@ -1147,11 +1176,12 @@ export default function (pi: ExtensionAPI) { return; } if (!generationIsLive(owner)) return; - if (owner.restoring) { - // The pipeline is still delivering the wake this successor was - // started for. A verified successor that failed on its own keeps its - // bounded retry for the end of that delivery; an unready child closing - // here was retired by the restoration itself. + if (owner.restoring || owner.delivering) { + // Restore still owns this child, or the serialized pump is still + // delivering the wake this successor was started for. A verified + // successor that failed on its own keeps its bounded retry for the + // end of that delivery; an unready child closing here was retired by + // the restoration itself. if (verified && !armRetired.has(armChild)) { owner.deferredClose = { message: classification.message, predecessorArmPid: predecessor }; } @@ -1166,7 +1196,7 @@ export default function (pi: ExtensionAPI) { settleReadiness(false); releaseChild(); if (!generationIsLive(owner)) return; - if (owner.restoring) return; + if (owner.restoring || owner.delivering) return; scheduleRetry(owner, `watcher: FAILED - Pi extension arm child ${id} failed: ${error.message}`, String(armChild.pid ?? "")); }); return { diff --git a/bin/fm-test-run.sh b/bin/fm-test-run.sh index ff9aa16386c..4273394f781 100755 --- a/bin/fm-test-run.sh +++ b/bin/fm-test-run.sh @@ -335,6 +335,7 @@ family_for_basename() { fm-rovo-signals-live-e2e.test.sh|\ fm-opencode-primary-live-e2e.test.sh|fm-pi-branch-live-e2e.test.sh|\ fm-pi-branch-responsiveness-live-e2e.test.sh|\ + fm-pi-hung-delivery-herdr-e2e.test.sh|\ fm-pi-primary-live-e2e.test.sh|fm-omp-primary-live-e2e.test.sh|fm-stow-horizon-live-e2e.test.sh|\ fm-sessionstart-hook-live-e2e.test.sh|fm-sessionstart-instruction-refresh-live-e2e.test.sh|\ fm-quota-array-dispatch-live-e2e.test.sh|fm-send-secondmate-marker-herdr-e2e.test.sh|\ @@ -692,6 +693,7 @@ tests/fm-pending-reply.test.sh 24679 tests/fm-pi-branch-extension.test.sh 22239 tests/fm-pi-branch-live-e2e.test.sh 56 tests/fm-pi-branch-responsiveness-live-e2e.test.sh 21 +tests/fm-pi-hung-delivery-herdr-e2e.test.sh 23 tests/fm-pi-primary-live-e2e.test.sh 20 tests/fm-pi-watch-extension.test.sh 42970 tests/fm-pointer-check.test.sh 1314 diff --git a/docs/verification/supervision.md b/docs/verification/supervision.md index 44950bf4bc2..683924616be 100644 --- a/docs/verification/supervision.md +++ b/docs/verification/supervision.md @@ -469,6 +469,30 @@ grok 0.2.103 (89c3d36fb6f1) [stable] Pi 0.81.1 repeated the continuity and clean-exit lifecycle on 2026-07-23 after the Calm presentation changes. +Pi later-cycle restore while an earlier branch settlement is still hung is covered by the portable watcher suite plus this isolated-lab live guard: + +```sh +tests/fm-pi-watch-extension.test.sh +FM_PI_HUNG_DELIVERY_HERDR_E2E=1 tests/fm-pi-hung-delivery-herdr-e2e.test.sh +``` + +The live guard launches a real Pi in a named non-default Herdr lab, hangs branch settlement, requires two further worker completions to keep a live successor, and checks the watcher beacon across an unattended interval. +It does not claim turn-end marker identity, branch-session rotation, or ready-to-validate classification. + +On 2026-09-11 that live guard was run against Pi 0.85.1 in a named non-default Herdr 0.8.2 lab: + +```sh +FM_PI_HUNG_DELIVERY_HERDR_E2E=1 tests/fm-pi-hung-delivery-herdr-e2e.test.sh +``` + +Observed output: + +```text +ok - isolated Pi hung-settlement Herdr lab restored later-cycle successors and kept a fresh beacon + +all fm-pi-hung-delivery-herdr-e2e tests passed +``` + Pi same-process session-transition ownership was verified on 2026-07-27 against the tracked extension with a faithful in-process factory rebind (module cache retained, real arm children): ```sh diff --git a/docs/watcher-continuity.md b/docs/watcher-continuity.md index 3adb6bc02bf..04b84780832 100644 --- a/docs/watcher-continuity.md +++ b/docs/watcher-continuity.md @@ -38,6 +38,7 @@ The omp adapter uses the bounded readiness and retirement deadlines owned by its When that retained arm later closes, its actual close is classified as a new supervised event without replaying the earlier fallback. After the configured retry bound is exhausted, it delivers the original wake with a typed continuity-restoration failure even if every successor arm hung without reporting readiness. This is deliberate Option B ordering: the fleet is protected before the model handles the wake whenever restoration succeeds, but the model is never left blind when it does not. +Pi's restore-versus-delivery split is owned by the generation contract in `.pi/extensions/fm-primary-pi-watch.ts`: successor restore is not held across wake delivery. Claude's Stop hook starts the successor arm at the next Stop after the handling turn, rather than before notification as Pi, omp, and OpenCode do. The durable wake queue preserves actionable events during the residual active-turn window, and the bounded turn-end guard enforces recovery at Stop when no watcher is live and no open generation claim is still deciding, so a finished, hung, or identity-mismatched claim cannot suppress it ([`turnend-guard.md`](turnend-guard.md#harness-integrations) owns that boundary). diff --git a/tests/fm-pi-hung-delivery-herdr-e2e.test.sh b/tests/fm-pi-hung-delivery-herdr-e2e.test.sh new file mode 100755 index 00000000000..1adceb44241 --- /dev/null +++ b/tests/fm-pi-hung-delivery-herdr-e2e.test.sh @@ -0,0 +1,164 @@ +#!/usr/bin/env bash +# Opt-in real Pi/Herdr regression for later-cycle restore while branch +# settlement is hung. Every Herdr call is routed through fm-herdr-lab.sh. +# The named lab is never default. Isolated FM_HOME. Abort before any provider +# call. The portable suite pins the same restore/delivery split without Pi. +set -u + +# shellcheck source=tests/lib.sh +. "$(dirname "${BASH_SOURCE[0]}")/lib.sh" +# shellcheck source=tests/herdr-test-safety.sh +. "$(dirname "${BASH_SOURCE[0]}")/herdr-test-safety.sh" + +if [ "${FM_PI_HUNG_DELIVERY_HERDR_E2E:-0}" != 1 ]; then + echo "skip: set FM_PI_HUNG_DELIVERY_HERDR_E2E=1 to run the isolated Pi hung-settlement Herdr regression" + exit 0 +fi + +command -v pi >/dev/null 2>&1 || fail "pi not found" +command -v herdr >/dev/null 2>&1 || fail "herdr not found" +command -v jq >/dev/null 2>&1 || fail "jq not found" + +herdr_forget_inherited_pane + +LAB_HELPER=${HERDR_LAB_HELPER:-$ROOT/bin/fm-herdr-lab.sh} +SESSION=$("$LAB_HELPER" name fm-pi-hung-deliv) +TMP_ROOT=$(fm_test_tmproot fm-pi-hung-delivery-herdr-e2e) +HOME_DIR="$TMP_ROOT/home" +PROJECT="$TMP_ROOT/project" +PI_DIR="$TMP_ROOT/pi-agent" +HANG_EXT="$TMP_ROOT/hang-dispatch.ts" +LAUNCH="$TMP_ROOT/launch-pi.sh" +PANE= +UNATTENDED_SECS=8 + +cleanup() { + local rc=$? + trap - EXIT + if [ -n "$PANE" ]; then + "$LAB_HELPER" run "$SESSION" pane send-text "$PANE" '/quit' >/dev/null 2>&1 || true + sleep 1 + fi + if ! "$LAB_HELPER" teardown "$SESSION"; then + rc=1 + fi + rm -rf "$TMP_ROOT" + exit "$rc" +} +trap cleanup EXIT + +"$LAB_HELPER" provision "$SESSION" + +mkdir -p "$HOME_DIR"/{state,config,data} "$PROJECT" "$PI_DIR" +printf '# Synthetic isolated hung-settlement lab\n' > "$PROJECT/AGENTS.md" + +cat > "$HANG_EXT" <<'EOF' +import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; +export default function (pi: ExtensionAPI) { + pi.on("project_trust", () => ({ trusted: "yes", remember: false })); + pi.events?.on?.("fm-branch-supervision:dispatch", (data: { accept?: (p: Promise) => void }) => { + data.accept?.(new Promise(() => {})); + }); + pi.on("before_agent_start", (_event, ctx) => { + ctx.abort(); + }); +} +EOF + +wait_markers() { + local i=0 + while [ "$i" -lt 120 ]; do + if [ -f "$HOME_DIR/state/.pi-watch-extension-loaded" ] && [ -f "$HOME_DIR/state/.pi-turnend-extension-loaded" ]; then + return 0 + fi + sleep 0.5 + i=$((i + 1)) + done + return 1 +} + +cycle_count() { # + local needle=$1 file=$HOME_DIR/state/.watch-cycle-exits.log count=0 + if [ -f "$file" ]; then + count=$(grep -cE "$needle" "$file" || true) + fi + printf '%s\n' "${count:-0}" +} + +wait_started() { # + local need=$1 i=0 + while [ "$i" -lt 80 ]; do + [ "$(cycle_count 'successor=started:')" -ge "$need" ] && return 0 + sleep 0.5 + i=$((i + 1)) + done + return 1 +} + +beacon_age_s() { + local beat=$HOME_DIR/state/.last-watcher-beat now mtime + [ -f "$beat" ] || return 1 + now=$(date +%s) + mtime=$(stat -c %Y "$beat") + printf '%s\n' "$((now - mtime))" +} + +OUT=$("$LAB_HELPER" run "$SESSION" workspace create --cwd "$PROJECT" --label fm-pi-hung-deliv --no-focus) \ + || fail "Herdr lab workspace create failed" +PANE=$(printf '%s' "$OUT" | jq -er '.result.root_pane.pane_id') \ + || fail "Herdr lab workspace create omitted pane id" + +cat > "$LAUNCH" < $(printf %q "$HOME_DIR/state/.lock") +exec env \\ + FM_HOME=$(printf %q "$HOME_DIR") \\ + FM_ROOT_OVERRIDE=$(printf %q "$ROOT") \\ + PI_CODING_AGENT_DIR=$(printf %q "$PI_DIR") \\ + FM_POLL=1 \\ + FM_SIGNAL_GRACE=0 \\ + FM_HEARTBEAT=600 \\ + FM_CHECK_INTERVAL=999999 \\ + pi --approve --no-session --no-context-files --no-extensions \\ + -e $(printf %q "$HANG_EXT") \\ + -e $(printf %q "$ROOT/.pi/extensions/fm-primary-turnend-guard.ts") \\ + -e $(printf %q "$ROOT/.pi/extensions/fm-primary-pi-watch.ts") \\ + --model openai-codex/gpt-5.6-sol --thinking low +EOF +chmod +x "$LAUNCH" + +"$LAB_HELPER" run "$SESSION" pane run "$PANE" "$LAUNCH" >/dev/null \ + || fail "could not launch isolated Pi in the named Herdr lab" + +wait_markers || fail "Pi watch and turn-end extensions did not load in the isolated lab" + +cat > "$HOME_DIR/state/hung-cycle.meta" <> "$HOME_DIR/state/hung-cycle.status" +wait_started 1 || fail "first actionable close did not start a successor while settlement hung" + +printf 'done: simulated hung completion 2\n' >> "$HOME_DIR/state/hung-cycle.status" +wait_started 2 || fail "second actionable close did not restore a successor while settlement hung" + +printf 'done: simulated hung completion 3\n' >> "$HOME_DIR/state/hung-cycle.status" +wait_started 3 || fail "third actionable close did not restore a successor while settlement hung" + +none=$(cycle_count 'successor=none') +[ "$none" -eq 0 ] || fail "hung settlement left successor=none after later-cycle restores (none=$none)" + +sleep "$UNATTENDED_SECS" +age=$(beacon_age_s) || fail "watcher beacon missing after the unattended interval" +[ "$age" -le $((UNATTENDED_SECS + 5)) ] || fail "watcher beacon went stale during the unattended interval (age=${age}s)" +[ "$(cycle_count 'successor=started:')" -ge 3 ] || fail "unattended interval lost restored successors" +[ "$(cycle_count 'successor=none')" -eq 0 ] || fail "unattended interval recorded successor=none" + +printf 'ok - isolated Pi hung-settlement Herdr lab restored later-cycle successors and kept a fresh beacon\n' +printf '\nall fm-pi-hung-delivery-herdr-e2e tests passed\n' diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index 6ac6934f393..c716f4a6425 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -2571,7 +2571,7 @@ if (previous.prompts.length !== 0) { } await waitFor(() => liveArms().length === 1 && armRows().length >= 2, "old-session successor"); writeFileSync(process.env.FM_TRIGGER_FILE, "replacement-successor actionable outcome\n"); -await waitFor(() => liveArms().length === 0, "mid-delivery successor actionable close"); +await waitFor(() => liveArms().length === 1 && armRows().length >= 3, "mid-delivery successor restored after second actionable close"); await previous.handlers.get("session_shutdown")?.({ type: "session_shutdown", reason: "new" }, {}); await waitFor(() => liveArms().length === 0, "retired old-session successor"); @@ -2585,7 +2585,7 @@ const replacementStart = replacement.handlers.get("session_start")?.({ previousSessionFile: "/tmp/previous.jsonl", }, {}); await new Promise((resolve) => setTimeout(resolve, 50)); -await waitFor(() => liveArms().length === 1 && armRows().length >= 3, "replacement arm before old delivery settlement"); +await waitFor(() => liveArms().length === 1 && armRows().length >= 4, "replacement arm before old delivery settlement"); if (replacement.prompts.some((message) => message.includes("signal: replacement-race actionable outcome"))) { throw new Error(`replacement raced the accepted old-session delivery: ${replacement.prompts.join(" | ")}`); } @@ -2612,7 +2612,7 @@ await new Promise((resolve) => setTimeout(resolve, 700)); if (replacement.prompts.filter((message) => message.includes("could not clear a delivered replacement-session actionable wake")).length !== 1) { throw new Error(`persistent handoff cleanup failure repeated alerts: ${replacement.prompts.join(" | ")}`); } -await waitFor(() => liveArms().length === 1 && armRows().length >= 3, "replacement live arm"); +await waitFor(() => liveArms().length === 1 && armRows().length >= 4, "replacement live arm"); const redundant = await replacement.getTool().execute("replacement-redundant", {}, undefined, undefined, {}); if (!redundant.details?.ok || !String(redundant.details.message).includes("unchanged")) { throw new Error(`replacement did not retain automatic arm ownership: ${JSON.stringify(redundant.details)}`); @@ -2966,6 +2966,156 @@ EOF pass "Pi retries a verified successor that failed during wake delivery once that delivery settles" } +test_pi_hung_settlement_later_cycles_restore_successor() { + local repo home plugin log marker_root trigger stop out status + repo="$TMP_ROOT/pi-hung-settlement-second-cycle-root" + home="$TMP_ROOT/pi-hung-settlement-second-cycle-home" + log="$TMP_ROOT/pi-hung-settlement-second-cycle.log" + marker_root="$TMP_ROOT/pi-hung-settlement-second-cycle-markers" + trigger="$TMP_ROOT/pi-hung-settlement-second-cycle.trigger" + stop="$TMP_ROOT/pi-hung-settlement-second-cycle.stop" + mkdir -p "$repo/bin" "$home/state" "$home/config" "$marker_root" + 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 +if [ "${1:-}" = --handling-delivered ]; then + printf 'confirmed generation=%s watcher=%s\n' "$2" "$4" >> "${FM_ARM_LOG:?}" + exit 0 +fi +marker=$(mktemp "${FM_MARKER_ROOT:?}/arm.XXXXXX") || exit 1 +cleanup() { rm -f "$marker"; } +trap cleanup EXIT +trap 'exit 0' TERM INT +printf 'arm pid=%s marker=%s\n' "$$" "$marker" >> "${FM_ARM_LOG:?}" +printf 'watcher: started pid=%s (beacon fresh) recovery-generation=hung-settlement-fixture\n' "$$" +while :; do + if [ -e "$FM_TRIGGER_FILE" ]; then + outcome=$(cat "$FM_TRIGGER_FILE") + rm -f "$FM_TRIGGER_FILE" + printf 'signal: ' + sleep 0.02 + printf '%s\n' "$outcome" + exit 0 + fi + [ ! -e "$FM_STOP_FILE" ] || exit 0 + 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_MARKER_ROOT="$marker_root" FM_TRIGGER_FILE="$trigger" FM_STOP_FILE="$stop" node --input-type=module 2>&1 <<'EOF' +import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { pathToFileURL } from "node:url"; + +let deliveryStarted = false; + +function makePi() { + const eventHandlers = new Map(); + let tool = null; + const prompts = []; + const pi = { + on() {}, + registerCommand() {}, + registerTool(candidate) { + if (candidate.name === "fm_watch_arm_pi") tool = candidate; + }, + sendUserMessage: async (message) => { + prompts.push(message); + }, + events: { + on(event, handler) { + eventHandlers.set(event, [...(eventHandlers.get(event) ?? []), handler]); + }, + emit(event, data) { + if (event === "fm-branch-supervision:dispatch") { + deliveryStarted = true; + data.accept(new Promise(() => {})); + } + for (const handler of eventHandlers.get(event) ?? []) handler(data); + }, + }, + }; + return { pi, getTool: () => tool, prompts }; +} + +function pidAlive(pid) { + try { + process.kill(Number(pid), 0); + return true; + } catch { + return false; + } +} + +function armRows() { + if (!existsSync(process.env.FM_ARM_LOG)) return []; + return readFileSync(process.env.FM_ARM_LOG, "utf8") + .trim() + .split(/\n/) + .filter((row) => row.startsWith("arm ")) + .map((row) => { + const match = /pid=(\d+) marker=(\S+)/.exec(row); + return match ? { pid: match[1], marker: match[2] } : { pid: "", marker: "" }; + }); +} + +function liveArms() { + return armRows().filter((arm) => arm.pid && arm.marker && existsSync(arm.marker) && pidAlive(arm.pid)); +} + +async function waitFor(pred, label, attempts = 500) { + for (let i = 0; i < attempts; i += 1) { + if (pred()) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + throw new Error(`timeout waiting for ${label}: live=${liveArms().length} rows=${armRows().length}`); +} + +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +writeFileSync(`${process.env.FM_HOME}/state/hung-cycle.meta`, "project=/projects/hung-cycle\nwindow=fm-hung-cycle\n"); +writeFileSync(`${process.env.FM_HOME}/state/.wake-queue`, "1\t1\tsignal\thung-cycle.status\tsignal: hung-cycle first outcome\n"); +const mod = await import(pathToFileURL(process.env.PLUGIN).href); +const session = makePi(); +mod.default(session.pi); +const initial = await session.getTool().execute("initial-arm", {}, undefined, undefined, {}); +if (!initial.details?.ok || !String(initial.details.message).includes("started Pi extension arm child")) { + throw new Error(`initial arm failed: ${JSON.stringify(initial.details)}`); +} +await waitFor(() => liveArms().length === 1, "initial live arm"); + +writeFileSync(process.env.FM_TRIGGER_FILE, "hung-cycle first outcome\n"); +await waitFor(() => deliveryStarted, "branch accepted hung first delivery"); +await waitFor(() => liveArms().length === 1 && armRows().length >= 2, "first successor during hung settlement"); +if (session.prompts.length !== 0) { + throw new Error(`hung first wake reached main: ${session.prompts.join(" | ")}`); +} + +writeFileSync(process.env.FM_TRIGGER_FILE, "hung-cycle second outcome\n"); +await waitFor(() => liveArms().length === 1 && armRows().length >= 3, "second successor during hung settlement"); + +writeFileSync(process.env.FM_TRIGGER_FILE, "hung-cycle third outcome\n"); +await waitFor(() => liveArms().length === 1 && armRows().length >= 4, "third successor during hung settlement"); + +const liveBeforeHold = liveArms().length; +const rowsBeforeHold = armRows().length; +await new Promise((resolve) => setTimeout(resolve, 200)); +if (liveArms().length !== 1 || armRows().length !== rowsBeforeHold) { + throw new Error(`unattended hold lost the restored successor: live=${liveArms().length} rows=${armRows().length} before=${liveBeforeHold}/${rowsBeforeHold}`); +} +if (session.prompts.length !== 0) { + throw new Error(`hung later wakes reached main: ${session.prompts.join(" | ")}`); +} + +writeFileSync(process.env.FM_STOP_FILE, "stop\n"); +process.exit(0); +EOF +) + status=$? + expect_code 0 "$status" "Pi hung settlement must still restore later-cycle successors" + [ -z "$out" ] || fail "Pi hung-settlement later-cycle test printed output: $out" + pass "Pi restores later-cycle successors while an earlier branch settlement is still hung" +} + test_pi_late_retiring_actionable_reaches_replacement() { local repo home plugin count out status repo="$TMP_ROOT/pi-late-retiring-actionable-root" @@ -4946,6 +5096,7 @@ test_pi_streaming_followup_is_replayed_after_replacement test_pi_streaming_followup_is_replayed_after_replacement away test_pi_streaming_time_delivery_keeps_the_successor_chain test_pi_successor_failure_during_delivery_is_retried_after_delivery +test_pi_hung_settlement_later_cycles_restore_successor test_pi_late_retiring_actionable_reaches_replacement test_pi_replacement_tokens_are_process_unique test_pi_replacement_persistence_failure_stops_arm_child From 23af31e31a0dfa7b3201e06805afa99fabec9148 Mon Sep 17 00:00:00 2001 From: Sungin Kim Date: Sat, 12 Sep 2026 00:21:23 +0000 Subject: [PATCH 2/4] no-mistakes(review): Preserve concurrent restore results and surface monitoring failures --- .pi/extensions/fm-primary-pi-watch.ts | 34 +++++++-------- tests/fm-pi-watch-extension.test.sh | 61 +++++++++++++++++++++++---- 2 files changed, 69 insertions(+), 26 deletions(-) diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index d90f6d23b35..d1dc299c114 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -32,7 +32,6 @@ // replacement handoff. // // Restore versus delivery (stated once here): -// restoring is true only while a successor arm is started and verified. // delivering is true while the serialized pending-wake pump is in flight. // A later actionable close restores a successor even when the previous wake's // branch settlement has not resolved. Failure closes during delivering still @@ -98,6 +97,11 @@ type UnconsumedWake = { pending: PendingActionableClose; }; +type ContinuityRestoration = { + failure: string; + recovery?: { generation: string; watcherPid: string }; +}; + type SessionGeneration = { id: number; stopping: boolean; @@ -106,7 +110,7 @@ type SessionGeneration = { retryTimer: ReturnType | null; cleanupTimer: ReturnType | null; retryFailures: number; - restoring: boolean; + restoring: Promise | null; delivering: boolean; seq: number; awayStandby: boolean; @@ -442,7 +446,7 @@ function createGeneration(): SessionGeneration { retryTimer: null, cleanupTimer: null, retryFailures: 0, - restoring: false, + restoring: null, delivering: false, seq: 0, awayStandby: false, @@ -790,17 +794,13 @@ export default function (pi: ExtensionAPI) { owner.cleanupTimer = timer; } - async function restoreContinuity(owner: SessionGeneration, predecessorArmPid: string): Promise<{ - failure: string; - recovery?: { generation: string; watcherPid: string }; - }> { - if (owner.restoring) return { failure: "" }; - owner.restoring = true; - try { - return await restoreAfterActionableClose(owner, predecessorArmPid); - } finally { - owner.restoring = false; + async function restoreContinuity(owner: SessionGeneration, predecessorArmPid: string): Promise { + if (!owner.restoring) { + owner.restoring = restoreAfterActionableClose(owner, predecessorArmPid).finally(() => { + owner.restoring = null; + }); } + return await owner.restoring; } async function processPendingActionables(owner: SessionGeneration): Promise { @@ -812,7 +812,8 @@ export default function (pi: ExtensionAPI) { if (owner.delivering) { const newest = owner.pendingActionables[owner.pendingActionables.length - 1]; if (newest && !newest.delivered && !owner.unconsumedWakes.has(newest.token)) { - void restoreContinuity(owner, newest.predecessorArmPid); + const restoration = await restoreContinuity(owner, newest.predecessorArmPid); + if (restoration.failure) surfaceFailure(owner, restoration.failure); } return; } @@ -1007,10 +1008,7 @@ export default function (pi: ExtensionAPI) { }); } - async function restoreAfterActionableClose(owner: SessionGeneration, predecessorArmPid: string): Promise<{ - failure: string; - recovery?: { generation: string; watcherPid: string }; - }> { + async function restoreAfterActionableClose(owner: SessionGeneration, predecessorArmPid: string): Promise { let failure = ""; for (let attempt = 0; attempt <= retryLimit; attempt += 1) { if (!generationIsLive(owner)) return { failure: "" }; diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index c716f4a6425..e952abcafc0 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -2968,12 +2968,13 @@ EOF test_pi_hung_settlement_later_cycles_restore_successor() { local repo home plugin log marker_root trigger stop out status - repo="$TMP_ROOT/pi-hung-settlement-second-cycle-root" - home="$TMP_ROOT/pi-hung-settlement-second-cycle-home" - log="$TMP_ROOT/pi-hung-settlement-second-cycle.log" - marker_root="$TMP_ROOT/pi-hung-settlement-second-cycle-markers" - trigger="$TMP_ROOT/pi-hung-settlement-second-cycle.trigger" - stop="$TMP_ROOT/pi-hung-settlement-second-cycle.stop" + local mode=${1:-consecutive} + repo="$TMP_ROOT/pi-hung-settlement-second-cycle-$mode-root" + home="$TMP_ROOT/pi-hung-settlement-second-cycle-$mode-home" + log="$TMP_ROOT/pi-hung-settlement-second-cycle-$mode.log" + marker_root="$TMP_ROOT/pi-hung-settlement-second-cycle-$mode-markers" + trigger="$TMP_ROOT/pi-hung-settlement-second-cycle-$mode.trigger" + stop="$TMP_ROOT/pi-hung-settlement-second-cycle-$mode.stop" mkdir -p "$repo/bin" "$home/state" "$home/config" "$marker_root" install_pi_watch_extension_fixture "$repo" plugin="$repo/.pi/extensions/fm-primary-pi-watch.ts" @@ -2988,6 +2989,15 @@ cleanup() { rm -f "$marker"; } trap cleanup EXIT trap 'exit 0' TERM INT printf 'arm pid=%s marker=%s\n' "$$" "$marker" >> "${FM_ARM_LOG:?}" +count=$(find "$FM_MARKER_ROOT" -name 'attempt.*' | wc -l) +: > "$FM_MARKER_ROOT/attempt.$$" +if [ "$count" -ge 2 ] && [ "$FM_RESTORE_MODE" != consecutive ]; then + while [ ! -e "$FM_MARKER_ROOT/release" ]; do sleep 0.02; done + if [ "$FM_RESTORE_MODE" != slow-success ]; then + printf 'watcher: FAILED - fixture successor failed\n' + exit 1 + fi +fi printf 'watcher: started pid=%s (beacon fresh) recovery-generation=hung-settlement-fixture\n' "$$" while :; do if [ -e "$FM_TRIGGER_FILE" ]; then @@ -3003,11 +3013,13 @@ while :; do 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_MARKER_ROOT="$marker_root" FM_TRIGGER_FILE="$trigger" FM_STOP_FILE="$stop" node --input-type=module 2>&1 <<'EOF' + out=$(FM_RESTORE_MODE="$mode" FM_WATCH_REARM_RETRY_LIMIT=1 FM_WATCH_REARM_RETRY_BASE_MS=10 FM_PI_ARM_READY_TIMEOUT_MS=5000 PLUGIN="$plugin" FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_ARM_LOG="$log" FM_MARKER_ROOT="$marker_root" FM_TRIGGER_FILE="$trigger" FM_STOP_FILE="$stop" node --input-type=module 2>&1 <<'EOF' import { existsSync, readFileSync, writeFileSync } from "node:fs"; import { pathToFileURL } from "node:url"; let deliveryStarted = false; +let finishFirst; +let offers = 0; function makePi() { const eventHandlers = new Map(); @@ -3029,7 +3041,8 @@ function makePi() { emit(event, data) { if (event === "fm-branch-supervision:dispatch") { deliveryStarted = true; - data.accept(new Promise(() => {})); + offers += 1; + data.accept(offers === 1 ? new Promise((resolve) => { finishFirst = resolve; }) : Promise.resolve()); } for (const handler of eventHandlers.get(event) ?? []) handler(data); }, @@ -3093,6 +3106,35 @@ if (session.prompts.length !== 0) { writeFileSync(process.env.FM_TRIGGER_FILE, "hung-cycle second outcome\n"); await waitFor(() => liveArms().length === 1 && armRows().length >= 3, "second successor during hung settlement"); +if (process.env.FM_RESTORE_MODE !== "consecutive") { + const mode = process.env.FM_RESTORE_MODE; + if (mode.startsWith("slow-")) { + finishFirst(); + await new Promise((resolve) => setTimeout(resolve, 150)); + if (offers !== 1 || session.prompts.length !== 0) { + throw new Error("later wake delivered before the in-flight restore settled"); + } + } + writeFileSync(`${process.env.FM_MARKER_ROOT}/release`, ""); + if (mode === "slow-success") { + await waitFor(() => offers === 2, "second wake delivered after readiness"); + const successor = armRows()[2]; + const confirmations = readFileSync(process.env.FM_ARM_LOG, "utf8").split("\n") + .filter((row) => row.startsWith("confirmed ") && row.includes(`watcher=${successor.pid}`)); + if (confirmations.length !== 1) throw new Error("shared restore lost its recovery confirmation"); + } else { + await waitFor(() => session.prompts.some((message) => message.includes("after 1 retries")), "typed side restore failure"); + if (mode === "slow-failure") { + await waitFor(() => session.prompts.some((message) => message.includes("second outcome") && message.includes("after 1 retries")), "later wake annotated with shared restore failure"); + } else if (offers !== 1 || session.prompts.some((message) => message.includes("second outcome"))) { + throw new Error("failure reporting delivered a later wake during hung settlement"); + } + if (armRows().length !== 4) throw new Error(`restore did not honor retry bound: ${armRows().length}`); + } + writeFileSync(process.env.FM_STOP_FILE, "stop\n"); + process.exit(0); +} + writeFileSync(process.env.FM_TRIGGER_FILE, "hung-cycle third outcome\n"); await waitFor(() => liveArms().length === 1 && armRows().length >= 4, "third successor during hung settlement"); @@ -5097,6 +5139,9 @@ test_pi_streaming_followup_is_replayed_after_replacement away test_pi_streaming_time_delivery_keeps_the_successor_chain test_pi_successor_failure_during_delivery_is_retried_after_delivery test_pi_hung_settlement_later_cycles_restore_successor +test_pi_hung_settlement_later_cycles_restore_successor held-failure +test_pi_hung_settlement_later_cycles_restore_successor slow-success +test_pi_hung_settlement_later_cycles_restore_successor slow-failure test_pi_late_retiring_actionable_reaches_replacement test_pi_replacement_tokens_are_process_unique test_pi_replacement_persistence_failure_stops_arm_child From 690534f6cf8599a43fea8eda68fd523828adffd9 Mon Sep 17 00:00:00 2001 From: Sungin Kim Date: Sat, 12 Sep 2026 00:30:37 +0000 Subject: [PATCH 3/4] no-mistakes(review): Catch side-path restore exceptions without interrupting hung delivery --- .pi/extensions/fm-primary-pi-watch.ts | 9 +++++-- tests/fm-pi-watch-extension.test.sh | 36 +++++++++++++++++++++++---- 2 files changed, 38 insertions(+), 7 deletions(-) diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index d1dc299c114..ad1fc8ba7d1 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -812,8 +812,13 @@ export default function (pi: ExtensionAPI) { if (owner.delivering) { const newest = owner.pendingActionables[owner.pendingActionables.length - 1]; if (newest && !newest.delivered && !owner.unconsumedWakes.has(newest.token)) { - const restoration = await restoreContinuity(owner, newest.predecessorArmPid); - if (restoration.failure) surfaceFailure(owner, restoration.failure); + try { + const restoration = await restoreContinuity(owner, newest.predecessorArmPid); + if (restoration.failure) surfaceFailure(owner, restoration.failure); + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + surfaceFailure(owner, `watcher: FAILED - Pi extension could not restore watcher continuity\n${detail}`); + } } return; } diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index e952abcafc0..fcf1b76d54b 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -3013,8 +3013,8 @@ while :; do done SH chmod +x "$repo/bin/fm-watch-arm.sh" - out=$(FM_RESTORE_MODE="$mode" FM_WATCH_REARM_RETRY_LIMIT=1 FM_WATCH_REARM_RETRY_BASE_MS=10 FM_PI_ARM_READY_TIMEOUT_MS=5000 PLUGIN="$plugin" FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_ARM_LOG="$log" FM_MARKER_ROOT="$marker_root" FM_TRIGGER_FILE="$trigger" FM_STOP_FILE="$stop" node --input-type=module 2>&1 <<'EOF' -import { existsSync, readFileSync, writeFileSync } from "node:fs"; + out=$(FM_RESTORE_MODE="$mode" FM_WATCH_REARM_RETRY_LIMIT=1 FM_WATCH_REARM_RETRY_BASE_MS=10 FM_PI_ARM_READY_TIMEOUT_MS=5000 PLUGIN="$plugin" FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_ARM_LOG="$log" FM_MARKER_ROOT="$marker_root" FM_TRIGGER_FILE="$trigger" FM_STOP_FILE="$stop" node --unhandled-rejections=strict --input-type=module 2>&1 <<'EOF' +import { existsSync, mkdirSync, readFileSync, unlinkSync, writeFileSync } from "node:fs"; import { pathToFileURL } from "node:url"; let deliveryStarted = false; @@ -3023,10 +3023,11 @@ let offers = 0; function makePi() { const eventHandlers = new Map(); + const handlers = new Map(); let tool = null; const prompts = []; const pi = { - on() {}, + on(event, handler) { handlers.set(event, handler); }, registerCommand() {}, registerTool(candidate) { if (candidate.name === "fm_watch_arm_pi") tool = candidate; @@ -3048,7 +3049,7 @@ function makePi() { }, }, }; - return { pi, getTool: () => tool, prompts }; + return { pi, getTool: () => tool, prompts, handlers }; } function pidAlive(pid) { @@ -3103,7 +3104,31 @@ if (session.prompts.length !== 0) { throw new Error(`hung first wake reached main: ${session.prompts.join(" | ")}`); } +if (process.env.FM_RESTORE_MODE === "held-exception") { + const marker = `${process.env.FM_HOME}/state/.pi-watch-extension-loaded`; + unlinkSync(marker); + mkdirSync(marker); +} + writeFileSync(process.env.FM_TRIGGER_FILE, "hung-cycle second outcome\n"); +if (process.env.FM_RESTORE_MODE === "held-exception") { + await waitFor(() => session.prompts.some((message) => + message.includes("watcher: FAILED - Pi extension could not restore watcher continuity") && message.includes("EISDIR")), + "typed side restore exception while first delivery remains hung"); + await new Promise((resolve) => setTimeout(resolve, 150)); + if (offers !== 1 || session.prompts.length !== 1 || session.prompts[0].includes("second outcome")) { + throw new Error("restore exception delivered a later wake during hung settlement"); + } + await session.handlers.get("session_shutdown")({ type: "session_shutdown", reason: "new" }, {}); + const handoff = JSON.parse(readFileSync(`${process.env.FM_HOME}/state/extensions/pi-primary-watch/session-replacement-actionable.json`, "utf8")); + if (handoff.pending.length !== 2 || handoff.pending.some((pending) => pending.delivered)) { + throw new Error(`restore exception lost pending wakes: ${JSON.stringify(handoff)}`); + } + if (armRows().length !== 2 || liveArms().length !== 0) throw new Error("restore exception launched an unexpected successor"); + writeFileSync(process.env.FM_STOP_FILE, "stop\n"); + process.exit(0); +} + await waitFor(() => liveArms().length === 1 && armRows().length >= 3, "second successor during hung settlement"); if (process.env.FM_RESTORE_MODE !== "consecutive") { @@ -3153,7 +3178,7 @@ process.exit(0); EOF ) status=$? - expect_code 0 "$status" "Pi hung settlement must still restore later-cycle successors" + expect_code 0 "$status" "Pi hung settlement must still restore later-cycle successors ($mode): $out" [ -z "$out" ] || fail "Pi hung-settlement later-cycle test printed output: $out" pass "Pi restores later-cycle successors while an earlier branch settlement is still hung" } @@ -5140,6 +5165,7 @@ test_pi_streaming_time_delivery_keeps_the_successor_chain test_pi_successor_failure_during_delivery_is_retried_after_delivery test_pi_hung_settlement_later_cycles_restore_successor test_pi_hung_settlement_later_cycles_restore_successor held-failure +test_pi_hung_settlement_later_cycles_restore_successor held-exception test_pi_hung_settlement_later_cycles_restore_successor slow-success test_pi_hung_settlement_later_cycles_restore_successor slow-failure test_pi_late_retiring_actionable_reaches_replacement From d846ae520a1acc80b1df3389666e53f895f25d53 Mon Sep 17 00:00:00 2001 From: Sungin Kim Date: Sat, 12 Sep 2026 00:43:42 +0000 Subject: [PATCH 4/4] no-mistakes(document): state shared restore and side-path failure invariants --- .pi/extensions/fm-primary-pi-watch.ts | 9 +++++++++ tests/fm-pi-watch-extension.test.sh | 2 +- 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index ad1fc8ba7d1..cff3962f9d3 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -36,6 +36,15 @@ // A later actionable close restores a successor even when the previous wake's // branch settlement has not resolved. Failure closes during delivering still // defer their bounded retry until that delivery settles. +// restoring holds the one in-flight restore promise for the generation, so a +// concurrent caller awaits that restore's real result instead of reading a +// synthetic success: the pump needs the recovery marker to confirm handling +// and the failure text to annotate the wake it delivers. +// Because that side-path restore runs outside the pump's own try/catch, it +// surfaces its typed failure - and any thrown persistence error - on its own +// rather than waiting for the hung delivery to settle. A shared failing +// restore therefore reaches main twice, standalone and appended to the later +// wake; re-arming is idempotent, and silent stopped monitoring is not. import { spawn, spawnSync, type ChildProcess } from "node:child_process"; import { createHash } from "node:crypto"; import { existsSync, mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs"; diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index fcf1b76d54b..a466bf011f6 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -3180,7 +3180,7 @@ EOF status=$? expect_code 0 "$status" "Pi hung settlement must still restore later-cycle successors ($mode): $out" [ -z "$out" ] || fail "Pi hung-settlement later-cycle test printed output: $out" - pass "Pi restores later-cycle successors while an earlier branch settlement is still hung" + pass "Pi restores later-cycle successors while an earlier branch settlement is still hung ($mode)" } test_pi_late_retiring_actionable_reaches_replacement() {