Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 51 additions & 3 deletions .omp/extensions/fm-primary-omp-watch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")));
Expand All @@ -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 {
Expand Down Expand Up @@ -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}`;
Expand Down
155 changes: 154 additions & 1 deletion tests/fm-omp-harness.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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}`,
Expand Down Expand Up @@ -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
Expand All @@ -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