diff --git a/.omp/extensions/fm-primary-omp-watch.ts b/.omp/extensions/fm-primary-omp-watch.ts index 44ae7b81178..f401f3942bf 100644 --- a/.omp/extensions/fm-primary-omp-watch.ts +++ b/.omp/extensions/fm-primary-omp-watch.ts @@ -43,7 +43,7 @@ // record, and a still-unconsumed doorbell rides the replacement handoff. import { spawn, spawnSync, type ChildProcess } from "node:child_process"; import { createHash } from "node:crypto"; -import { mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs"; +import { existsSync, mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs"; import { dirname, resolve } from "node:path"; import { fileURLToPath } from "node:url"; // typebox resolves inside omp's extension loader (verified, omp 18.1.11); the @@ -327,7 +327,6 @@ function persistReplacementHandoff(pending: PendingActionableClose[]): void { if (pending.length === 0) return; writeReplacementHandoff(pending); } - function loadReplacementHandoff(): PendingActionableClose[] { try { const pending = validateReplacementHandoff(JSON.parse(readFileSync(actionableHandoff, "utf8"))); @@ -342,6 +341,55 @@ function loadReplacementHandoff(): PendingActionableClose[] { } } +function loadAndFilterReplacementHandoff(): PendingActionableClose[] { + try { + const pending = validateReplacementHandoff(JSON.parse(readFileSync(actionableHandoff, "utf8"))); + const filtered = pending.filter((item) => isTaskAlive(item.message)); + if (filtered.length !== pending.length) { + writeReplacementHandoff(filtered); + } + replacementHandoff = filtered; + return [...filtered]; + } catch (error) { + if (nodeErrorCode(error) === "ENOENT") { + replacementHandoff = null; + return []; + } + throw error; + } +} + +function isTaskAlive(message: string): boolean { + const taskId = extractTaskId(message); + if (!taskId) return true; + const stateDir = state; + const wakeQueue = `${stateDir}/.wake-queue`; + const candidates = [taskId]; + if (taskId.startsWith("fm-")) candidates.push(taskId.slice(3)); + const hasStateFile = candidates.some((id) => + existsSync(`${stateDir}/${id}.meta`) || + existsSync(`${stateDir}/${id}.status`) || + existsSync(`${stateDir}/${id}.inbox`) || + existsSync(`${stateDir}/${id}.progress`) || + existsSync(`${stateDir}/${id}.turn-ended`), + ); + if (hasStateFile) return true; + if (!existsSync(wakeQueue)) return false; + const queueContent = readFileSync(wakeQueue, "utf8"); + const lines = queueContent.trim().split("\n").filter((l) => l.length > 0); + for (const line of lines) { + const parts = line.split("\t"); + if (parts.length >= 4 && candidates.includes(parts[3])) return true; + } + return false; +} + +function extractTaskId(message: string): string | null { + const match = message.match(/^(?:stale|signal|check|heartbeat):\s*([^\s:\.]+)/); + if (match) return match[1]; + return null; +} + function mergeReplacementHandoff(pending: PendingActionableClose): void { let stored: PendingActionableClose[] = []; try { @@ -1030,7 +1078,7 @@ export default function (pi: ExtensionAPI) { let pending: PendingActionableClose[] = []; let loadFailure = ""; try { - pending = loadReplacementHandoff(); + pending = loadAndFilterReplacementHandoff(); } catch (error) { const detail = error instanceof Error ? error.message : String(error); loadFailure = `watcher: FAILED - omp extension could not load a replacement-session actionable wake\n${detail}`; diff --git a/tests/fm-omp-harness.test.sh b/tests/fm-omp-harness.test.sh index 125fe19ac78..510376d3197 100755 --- a/tests/fm-omp-harness.test.sh +++ b/tests/fm-omp-harness.test.sh @@ -728,12 +728,14 @@ SH import { pathToFileURL } from "node:url"; import { mkdirSync, writeFileSync, existsSync, readFileSync } from "node:fs"; writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +mkdirSync(`${process.env.FM_HOME}/state`, { recursive: true }); +writeFileSync(`${process.env.FM_HOME}/state/default.meta`, "task_id=default\n"); const handoffDir = `${process.env.FM_HOME}/state/extensions/omp-primary-watch`; const handoff = `${handoffDir}/session-replacement-actionable.json`; mkdirSync(handoffDir, { recursive: true }); writeFileSync(handoff, `${JSON.stringify({ version: 2, - pending: [1, 2, 3, 4].map((n) => ({ + pending: [1,2,3,4].map((n) => ({ version: 1, token: `900-1000-${n}`, message: `stale: default:wR:p${n}`, @@ -956,6 +958,155 @@ EOF [ -z "$out" ] || fail "omp watch extension genuine-wake test printed output: $out" pass ".omp watch extension: a genuine wake between failed cycles clears the consecutive-failure count" } +# --- Stale replacement handoff records are finalized and dropped ----------- + +test_watch_extension_drops_stale_replacement_records() { + local repo home out status + repo="$TMP_ROOT/watch-stale/repo"; home="$TMP_ROOT/watch-stale/home" + install_omp_extension_fixture "$repo" + mkdir -p "$home/state" + cat > "$repo/bin/fm-watch-arm.sh" <<'SH' +#!/usr/bin/env bash +if [ "${1:-}" = --handling-delivered ]; then + exit 0 +fi +state=${FM_HOME:?}/state +printf 'watcher: started pid=%s (beacon 0s) recovery-generation=gen-1\n' "$$" +exec sleep 30 +SH + chmod +x "$repo/bin/fm-watch-arm.sh" + out=$(FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_OMP_ARM_READY_TIMEOUT_MS=5000 \ + EXT="$repo/.omp/extensions/fm-primary-omp-watch.ts" node --input-type=module 2>&1 <<'EOF' +import { pathToFileURL } from "node:url"; +import { mkdirSync, writeFileSync, existsSync, readFileSync, unlinkSync } from "node:fs"; +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +const handoffDir = `${process.env.FM_HOME}/state/extensions/omp-primary-watch`; +const handoff = `${handoffDir}/session-replacement-actionable.json`; +mkdirSync(handoffDir, { recursive: true }); +writeFileSync(handoff, `${JSON.stringify({ + version: 2, + pending: [ + { version: 1, token: "1000-2000-1", message: "stale: ghost-task-1:wake-a", predecessorArmPid: "1111" }, + { version: 1, token: "1000-2000-2", message: "signal: ghost-task-2:wake-b", predecessorArmPid: "2222" }, + ], +})}\n`); +const waitUntil = async (predicate, label, timeoutMs = 5000) => { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) return; + await new Promise((resolve) => setTimeout(resolve, 25)); + } + throw new Error(`timed out waiting for ${label}`); +}; +const handlers = new Map(); let tool = null; const sent = []; +const pi = { + on(e, h) { handlers.set(e, h); }, + registerCommand() {}, + registerTool(t) { tool = t; }, + sendUserMessage(m, o) { sent.push({ m, o }); return undefined; }, +}; +const mod = await import(pathToFileURL(process.env.EXT).href); +mod.default(pi); +await tool.execute(); +await waitUntil(() => existsSync(handoff), "handoff file to be rewritten"); +const filtered = JSON.parse(readFileSync(handoff, "utf8")); +if (filtered.pending.length !== 0) throw new Error(`stale records should be dropped, got ${JSON.stringify(filtered)}`); +if (sent.length !== 0) throw new Error(`no follow-up should be sent for stale records, got ${JSON.stringify(sent)}`); +writeFileSync(`${process.env.FM_HOME}/state/valid-task.meta`, "task_id=valid-task\n"); +const handoff2 = `${handoffDir}/session-replacement-actionable.json`; +writeFileSync(handoff2, `${JSON.stringify({ + version: 2, + pending: [ + { version: 1, token: "1000-2000-3", message: "stale: valid-task:wake-c", predecessorArmPid: "3333" }, + ], +})}\n`); +await waitUntil(() => existsSync(handoff2), "second handoff to be processed"); +const filtered2 = JSON.parse(readFileSync(handoff2, "utf8")); +if (filtered2.pending.length !== 1) throw new Error(`valid record should remain, got ${JSON.stringify(filtered2)}`); +if (filtered2.pending[0].token !== "1000-2000-3") throw new Error(`wrong token: ${filtered2.pending[0].token}`); +process.exit(0); +EOF +) + status=$? + expect_code 0 "$status" "omp watch extension stale record filter: $out" + [ -z "$out" ] || fail "omp watch extension stale record test printed output: $out" + pass ".omp watch extension: stale replacement records without backing state or queue are dropped" +} + +test_watch_extension_consumption_clears_handoff() { + local repo home out status + repo="$TMP_ROOT/watch-consumption/repo"; home="$TMP_ROOT/watch-consumption/home" + install_omp_extension_fixture "$repo" + mkdir -p "$home/state" + cat > "$repo/bin/fm-watch-arm.sh" <<'SH' +#!/usr/bin/env bash +if [ "${1:-}" = --handling-delivered ]; then + exit 0 +fi +state=${FM_HOME:?}/state +count=$(cat "$state/.arm-count" 2>/dev/null || printf 0) +count=$((count + 1)) +printf '%s\n' "$count" > "$state/.arm-count" +printf 'watcher: started pid=%s (beacon 0s) recovery-generation=gen-%s\n' "$$" "$count" +if [ "$count" -eq 1 ]; then + while [ ! -e "$state/.release-first-arm" ]; do sleep 0.05; done + printf 'stale: valid-task:wake-2\n' + exit 0 +fi +exec sleep 30 +SH + chmod +x "$repo/bin/fm-watch-arm.sh" + out=$(FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_OMP_ARM_READY_TIMEOUT_MS=10000 FM_WATCH_REARM_RETRY_LIMIT=1 FM_WATCH_REARM_RETRY_BASE_MS=5 FM_WATCH_REARM_RETRY_MAX_MS=10 \ + EXT="$repo/.omp/extensions/fm-primary-omp-watch.ts" node --input-type=module 2>&1 <<'EOF' +import { pathToFileURL } from "node:url"; +import { mkdirSync, writeFileSync, existsSync, readFileSync } from "node:fs"; +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +mkdirSync(`${process.env.FM_HOME}/state`, { recursive: true }); +writeFileSync(`${process.env.FM_HOME}/state/valid-task.meta`, "task_id=valid-task\n"); +const handoffDir = `${process.env.FM_HOME}/state/extensions/omp-primary-watch`; +const handoff = `${handoffDir}/session-replacement-actionable.json`; +mkdirSync(handoffDir, { recursive: true }); +writeFileSync(handoff, `${JSON.stringify({ + version: 2, + pending: [ + { version: 1, token: "1000-2000-4", message: "stale: valid-task:wake-1", predecessorArmPid: "4444" }, + ], +})}\n`); +const waitUntil = async (predicate, label, timeoutMs = 10000) => { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) return; + await new Promise((resolve) => setTimeout(resolve, 25)); + } + throw new Error(`timed out waiting for ${label}`); +}; +const handlers = new Map(); let tool = null; const sent = []; +const pi = { + on(e, h) { handlers.set(e, h); }, + registerCommand() {}, + registerTool(t) { tool = t; }, + sendUserMessage(m, o) { sent.push({ m, o }); return undefined; }, +}; +const mod = await import(pathToFileURL(process.env.EXT).href); +mod.default(pi); +await tool.execute(); +await waitUntil(() => sent.length >= 1, "replacement follow-up delivered"); +if (sent.length !== 1) throw new Error(`expected one follow-up, got ${sent.length}`); +if (!sent[0].m.includes("stale: valid-task:wake-1")) throw new Error(`unexpected wake: ${sent[0].m}`); +await handlers.get("before_agent_start")({ type: "before_agent_start", prompt: sent[0].m }, {}); +await waitUntil(() => !existsSync(handoff), "handoff cleared after consumption"); +writeFileSync(`${process.env.FM_HOME}/state/.release-first-arm`, "\n"); +await waitUntil(() => sent.length >= 2, "post-consumption actionable"); +if (sent.length !== 2) throw new Error(`expected second follow-up, got ${sent.length}`); +if (existsSync(handoff)) throw new Error("handoff should not be recreated for post-consumption actionable"); +process.exit(0); +EOF +) + status=$? + expect_code 0 "$status" "omp watch extension consumption clears handoff: $out" + [ -z "$out" ] || fail "omp watch extension consumption test printed output: $out" + pass ".omp watch extension: consumption at before_agent_start clears replacement handoff so it is not replayed" +} test_detection_anchored_name_and_marker_precedence test_isolation_hides_a_genuine_omp_ancestor @@ -974,3 +1125,5 @@ test_watch_extension_arms_and_delivers test_watch_extension_retry_arms_as_a_cold_start test_watch_extension_resurface_cycles_still_reach_the_retry_limit test_watch_extension_genuine_wake_clears_the_failure_count +test_watch_extension_drops_stale_replacement_records +test_watch_extension_consumption_clears_handoff