From 22c8de712df08a3c2947347469ca6f4e27296c18 Mon Sep 17 00:00:00 2001 From: dnth Date: Thu, 3 Sep 2026 19:03:46 +0800 Subject: [PATCH 1/6] fix(watcher): preserve continuity across session replacement --- .omp/extensions/fm-branch-supervision-omp.ts | 54 +- .omp/extensions/fm-primary-omp.ts | 18 +- .omp/extensions/lib/fm-branch-dispatch.ts | 12 +- .pi/extensions/fm-primary-pi-watch.ts | 9 +- bin/fm-primary-watch-core.ts | 540 +++++++++++++++++-- docs/omp-supervision-branch.md | 8 +- docs/supervision-protocols/omp.md | 5 +- docs/supervision-protocols/pi.md | 6 +- docs/verification/supervision.md | 32 +- docs/watcher-continuity.md | 5 +- tests/fm-omp-primary.test.sh | 142 +++++ tests/fm-pi-primary-live-e2e.test.sh | 33 +- tests/fm-pi-watch-extension.test.sh | 140 ++++- 13 files changed, 854 insertions(+), 150 deletions(-) diff --git a/.omp/extensions/fm-branch-supervision-omp.ts b/.omp/extensions/fm-branch-supervision-omp.ts index c2e6b5eaab5..14842fac936 100644 --- a/.omp/extensions/fm-branch-supervision-omp.ts +++ b/.omp/extensions/fm-branch-supervision-omp.ts @@ -43,10 +43,11 @@ // process; and a secondary read-only OMP session that never owns the lock must // never write markers, clean leases, or accept wakes. // -// Failure direction: every path that cannot reach a working branch falls back -// to delivering the wake to MAIN exactly as before the branch existed - a -// broken branch degrades to today's behavior, never to a lost wake. The wake -// queue itself stays durable until the handler runs the drain's +// Failure direction: every accepted path that cannot reach a working branch +// rejects its settlement to the watcher, which retains delivery ownership and +// routes the wake to MAIN through its consumption-acknowledged path. A broken +// branch declines later offers, so they take that same watcher path directly. +// The wake queue itself stays durable until the handler runs the drain's // acknowledgement, so a branch that dies mid-handling re-presents its rows at // the next drain exactly as a mid-handling main crash always has. // @@ -845,30 +846,8 @@ ${context.command} } } - function fallbackToMain(message: string, detail: string): void { - const body = `FIRSTMATE WATCHER WAKE: ${message}\n\nRun bin/fm-wake-drain.sh first and handle the queued wake. (Supervision branch unavailable, falling back to main: ${detail})`; - // Marked operational like every watcher injection, so the wake is never - // mistaken for captain input (away-mode return semantics, mirror filter). - const content = encodeOperationalInput(body); - // Deliver through the exact main-wake mechanism the primary OMP adapter uses - // (fm-primary-omp.ts sendFollowUp): a custom watcher-wake message delivered - // as a steer with triggerTurn. Unlike sendUserMessage/followUp, this - // reliably wakes an idle or interrupted main under OMP steer/continuation - // semantics, so a broken branch never strands the wake. - pi.sendMessage( - { - customType: "firstmate-watcher-wake", - content, - display: false, - attribution: "agent", - details: { kind: "watcher", runtime: "omp" }, - }, - { deliverAs: "steer", triggerTurn: true }, - ); - } - - function enqueueWake(message: string, acceptedGeneration: number): void { - branchChain = branchChain + function enqueueWake(message: string, acceptedGeneration: number): Promise { + const delivery = branchChain .then(async () => { if (shuttingDown || acceptedGeneration !== generation) { throw new Error("supervision session was replaced before handling the accepted wake"); @@ -902,20 +881,12 @@ ${context.command} throw new Error("could not release the branch's settled wake-row grant"); } }) - .catch(async (error: unknown) => { + .catch((error: unknown) => { releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration)); - // Only fall the wake back to main when this generation still owns the - // fleet lock and is not shutting down. If the session lost ownership or a - // cold-start re-arm advanced the generation, the durable row stays queued - // for the owning session to reclaim - the row is never lost - and falling - // back here would risk a stale extra main turn for a row a fresh arm may - // already handle. - if (shuttingDown || acceptedGeneration !== generation) return; - if (!actingAsOwner(acceptedGeneration)) return; - try { - fallbackToMain(message, error instanceof Error ? error.message : String(error)); - } catch {} + throw error; }); + branchChain = delivery.catch(() => {}); + return delivery; } function enqueueMirrorFlush(): void { @@ -943,8 +914,7 @@ ${context.command} if (!actingAsOwner()) return; // cold start pre-lock, secondary session, or shutdown if (afkActive()) return; // the away daemon owns supervision while afk if (branchBroken) return; // fail back to today's wake-to-main path - offer.accept(); - enqueueWake(offer.message, generation); + offer.accept(enqueueWake(offer.message, generation)); }); pi.on?.("agent_start", () => { diff --git a/.omp/extensions/fm-primary-omp.ts b/.omp/extensions/fm-primary-omp.ts index d444774b705..1f48159a202 100644 --- a/.omp/extensions/fm-primary-omp.ts +++ b/.omp/extensions/fm-primary-omp.ts @@ -212,14 +212,14 @@ export default function (omp: ExtensionAPI) { // credential/auth failures) is never offered - it stays main-owned - and with // no branch extension loaded no one accepts, so every wake falls through to // main exactly as before. - const offerWakeToBranch = (message: string): boolean => { + const offerWakeToBranch = (message: string): Promise | null => { const heartbeat = /^heartbeat($|:)/.test(message); const isCheckTrigger = /^check:/.test(message); const scope = scopeForUnreadWake(state, heartbeat); const eligible = !isCheckTrigger && scope.eligible; const offer = createBranchDispatchOffer(message, scope.projects, heartbeat, eligible); omp.events?.emit?.(FM_BRANCH_DISPATCH_EVENT, offer); - return offer.accepted; + return offer.accepted ? offer.settlement : null; }; const watch = createPrimaryWatchCore({ @@ -266,15 +266,15 @@ export default function (omp: ExtensionAPI) { publishTaskTurnStarted(); }); - omp.on("session_switch", (event, ctx) => { - watch.sessionShutdown(); - watch.sessionStart(); + omp.on("session_switch", async (event, ctx) => { + await watch.sessionShutdown(true); publishSecondmateSession(ctx); deliverSessionstartNudge(event.reason === "new" || event.reason === "resume"); - watch.arm(); + watch.sessionStart(); }); - omp.on("before_agent_start", (): BeforeAgentStartEventResult | undefined => { + omp.on("before_agent_start", (event): BeforeAgentStartEventResult | undefined => { + watch.acknowledgeWake(event.prompt); if (!pendingStartupNudge) return undefined; const content = pendingStartupNudge; pendingStartupNudge = ""; @@ -322,9 +322,9 @@ export default function (omp: ExtensionAPI) { return {}; }); - omp.on("session_shutdown", () => { + omp.on("session_shutdown", async () => { taskInboxDoorbell.retire(); - watch.sessionShutdown(); + await watch.sessionShutdown(false); }); omp.registerCommand("fm-watch-arm-omp", { diff --git a/.omp/extensions/lib/fm-branch-dispatch.ts b/.omp/extensions/lib/fm-branch-dispatch.ts index f2d88f2fe35..3bcdaf9a5a3 100644 --- a/.omp/extensions/lib/fm-branch-dispatch.ts +++ b/.omp/extensions/lib/fm-branch-dispatch.ts @@ -9,8 +9,9 @@ import { readdirSync, readFileSync } from "node:fs"; // FM_BRANCH_DISPATCH_EVENT. A live, enabled branch extension calls accept() // SYNCHRONOUSLY inside its handler (the event bus invokes handlers // synchronously up to their first await), so after emit returns the watcher -// reads `accepted`: true means the branch now owns delivering and handling the -// wake (including its own fallback back to main on a later failure); false +// reads `accepted`: true means the branch owns handling the wake, and its +// settlement promise keeps the watcher outcome pending until handling finishes +// or rejects back to the watcher's consumption-acknowledged main path; false // means no branch took it and the watcher delivers to main exactly as it did // before the branch existed. Watcher-failure alarms are never offered - only // main can repair the watcher cycle (fm_watch_arm_pi lives on main). @@ -229,7 +230,8 @@ export interface BranchDispatchOffer { eligible: boolean; /** Set by accept(); read by the watcher after emit returns. */ accepted: boolean; - accept(): void; + settlement: Promise; + accept(settlement?: Promise): void; } export function createBranchDispatchOffer( @@ -244,8 +246,10 @@ export function createBranchDispatchOffer( heartbeat, eligible, accepted: false, - accept() { + settlement: Promise.resolve(), + accept(settlement = Promise.resolve()) { offer.accepted = true; + offer.settlement = settlement; }, }; return offer; diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index 309f43f4d2d..6cef65c16b4 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -89,11 +89,16 @@ export default function (pi: ExtensionAPI) { }, }); + pi.on?.("before_agent_start", (event) => { + watch.acknowledgeWake(event.prompt); + }); + pi.on?.("session_start", () => { watch.sessionStart(); }); - pi.on?.("session_shutdown", () => { - watch.sessionShutdown(); + pi.on?.("session_shutdown", async (event) => { + const replacement = ["reload", "new", "resume", "fork"].includes(String(event.reason ?? "")); + await watch.sessionShutdown(replacement); }); pi.registerCommand?.("fm-watch-arm-pi", { diff --git a/bin/fm-primary-watch-core.ts b/bin/fm-primary-watch-core.ts index 64b48694ace..4f1587abed9 100644 --- a/bin/fm-primary-watch-core.ts +++ b/bin/fm-primary-watch-core.ts @@ -4,12 +4,14 @@ // // Session-generation ownership (stated once here): one generation is bound per // runtime session activation. Only the active live generation may start, stop, -// rearm, or clear the arm child. A runtime-specific replacement event or a -// fresh factory bind activates a new live generation so monitoring can arm -// without restarting the process. Terminal shutdown leaves the final generation -// stopped so late callbacks cannot rearm. Stale callbacks from an earlier -// generation, including callbacks retained by a superseded core instance, are -// no-ops against the active replacement. The active generation and the +// rearm, or clear the arm child. An owning replacement activation arms its new +// generation without a model turn. A replacement handoff carries actionable +// closes that were still pending delivery; its durable state lives below +// state/extensions/-primary-watch/. Terminal shutdown leaves the final +// generation stopped so late callbacks cannot rearm. Stale callbacks from an +// earlier generation, including callbacks retained by a superseded core instance, +// are no-ops against the active replacement. Compaction is not a session +// replacement and never enters this lifecycle. The active generation and the // process-exit fallback are process-wide, not per-core-instance. import { spawn, spawnSync, type ChildProcess } from "node:child_process"; import { createHash, randomUUID } from "node:crypto"; @@ -37,6 +39,19 @@ type CloseClassification = { message: string; }; +type PendingActionableClose = { + version: 1; + token: string; + message: string; + predecessorArmPid: string; + delivered?: true; +}; + +type ReplacementActionableHandoff = { + version: 2; + pending: PendingActionableClose[]; +}; + // One outstanding watcher recovery episode, as reported by an established // successor arm. bin/fm-wake-lib.sh owns the marker grammar; this adapter only // carries the generation back so the handling handshake can name it. @@ -53,11 +68,16 @@ type RestorationResult = { type SessionGeneration = { id: number; stopping: boolean; + replacement: boolean; child: ChildProcess | null; retryTimer: ReturnType | null; + cleanupTimer: ReturnType | null; retryFailures: number; restoring: boolean; seq: number; + pendingActionables: PendingActionableClose[]; + cleanupFailure: string; + wakeAcknowledgements: Map void }>; }; export type ArmResult = { @@ -78,22 +98,20 @@ export type PrimaryWatchCoreOptions = { repairToolName: string; encodeOperationalInput: (kind: "watcher", content: string) => string; sendFollowUp: (content: string) => Promise; - // Optional supervision-branch dispatch handshake. When supplied (the OMP - // adapter with its branch extension loaded), the core offers each ordinary - // actionable wake to the branch before delivering it to main; a synchronous - // true means the branch now owns handling the wake and main is not woken. An - // adapter without a supervision branch omits this, and every wake goes to - // main exactly as before (docs/omp-supervision-branch.md). Never offered for - // a repair-failed delivery: only main can repair the watcher cycle. - offerWakeToBranch?: (message: string) => boolean; + // Optional supervision-branch dispatch handshake. A synchronous non-null + // settlement means the branch accepted handling, while rejection returns + // delivery ownership to the core's consumption-acknowledged main path. + // Never offered for repair-failed delivery: only main can repair the cycle. + offerWakeToBranch?: (message: string) => Promise | null; }; export type PrimaryWatchCore = { readonly runtime: "pi" | "omp"; arm: () => ArmResult; armAndWait: () => Promise; + acknowledgeWake: (content: string) => void; markLoaded: () => void; - sessionShutdown: () => void; + sessionShutdown: (replacement?: boolean) => Promise; sessionStart: () => void; }; @@ -152,6 +170,35 @@ function actionableLine(output: string): string { const lines = output.split(/\r?\n/); return lines.find((line) => /^(signal:|stale:|check:|heartbeat($|:))/.test(line)) || ""; } + +function completedActionableLine(output: string): string { + const newline = output.lastIndexOf("\n"); + return newline < 0 ? "" : actionableLine(output.slice(0, newline + 1)); +} + +function nodeErrorCode(error: unknown): string { + return typeof error === "object" && error !== null && "code" in error + ? String((error as { code?: unknown }).code ?? "") + : ""; +} + +type ReplacementActionableReceiver = (pending: PendingActionableClose) => void; +type ActionableDeliveryClaim = { + owner: SessionGeneration; + settlement: Promise<"delivered" | "failed">; +}; +type ReplacementCoordinator = { + receiver: ReplacementActionableReceiver | null; + pending: PendingActionableClose[]; + nextTokenId: number; + deliveries: Map; +}; +type ReplacementCoordinatorGlobal = typeof globalThis & { + __firstmatePrimaryWatchReplacements?: Map; +}; +const replacementCoordinatorGlobal = globalThis as ReplacementCoordinatorGlobal; +const replacementCoordinators = replacementCoordinatorGlobal.__firstmatePrimaryWatchReplacements ??= new Map(); + let nextGenerationId = 0; let activeGeneration: SessionGeneration | null = null; let activeBinding: symbol | null = null; @@ -161,11 +208,16 @@ function createGeneration(): SessionGeneration { return { id: ++nextGenerationId, stopping: false, + replacement: false, child: null, retryTimer: null, + cleanupTimer: null, retryFailures: 0, restoring: false, seq: 0, + pendingActionables: [], + cleanupFailure: "", + wakeAcknowledgements: new Map(), }; } @@ -178,12 +230,16 @@ function generationIsLive(owner: SessionGeneration): boolean { return activeGeneration === owner && !owner.stopping; } -function stopGeneration(owner: SessionGeneration): void { +function stopGeneration(owner: SessionGeneration): ChildProcess | null { owner.stopping = true; clearTimeout(owner.retryTimer ?? undefined); + clearTimeout(owner.cleanupTimer ?? undefined); owner.retryTimer = null; - if (owner.child) owner.child.kill("SIGTERM"); + owner.cleanupTimer = null; + const child = owner.child; + if (child) child.kill("SIGTERM"); owner.child = null; + return child; } function cleanupOnProcessExit(): void { @@ -216,6 +272,20 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar offerWakeToBranch, } = options; const armScript = `${fmRoot}/bin/fm-watch-arm.sh`; + const handoffDir = `${state}/extensions/${runtime}-primary-watch`; + const actionableHandoff = `${handoffDir}/session-replacement-actionable.json`; + let nextHandoffId = 0; + let replacementHandoff: PendingActionableClose[] | null = null; + let replacementCoordinator = replacementCoordinators.get(actionableHandoff); + if (!replacementCoordinator) { + replacementCoordinator = { + receiver: null, + pending: [], + nextTokenId: 0, + deliveries: new Map(), + }; + replacementCoordinators.set(actionableHandoff, replacementCoordinator); + } const extensionVersionHash = createHash("sha256") .update(readFileSync(extensionFile)) .update(readFileSync(coreFile)); @@ -240,12 +310,120 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar let generation = createGeneration(); const armReadiness = new WeakMap>(); const armClose = new WeakMap>(); + const armPendingActionable = new WeakMap(); // The recovery generation an established successor arm reported, keyed by the // arm child that reported it. bin/fm-watch-arm.sh only prints it while a // recovery episode is outstanding, so its absence means there is nothing to // hand off and the ordinary wake path applies unchanged. const armRecovery = new WeakMap(); + function createPendingActionable(message: string, predecessorArmPid: string): PendingActionableClose { + return { + version: 1, + token: `${process.pid}-${Date.now()}-${++replacementCoordinator.nextTokenId}`, + message, + predecessorArmPid, + }; + } + + function validatePendingActionable(value: unknown): PendingActionableClose { + if ( + typeof value !== "object" || value === null || + (value as { version?: unknown }).version !== 1 || + typeof (value as { token?: unknown }).token !== "string" || + !/^[0-9]+-[0-9]+-[0-9]+$/.test((value as { token: string }).token) || + typeof (value as { message?: unknown }).message !== "string" || + !actionableLine((value as { message: string }).message) || + typeof (value as { predecessorArmPid?: unknown }).predecessorArmPid !== "string" || + !/^[0-9]*$/.test((value as { predecessorArmPid: string }).predecessorArmPid) || + ((value as { delivered?: unknown }).delivered !== undefined && + (value as { delivered?: unknown }).delivered !== true) + ) { + throw new Error(`invalid ${runtimeLabel} replacement actionable handoff at ${actionableHandoff}`); + } + return value as PendingActionableClose; + } + + function validateReplacementHandoff(value: unknown): PendingActionableClose[] { + if ( + typeof value !== "object" || value === null || + (value as { version?: unknown }).version !== 2 || + !Array.isArray((value as { pending?: unknown }).pending) || + (value as { pending: unknown[] }).pending.length === 0 + ) { + throw new Error(`invalid ${runtimeLabel} replacement actionable handoff at ${actionableHandoff}`); + } + const pending = (value as { pending: unknown[] }).pending.map(validatePendingActionable); + if (new Set(pending.map((item) => item.token)).size !== pending.length) { + throw new Error(`invalid ${runtimeLabel} replacement actionable handoff at ${actionableHandoff}`); + } + return pending; + } + + function writeReplacementHandoff(pending: PendingActionableClose[]): void { + replacementHandoff = [...pending]; + mkdirSync(handoffDir, { recursive: true }); + const temporary = `${actionableHandoff}.tmp-${process.pid}-${++nextHandoffId}`; + const handoff: ReplacementActionableHandoff = { version: 2, pending }; + try { + writeFileSync(temporary, `${JSON.stringify(handoff)}\n`, { mode: 0o600 }); + renameSync(temporary, actionableHandoff); + } catch (error) { + try { + unlinkSync(temporary); + } catch { + // Preserve the original handoff publication error. + } + throw error; + } + } + + function persistReplacementHandoff(pending: PendingActionableClose[]): void { + if (pending.length === 0) return; + writeReplacementHandoff(pending); + } + + function loadReplacementHandoff(): PendingActionableClose[] { + try { + const pending = validateReplacementHandoff(JSON.parse(readFileSync(actionableHandoff, "utf8"))); + replacementHandoff = pending; + return [...pending]; + } catch (error) { + if (nodeErrorCode(error) === "ENOENT") { + replacementHandoff = null; + return []; + } + throw error; + } + } + + function mergeReplacementHandoff(pending: PendingActionableClose): void { + let stored: PendingActionableClose[] = []; + try { + stored = validateReplacementHandoff(JSON.parse(readFileSync(actionableHandoff, "utf8"))); + } catch (error) { + if (nodeErrorCode(error) !== "ENOENT") throw error; + } + if (!stored.some((item) => item.token === pending.token)) stored.push(pending); + writeReplacementHandoff(stored); + } + + function clearReplacementHandoff(pending: PendingActionableClose): void { + try { + const stored = validateReplacementHandoff(JSON.parse(readFileSync(actionableHandoff, "utf8"))); + const remaining = stored.filter((item) => item.token !== pending.token); + if (remaining.length === stored.length) return; + if (remaining.length > 0) { + writeReplacementHandoff(remaining); + } else { + replacementHandoff = null; + unlinkSync(actionableHandoff); + } + } catch (error) { + if (nodeErrorCode(error) !== "ENOENT") throw error; + } + } + function lockOwnership(): LockOwnership { let lockPid = ""; try { @@ -263,6 +441,49 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar return pidAlive(lockPid) ? "other" : "missing"; } + async function waitForGenerationChildClose(armChild: ChildProcess | null): Promise { + if (!armChild) return; + const closed = armClose.get(armChild); + if (!closed) return; + await new Promise((resolveWait) => { + const timer = setTimeout(resolveWait, armRetireTimeoutMs); + void closed.then(() => { + clearTimeout(timer); + resolveWait(); + }); + }); + } + + async function stopSessionGeneration(owner: SessionGeneration, replacement: boolean): Promise { + owner.replacement = replacement; + let persistedTokens = ""; + try { + if (replacement && owner.pendingActionables.length > 0) { + persistReplacementHandoff(owner.pendingActionables); + persistedTokens = owner.pendingActionables.map((pending) => pending.token).join("\n"); + } + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + for (const pending of owner.pendingActionables) { + if (replacementCoordinator.pending.some((item: PendingActionableClose) => item.token === pending.token)) continue; + replacementCoordinator.pending.push({ + ...pending, + message: + `${pending.message}\n\nwatcher: FAILED - ${runtimeLabel} extension could not persist ` + + `a replacement-session actionable wake\n${detail}`, + }); + } + throw error; + } finally { + const child = stopGeneration(owner); + await waitForGenerationChildClose(child); + } + const currentTokens = owner.pendingActionables.map((pending) => pending.token).join("\n"); + if (replacement && currentTokens && currentTokens !== persistedTokens) { + persistReplacementHandoff(owner.pendingActionables); + } + } + function markLoaded(): void { if (lockOwnership() === "other") return; mkdirSync(state, { recursive: true }); @@ -337,13 +558,29 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar }; } - async function sendWake(owner: SessionGeneration, message: string): Promise { - if (!generationIsLive(owner)) return; + async function sendWake(owner: SessionGeneration, message: string, token?: string): Promise { + if (!generationIsLive(owner)) return false; const content = encodeOperationalInput( "watcher", `FIRSTMATE WATCHER WAKE: ${message}\n\nRun bin/fm-wake-drain.sh first and handle the queued wake. Watcher continuity is extension-owned.`, ); - await sendFollowUp(content); + if (!token) { + await sendFollowUp(content); + return generationIsLive(owner); + } + let settleConsumption: (consumed: boolean) => void = () => {}; + const consumption = new Promise((resolveConsumption) => { + settleConsumption = resolveConsumption; + }); + owner.wakeAcknowledgements.set(token, { content, settle: settleConsumption }); + try { + await sendFollowUp(content); + return await consumption; + } catch (error) { + owner.wakeAcknowledgements.delete(token); + settleConsumption(false); + throw error; + } } function confirmHandlingDelivery(recovery: RecoveryHandoff): { ok: boolean; detail: string } { @@ -397,24 +634,29 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar owner: SessionGeneration, message: string, repairFailed: boolean, + token: string, recovery?: RecoveryHandoff, - ): Promise { - if (!generationIsLive(owner)) return; + ): Promise { + if (!generationIsLive(owner)) return false; if (recovery) { const confirmed = confirmHandlingDeliveryWithRetry(owner, recovery); if (!confirmed.ok) { if (!pidAlive(recovery.watcherPid)) await retireArm(owner.child); - await sendWake(owner, `${message}\n\n${confirmed.detail}`); - return; + return await sendWake(owner, `${message}\n\n${confirmed.detail}`, token); } } - // Offer an ordinary actionable wake to the supervision branch when one is - // wired in; a synchronous accept means the branch owns handling it and main - // stays quiet. A repair-failed delivery is never offered - only main can - // repair the watcher cycle. Absent a branch, this is a no-op and every wake - // goes to main exactly as before. - if (!repairFailed && offerWakeToBranch?.(message)) return; - await sendWake(owner, message); + if (!repairFailed && offerWakeToBranch) { + const branchDelivery = offerWakeToBranch(message); + if (branchDelivery) { + try { + await branchDelivery; + return true; + } catch { + // The core retains delivery ownership when branch settlement rejects. + } + } + } + return await sendWake(owner, message, token); } function surfaceFailure(owner: SessionGeneration, message: string): void { @@ -423,6 +665,152 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar }); } + function enqueuePendingActionable(owner: SessionGeneration, pending: PendingActionableClose): void { + if (owner.pendingActionables.some((item) => item.token === pending.token)) return; + owner.pendingActionables.push(pending); + if (owner.stopping && owner.replacement) { + let replacementPending = pending; + try { + mergeReplacementHandoff(pending); + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + replacementPending = { + ...pending, + message: + `${pending.message}\n\nwatcher: FAILED - ${runtimeLabel} extension could not persist ` + + `a late replacement-session actionable wake\n${detail}`, + }; + } + if (replacementCoordinator.receiver) { + replacementCoordinator.receiver(replacementPending); + } else if (replacementPending !== pending) { + replacementCoordinator.pending.push(replacementPending); + } + } + } + + function finishPendingActionable(owner: SessionGeneration, pending: PendingActionableClose): void { + clearReplacementHandoff(pending); + const index = owner.pendingActionables.findIndex((item) => item.token === pending.token); + if (index >= 0) owner.pendingActionables.splice(index, 1); + owner.cleanupFailure = ""; + } + + function surfaceCleanupFailure(owner: SessionGeneration, error: unknown): void { + const detail = error instanceof Error ? error.message : String(error); + if (owner.cleanupFailure === detail) return; + owner.cleanupFailure = detail; + surfaceFailure( + owner, + `watcher: FAILED - ${runtimeLabel} extension could not clear a delivered replacement-session actionable wake\n${detail}`, + ); + } + + function schedulePendingCleanup(owner: SessionGeneration): void { + if (!generationIsLive(owner) || owner.cleanupTimer) return; + const timer = setTimeout(() => { + if (owner.cleanupTimer === timer) owner.cleanupTimer = null; + void processPendingActionables(owner); + }, retryDelay(1)); + timer.unref(); + owner.cleanupTimer = timer; + } + + async function processPendingActionables(owner: SessionGeneration): Promise { + if (!generationIsLive(owner) || owner.restoring || owner.pendingActionables.length === 0) return; + owner.restoring = true; + const attemptedCleanup = new Set(); + try { + while (generationIsLive(owner) && owner.pendingActionables.length > 0) { + for (const delivered of owner.pendingActionables.filter( + (item) => item.delivered && !attemptedCleanup.has(item.token), + )) { + attemptedCleanup.add(delivered.token); + try { + finishPendingActionable(owner, delivered); + } catch (error) { + surfaceCleanupFailure(owner, error); + } + } + const pending = owner.pendingActionables.find((item) => !item.delivered); + if (!pending) break; + const existingClaim = replacementCoordinator.deliveries.get(pending.token); + if (existingClaim && existingClaim.owner !== owner) { + const settlement = await existingClaim.settlement; + if (!generationIsLive(owner)) return; + if (settlement === "delivered") { + pending.delivered = true; + continue; + } + if (replacementCoordinator.deliveries.get(pending.token) === existingClaim) { + replacementCoordinator.deliveries.delete(pending.token); + } + } + let settleClaim: (settlement: "delivered" | "failed") => void = () => {}; + const settlement = new Promise<"delivered" | "failed">((resolveSettlement) => { + settleClaim = resolveSettlement; + }); + const deliveryClaim = { owner, settlement }; + replacementCoordinator.deliveries.set(pending.token, deliveryClaim); + const releaseClaim = (): void => { + if (replacementCoordinator.deliveries.get(pending.token) === deliveryClaim) { + replacementCoordinator.deliveries.delete(pending.token); + } + }; + try { + const restoration = await restoreAfterActionableClose(owner, pending.predecessorArmPid); + if (!generationIsLive(owner)) { + settleClaim("failed"); + releaseClaim(); + return; + } + const message = restoration.failure + ? `${pending.message}\n\n${restoration.failure}` + : pending.message; + const delivered = await deliverActionableWake( + owner, + message, + Boolean(restoration.failure), + pending.token, + restoration.recovery, + ); + if (!delivered) { + settleClaim("failed"); + releaseClaim(); + return; + } + pending.delivered = true; + settleClaim("delivered"); + try { + finishPendingActionable(owner, pending); + } catch (error) { + surfaceCleanupFailure(owner, error); + } + releaseClaim(); + } catch (error) { + settleClaim("failed"); + releaseClaim(); + throw error; + } + } + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + surfaceFailure(owner, `watcher: FAILED - ${runtimeLabel} extension could not deliver an actionable wake\n${detail}`); + } finally { + if (generationIsLive(owner)) { + owner.restoring = false; + if (owner.pendingActionables.some((pending) => pending.delivered)) schedulePendingCleanup(owner); + if (!owner.child && !owner.retryTimer) startArm(owner); + } + } + } + + const receiveReplacementActionable: ReplacementActionableReceiver = (pending) => { + if (!generationIsLive(generation)) return; + enqueuePendingActionable(generation, pending); + void processPendingActionables(generation); + }; + function retryDelay(attempt: number): number { return Math.min(retryMaxMs, retryBaseMs * 2 ** Math.max(0, attempt - 1)); } @@ -611,6 +999,12 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar if (/^watcher: (?:started|attached)\b/m.test(combined)) { settleReadiness(true); } + const reason = completedActionableLine(stdout) || completedActionableLine(stderr); + if (reason && !armPendingActionable.has(armChild)) { + const pending = createPendingActionable(reason, String(armChild.pid ?? "")); + armPendingActionable.set(armChild, pending); + enqueuePendingActionable(owner, pending); + } }; const releaseChild = (): void => { if (owner.child === armChild) owner.child = null; @@ -629,34 +1023,18 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar resolveClosed(); settleReadiness(false); releaseChild(); - if (!generationIsLive(owner)) return; const classification = classifyClose(stdout, stderr, code, signal); const predecessor = String(armChild.pid ?? ""); if (classification.kind === "actionable") { - if (owner.restoring) return; + const pending = armPendingActionable.get(armChild) ?? + createPendingActionable(classification.message, predecessor); + enqueuePendingActionable(owner, pending); + if (!generationIsLive(owner)) return; owner.retryFailures = 0; - owner.restoring = true; - void (async () => { - try { - const restoration = await restoreAfterActionableClose(owner, predecessor); - if (!generationIsLive(owner)) return; - const message = restoration.failure - ? `${classification.message}\n\n${restoration.failure}` - : classification.message; - await deliverActionableWake(owner, message, Boolean(restoration.failure), restoration.recovery); - } catch (error) { - const detail = error instanceof Error ? error.message : String(error); - surfaceFailure( - owner, - `watcher: FAILED - ${runtimeLabel} extension could not deliver an actionable wake\n${detail}`, - ); - } finally { - if (generationIsLive(owner)) owner.restoring = false; - } - })(); + void processPendingActionables(owner); return; } - if (owner.restoring) return; + if (!generationIsLive(owner) || owner.restoring) return; scheduleRetry(owner, classification.message, predecessor); }); armChild.on("error", (error: Error) => { @@ -683,7 +1061,7 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar async function armAndWait(): Promise { const owner = generation; - const result = startArm(owner); + const result = activateOwnedWatch(owner); if (!result.ok) return result; const armChild = owner.child; if (!armChild || await waitForReadiness(armChild)) return result; @@ -695,16 +1073,63 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar }; } + function activateOwnedWatch(owner: SessionGeneration): ArmResult { + if (!generationIsLive(owner)) return { ok: false, message: shuttingDownMessage }; + if (lockOwnership() !== "owned") return startArm(owner); + replacementCoordinator.receiver = receiveReplacementActionable; + let pending: PendingActionableClose[] = []; + let loadFailure = ""; + try { + pending = loadReplacementHandoff(); + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + loadFailure = + `watcher: FAILED - ${runtimeLabel} extension could not load a replacement-session actionable wake\n${detail}`; + } + const inProcessPending = replacementCoordinator.pending.splice(0); + for (const actionable of [...pending, ...inProcessPending]) { + enqueuePendingActionable(owner, actionable); + } + if (owner.pendingActionables.length > 0) { + if (loadFailure) surfaceFailure(owner, loadFailure); + const armResult = startArm(owner, owner.pendingActionables[0].predecessorArmPid); + if (!armResult.ok) { + surfaceFailure( + owner, + `watcher: FAILED - ${runtimeLabel} extension could not arm before replacement wake delivery\n${armResult.message}`, + ); + } + void processPendingActionables(owner); + return armResult; + } + const result = startArm(owner); + if (loadFailure) surfaceFailure(owner, `${loadFailure}\n${result.message}`); + return result; + } + + function acknowledgeWake(content: string): void { + for (const [token, acknowledgement] of generation.wakeAcknowledgements) { + if (acknowledgement.content !== content) continue; + generation.wakeAcknowledgements.delete(token); + acknowledgement.settle(true); + break; + } + } + function sessionStart(): void { if (activeBinding !== binding) return; if (generation.stopping) generation = createGeneration(); activateGeneration(generation); markLoaded(); + if (lockOwnership() === "owned") activateOwnedWatch(generation); } - function sessionShutdown(): void { + async function sessionShutdown(replacement = false): Promise { if (activeBinding !== binding) return; - stopGeneration(generation); + for (const acknowledgement of generation.wakeAcknowledgements.values()) acknowledgement.settle(false); + generation.wakeAcknowledgements.clear(); + if (replacementCoordinator.receiver === receiveReplacementActionable) replacementCoordinator.receiver = null; + await stopSessionGeneration(generation, replacement); } activeBinding = binding; @@ -713,8 +1138,9 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar return { runtime, - arm: () => startArm(generation), + arm: () => activateOwnedWatch(generation), armAndWait, + acknowledgeWake, markLoaded, sessionShutdown, sessionStart, diff --git a/docs/omp-supervision-branch.md b/docs/omp-supervision-branch.md index d52709d4695..024c2ca7e8f 100644 --- a/docs/omp-supervision-branch.md +++ b/docs/omp-supervision-branch.md @@ -23,7 +23,7 @@ This feature is OMP-only by construction and changes nothing anywhere else: - The branch itself: `.omp/extensions/fm-branch-supervision-omp.ts` creates and reopens the persistent branch session, serializes wakes, mirrors dialog, and merges outcomes. It checks the current extension generation and `state/.lock` ownership before each guarded branch side effect, so a lost lock or a cold-start re-arm cannot let an old continuation mutate the fleet. The branch conversation is persistent across main's own `/new`, `/resume`, and `/fork` navigation: the mirror re-anchors on its own when main's session file changes, so no live re-arm runs, and `session_shutdown` only latches shutdown. - Every path that cannot reach a working branch falls back to delivering the wake to main - a broken branch degrades to today's behavior, never to a lost wake. + Every accepted path that cannot reach a working branch rejects its settlement to the shared watcher core, which retains delivery ownership until main consumes the follow-up; a broken branch declines later offers so they take that same watcher-owned path directly. - Branch model and effort selection: the same extension registers `/supervision-model`, which picks the branch's model and then its reasoning effort over OMP's portable select dialog, and saves both as pins applied at the next branch build, never to the branch already running; [configuration.md](configuration.md#omp-supervision-branch-model-and-effort-configsupervision-branch-model-configsupervision-branch-effort) owns the full activation boundary, operator-facing schema, and behavior. Model resolution reads OMP's live `ModelRegistry` (available models and their configured credentials) with pure lookups, so pinning the branch never moves main's own conversation; effort uses OMP's `Effort` catalog and the shim's `clampThinkingLevel`. - Branch system prompt: `bin/fm-branch-prompt.sh`; its header owns the byte-stable-prefix contract (no timestamps, no fleet snapshot, no per-wake content). @@ -50,7 +50,7 @@ Three capabilities are deliberately out of scope for this port and are a future The reason is structural: OMP coordinates the two in-process actors through shared filesystem locks (the wake-queue lock and the lease-command lock) acquired by synchronous subprocess calls with no timeout. A settlement that tried to displace a live or hung branch would have to acquire the very locks that branch's own in-flight work may hold, and on OMP's single-threaded event loop that can stall the whole process. Making mid-flight handoff safe needs a fenced, nonblocking rework of that shared lock substrate (which Pi and the rest of the fleet also use), so it is tracked separately rather than shipped here. -A broken or unreachable branch still degrades to the wake-to-main path with no lost wake; a hung branch stalls its own wake queue until the process is killed, and the durable rows survive to be re-presented on the next start. +A broken or unreachable branch still rejects to the watcher-owned main path with no lost wake; a hung branch stalls its own wake queue until the process is killed, and the durable rows survive to be re-presented on the next start. ## Per-actor acknowledgement @@ -64,7 +64,7 @@ This scoping engages only while a branch grant is, or recently was, in play: a h ## Main-fallback re-entry limitation -When the supervision branch is unavailable, the primary adapter falls eligible wakes back to MAIN through the ordinary operational notification path. +When the supervision branch is unavailable or rejects settlement, the shared watcher core retains eligible wakes and delivers them to MAIN through the ordinary consumption-acknowledged operational notification path. Per-actor queue ownership prevents MAIN and the branch from double-consuming a row, but it does not serialize fallback notifications while MAIN handles a claimed, unacknowledged row set. Each valid higher-sequence signal or stale row can therefore inject another priority operational notification during the same handling episode. OMP can preempt or skip the report reads, current-state reconciliation, cleanup, or generation-bound acknowledgement that would finish the active episode. @@ -123,7 +123,7 @@ What is new is only the attended path: outside away mode, the branch absorbs the ## Verification -Portable regressions: `tests/fm-omp-branch-supervision.test.sh` (prompt byte-stability, outcome store append-only, lease actor partition and guards, wake-grant lifecycle, non-branch-home invariance) and the per-actor consume regression in `tests/fm-wake-queue.test.sh` (branch-scoped acknowledgement never swallows a main-owned row, main excludes branch-granted rows, and the inert pre-branch path). +Portable regressions: `tests/fm-omp-branch-supervision.test.sh` covers prompt byte-stability, outcome store append-only, lease actor partition and guards, wake-grant lifecycle, and non-branch-home invariance; `tests/fm-omp-primary.test.sh` rejects an accepted branch settlement into the watcher-owned main path across replacement; and the per-actor consume regression in `tests/fm-wake-queue.test.sh` proves branch-scoped acknowledgement never swallows a main-owned row, main excludes branch-granted rows, and the inert pre-branch path. The versioned branch-marker closure is covered through `tests/fm-session-start.test.sh`, and the secondmate imported-helper trust boundary is covered through `tests/fm-spawn-dispatch-profile.test.sh`. The strict typecheck in `tests/fm-omp-branch-types.test.sh` pins the extension against the installed `@oh-my-pi/pi-coding-agent` package and fails on any renamed or removed named export or effort-level drift. Live guard: `FM_OMP_BRANCH_LIVE_E2E=1 tests/fm-omp-branch-live-e2e.test.sh` exercises the real installed OMP SDK; run it after every OMP upgrade and record the dated result in [docs/verification/runtime-backends.md](verification/runtime-backends.md). diff --git a/docs/supervision-protocols/omp.md b/docs/supervision-protocols/omp.md index 7f010b0654d..64be49f6df0 100644 --- a/docs/supervision-protocols/omp.md +++ b/docs/supervision-protocols/omp.md @@ -11,8 +11,9 @@ When this session owns supervision and away mode is not active: 5. If the extension says no live session holds the lock, run `bin/fm-session-start.sh` to reclaim the session lock, then call `fm_watch_arm_omp` again. 6. The extension starts `bin/fm-watch-arm.sh --restart`, keeps the child attached to the live OMP process, and owns every later successor launch. The tool and the fallback command return only after that child reports readiness, so a `watcher: FAILED` readiness timeout is a real failure to handle under step 11 rather than a slow success. -7. OMP `/new` and `/resume` events inject the session-start instruction exactly once for the new conversation, replace the prior extension generation, and restore the watcher without a foreground watcher command. -8. After an actionable child close, the shared watcher core rechecks session-lock ownership and verifies one successor before it delivers the follow-up notification. +7. OMP `/new`, `/resume`, `/fork`, and session reload emit `session_switch`, replace the prior extension generation, and restore the watcher without a foreground watcher command. + `/new` and `/resume` also inject the session-start instruction exactly once for the new conversation. +8. After an actionable child close, the shared watcher core rechecks session-lock ownership and verifies one successor before it delivers the follow-up notification; a replacement generation receives an actionable close whose prior delivery was not yet consumed. 9. Ordinary work, turn completion, and ordinary notification handling must not call `fm_watch_arm_omp` again because continuity is extension-owned. 10. An unexpected child close enters bounded exponential retry, and an exhausted retry or lost session lock is surfaced as a watcher failure. 11. Missing, failed, or unhealthy cycle only: drain queued notifications, inspect the failure, call `fm_watch_arm_omp`, and restart with the explicit `-e` fallback if the integration is missing or stale. diff --git a/docs/supervision-protocols/pi.md b/docs/supervision-protocols/pi.md index 102dd4ad1dc..3b8c172e2a0 100644 --- a/docs/supervision-protocols/pi.md +++ b/docs/supervision-protocols/pi.md @@ -4,13 +4,13 @@ When this session owns supervision and away mode is not active: 1. Drain first with `bin/fm-wake-drain.sh`. After handling all emitted wakes and reconciling open decisions and unread status lines, run the exact `--ack-through` command printed as `WAKE_ACK_REQUIRED`; until then the work remains durable for idempotent re-handling after interruption. 2. Confirm the Pi primary auto-loaded both project extensions (plain `pi` or `pi-signed`, after approving project trust once per clone); if not, restart the selected executable with `-e __FM_PI_TURNEND_EXT__ -e __FM_PI_EXT__` as a trust-free fallback. -3. First cycle only: make the one required `fm_watch_arm_pi` call. +3. Initial process cycle only: make the one required `fm_watch_arm_pi` call; if startup already owned the fleet lock, this is an ownership-based no-op. Use `/fm-watch-arm-pi` only as a human-entered fallback. Never run `bin/fm-watch-arm.sh` through Pi's bash tool because that foreground arm can wedge the agent and bypasses extension-owned cleanup. 4. If the extension says no live session holds the lock, run `bin/fm-session-start.sh` to reclaim the session lock, then call `fm_watch_arm_pi` again. 5. The extension starts `bin/fm-watch-arm.sh --restart`, keeps the child attached to the live Pi process, and owns every later successor launch. -6. Ordinary same-process session replacement (`/new`, `/resume`, `/fork`, reload) retires only the prior generation; call `fm_watch_arm_pi` once for the first cycle of the replacement session without restarting Pi. - The generation-owner contract lives in the shared `bin/fm-primary-watch-core.ts` that this extension binds. +6. Ordinary same-process session replacement (`/new`, `/resume`, `/fork`, reload) retires only the prior generation; when the replacement owns the fleet lock, its `session_start` arms the new generation without a model turn or another `fm_watch_arm_pi` call. + The shared `bin/fm-primary-watch-core.ts` owns the generation and in-flight actionable-close handoff contract. 7. After an actionable child close, the extension rechecks session-lock ownership and verifies one successor before it delivers the follow-up wake; its bounded fallback is defined in `docs/watcher-continuity.md`. 8. Ordinary work, turn completion, and ordinary signal, stale, check, heartbeat, or other wake handling: do not call `fm_watch_arm_pi` again because continuity is extension-owned rather than model-memory-owned. 9. An unexpected child close enters bounded exponential retry, and an exhausted retry or lost session lock is surfaced as a watcher failure instead of disappearing. diff --git a/docs/verification/supervision.md b/docs/verification/supervision.md index a880a4a2d48..24165234d17 100644 --- a/docs/verification/supervision.md +++ b/docs/verification/supervision.md @@ -311,18 +311,38 @@ 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 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): +Pi and OMP same-process session-transition ownership was verified on 2026-09-03 against the tracked adapters and shared core with provider-free public lifecycle events and real watcher child processes under Node v25.2.1: ```sh -pi --version -tests/fm-pi-watch-extension.test.sh -tests/fm-pi-primary-types.test.sh +bin/fm-test-run.sh tests/fm-pi-watch-extension.test.sh +bin/fm-test-run.sh tests/fm-omp-primary.test.sh ``` -Observed guarantee: after ordinary `session_shutdown` for `/new`, `/resume`, and `/fork`, plus same-instance shutdown-plus-start, the replacement generation armed again without a Pi restart and without the `watcher: not armed - Pi session is shutting down` refusal. +```text +ok - Pi session transitions auto-arm through a generation owner across /new /resume /fork/reload, stale callbacks, and quit +ok - Pi session replacement auto-arms and carries its in-flight actionable close +ok - OMP /new /resume /fork and reload session_switch paths auto-arm and carry an in-flight actionable close +``` + +For Pi, ordinary `session_shutdown` for `/new`, `/resume`, `/fork`, and reload plus same-instance shutdown-plus-start is followed by an owning `session_start` that arms before any model turn and without the `watcher: not armed - Pi session is shutting down` refusal. +The replacement consumed exactly once the actionable close whose first follow-up had begun but was not consumed before shutdown, while retaining one live successor. Stale prior-generation tool callbacks could not mutate the active child, repeated transitions kept exactly one live arm cycle, and terminal `quit` still refused late rearm. 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`. -OMP binds that same core through `.omp/extensions/fm-primary-omp.ts` and is covered by its own `tests/fm-omp-primary.test.sh` and the OMP row above; the remaining primary harnesses are not applicable because they do not use this extension generation lifecycle. +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-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 +FM_OMP_PRIMARY_LIVE_E2E=1 bin/fm-test-run.sh tests/fm-omp-primary-live-e2e.test.sh +``` + +```text +not ok - OMP /new did not inject exactly one startup instruction for the new session +``` + +The identical pre-port failure establishes a pre-existing OMP 18.1.5 session-start-nudge compatibility finding rather than a watcher-continuity regression. +The provider-free adapter regression remains the current proof that its `/new` and `/resume` replacement handlers each stage one nudge while re-arming, and the live guard remains intentionally red until that separate compatibility finding is resolved. +The remaining primary harnesses are not applicable because they do not bind the shared Pi-compatible extension generation lifecycle. The once-per-generation recovery bound and immediate handling-successor poll were verified on 2026-08-21 at revision `549dd1e0ff05f96607c5e7457b4d8e3d7396bd16` with the tracked Pi extension, real watcher processes, and an isolated home. The regression forced handling confirmation to fail, observed one recovery follow-up across the former repeat window, confirmed the successor remained live, and then proved a separate handling successor durably queued a crew event within the bounded poll window. diff --git a/docs/watcher-continuity.md b/docs/watcher-continuity.md index fb57bac2162..686362b46a1 100644 --- a/docs/watcher-continuity.md +++ b/docs/watcher-continuity.md @@ -11,7 +11,7 @@ Because both files decide watcher behavior, a loaded adapter's marker version ha Editing only the core therefore makes a running session's marker stale, and `bin/fm-session-start.sh` prints the exact restart fallback instead of accepting a live session that still runs the old core. Each adapter starts the next arm before delivering the wake prompt, checks current session-lock ownership at launch, preserves one child or scheduled retry at a time, and applies bounded exponential retry after an unexpected or failed close. A failed follow-up never cancels continuity restoration. -Same-process session replacement on Pi and OMP follows the generation-owner contract stated once in `bin/fm-primary-watch-core.ts`. +Same-process session replacement on Pi and OMP follows the generation-owner contract stated once in `bin/fm-primary-watch-core.ts`: an owning replacement activation arms without a model turn, and a state-scoped handoff carries every actionable close whose delivery overlaps replacement, including a follow-up not yet consumed and a retiring child that reports during shutdown. Claude's `.claude/settings.json` Stop `asyncRewake` hook (`bin/fm-claude-stop-autoarm.sh`) owns routine tokenless re-arm. The hook fires on every Stop, and an eligible primary with supervision need admits one home-scoped owner that foregrounds `bin/fm-watch-arm.sh` inside the hook-owned process tree. A numeric session-lock owner that fails the shared `fm_harness_pid_alive` predicate is reclaimed through `bin/fm-lock.sh` before auto-arm state changes, while a live owner, absent lock, or malformed lock keeps the competing hook inert. @@ -95,7 +95,8 @@ Only the watcher process touches `state/.last-watcher-beat`; no helper process c ## Regression coverage `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 ordinary same-process session replacement for `/new`, `/resume`, and `/fork`, same-instance shutdown-plus-start, 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. +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. 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 2850b6330b1..a824f712197 100755 --- a/tests/fm-omp-primary.test.sh +++ b/tests/fm-omp-primary.test.sh @@ -477,6 +477,14 @@ const api = { sendMessage(message, options) { watcherMessages.push({ message, options }); }, sendUserMessage(content, options) { customMessages.push({ content, options }); }, }; +async function waitForWatchCount(expected, label) { + const file = `${process.env.FM_STATE_OVERRIDE}/watch-count`; + for (let i = 0; i < 100; i += 1) { + if (existsSync(file) && Number(readFileSync(file, "utf8").trim()) === expected) return; + await new Promise(resolve => setTimeout(resolve, 10)); + } + throw new Error(`timeout waiting for ${label}`); +} process.argv[1] = process.env.EXTENSION; const extension = await import(`${pathToFileURL(process.env.EXTENSION).href}?test=${Date.now()}`); await extension.default(api); @@ -529,6 +537,7 @@ if (await handlers.get("before_agent_start")({ type: "before_agent_start" }, {}) } writeFileSync(`${process.env.FM_STATE_OVERRIDE}/.lock`, `${process.pid}\n`); await handlers.get("session_switch")({ type: "session_switch", reason: "new" }, extensionContext); +await waitForWatchCount(1, "in-process OMP /new automatic watcher arm"); const newStartup = await handlers.get("before_agent_start")({ type: "before_agent_start" }, {}); if (newStartup?.message?.customType !== "firstmate-sessionstart-nudge" || newStartup.message.attribution !== "agent") { throw new Error(`in-process OMP /new lost its once-only startup instruction: ${JSON.stringify(newStartup)}`); @@ -537,6 +546,7 @@ if (await handlers.get("before_agent_start")({ type: "before_agent_start" }, {}) throw new Error("in-process OMP /new repeated its startup instruction"); } await handlers.get("session_switch")({ type: "session_switch", reason: "resume" }, extensionContext); +await waitForWatchCount(2, "in-process OMP /resume automatic watcher arm"); const resumeStartup = await handlers.get("before_agent_start")({ type: "before_agent_start" }, {}); if (resumeStartup?.message?.customType !== "firstmate-sessionstart-nudge" || resumeStartup.message.attribution !== "agent") { throw new Error(`in-process OMP /resume lost its once-only startup instruction: ${JSON.stringify(resumeStartup)}`); @@ -940,6 +950,137 @@ JS pass "OMP surfaces a refused handling handshake as one typed wake" } +test_native_omp_session_switch_carries_inflight_actionable_close() { + local fixture out status=0 + fixture="$TMP_ROOT/native-session-replacement" + 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 'watcher: started pid=%s (beacon fresh)\n' "$$" +trap 'exit 0' TERM INT +if [ "$count" -eq 1 ]; then + while [ ! -e "$state/watch-trigger" ]; do sleep 0.02; done + printf 'signal: omp replacement in-flight actionable close\n' + 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" \ + 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 = []; +let withholdConsumption = true; +let rejectNextBranchOffer = true; +const api = { + zod: { object: () => ({}) }, + events: { + emit(name, offer) { + if (name !== "fm-branch-supervision:dispatch" || !rejectNextBranchOffer) return; + rejectNextBranchOffer = false; + offer.accept(Promise.reject(new Error("synthetic branch settlement rejection"))); + }, + }, + on(name, handler) { handlers.set(name, handler); }, + registerCommand() {}, + registerTool() {}, + sendMessage(message) { + steers.push(String(message?.content ?? "")); + if (withholdConsumption) { + withholdConsumption = false; + return; + } + handlers.get("before_agent_start")?.({ type: "before_agent_start", prompt: message.content }, {}); + }, +}; +const count = () => existsSync(`${process.env.FM_STATE_OVERRIDE}/watch-count`) + ? Number(readFileSync(`${process.env.FM_STATE_OVERRIDE}/watch-count`, "utf8").trim()) + : 0; +async function waitFor(pred, label) { + for (let i = 0; i < 500; i += 1) { + if (pred()) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + throw new Error(`timeout waiting for ${label}`); +} + +writeFileSync(`${process.env.FM_STATE_OVERRIDE}/.lock`, `${process.pid}\n`); +process.argv[1] = process.env.EXTENSION; +const extension = await import(`${pathToFileURL(process.env.EXTENSION).href}?replacement=${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(`${process.env.FM_STATE_OVERRIDE}/watch-trigger`, "trigger\n"); +await waitFor(() => steers.length === 1 && count() >= 2, "in-flight OMP delivery and successor"); + +await handlers.get("session_switch")({ type: "session_switch", reason: "new" }, context); +await waitFor(() => count() >= 3, "OMP /new replacement arm"); +await waitFor( + () => steers.filter((message) => message.includes("signal: omp replacement in-flight actionable close")).length === 2, + "OMP replacement actionable replay", +); +const handoff = `${process.env.FM_STATE_OVERRIDE}/extensions/omp-primary-watch/session-replacement-actionable.json`; +await waitFor(() => !existsSync(handoff), "OMP replacement handoff cleanup"); + +for (const [reason, label] of [["resume", "/resume"], ["fork", "/fork"], ["resume", "reload"]]) { + const before = count(); + await handlers.get("session_switch")({ type: "session_switch", reason }, context); + await waitFor(() => count() > before, `OMP ${label} replacement arm`); +} +if (steers.filter((message) => message.includes("signal: omp replacement in-flight actionable close")).length !== 2) { + throw new Error(`OMP replacement did not replay exactly once: ${steers.join(" | ")}`); +} +writeFileSync(`${process.env.FM_STATE_OVERRIDE}/watch-stop`, "stop\n"); +await handlers.get("session_shutdown")({ type: "session_shutdown" }, context); +console.log("omp-session-replacement-ok"); +JS + ) || status=$? + expect_code 0 "$status" "OMP session-switch replacement continuity" + assert_contains "$out" omp-session-replacement-ok "OMP replacement continuity test did not complete" + pass "OMP /new /resume /fork and reload session_switch paths auto-arm and carry an in-flight actionable close" +} + test_resolve_path_uses_node_when_readlink_f_is_unavailable test_exact_bun_omp_primary_identity test_standalone_omp_primary_identity @@ -951,3 +1092,4 @@ test_primary_marker_refuses_whitespace_identity 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 diff --git a/tests/fm-pi-primary-live-e2e.test.sh b/tests/fm-pi-primary-live-e2e.test.sh index 364a52d46d4..a8ab3c89cdb 100755 --- a/tests/fm-pi-primary-live-e2e.test.sh +++ b/tests/fm-pi-primary-live-e2e.test.sh @@ -117,6 +117,20 @@ wait_pid_dead() { return 1 } +wait_pid_change() { + local file=$1 old=$2 i=0 pid + while [ "$i" -lt 120 ]; do + pid=$(sed -n '1p' "$file" 2>/dev/null || true) + if [ -n "$pid" ] && [ "$pid" != "$old" ] && kill -0 "$pid" 2>/dev/null; then + printf '%s\n' "$pid" + return 0 + fi + sleep 0.5 + i=$((i + 1)) + done + return 1 +} + run_ahoy_case() { local label=$1 preceding=$2 expected=$3 out status=0 out=$( @@ -307,7 +321,7 @@ sleep 0.2 : > "$HOME_DIR/state/pi-e2e.meta" send_prompt "Start supervision with fm_watch_arm_pi and never use bash to arm supervision. After the watcher wake arrives, run bin/fm-wake-drain.sh and reply exactly HANDLED." -wait_for_text "watcher: started Pi extension arm child 1" || fail "Pi did not render the initial watcher tool result" +wait_for_text "Pi extension already owns an arm child" || fail "Pi did not render the startup-owned watcher tool result" printf 'done: pi live e2e watcher fire\n' > "$HOME_DIR/state/pi-e2e.status" i=0 @@ -336,6 +350,21 @@ watcher_pid=$(sed -n '1p' "$pid_file") arm_pid=$(ps -p "$watcher_pid" -o ppid= | tr -d ' ') [ -n "$arm_pid" ] || fail "re-armed watcher parent was not live" +send_prompt "/new" +replacement_watcher_pid=$(wait_pid_change "$pid_file" "$watcher_pid") \ + || fail "Pi /new did not auto-arm a replacement watcher generation" +replacement_arm_pid=$(ps -p "$replacement_watcher_pid" -o ppid= | tr -d ' ') +[ -n "$replacement_arm_pid" ] || fail "Pi /new replacement watcher parent was not live" +wait_pid_dead "$watcher_pid" || fail "Pi /new left the prior watcher alive" +wait_pid_dead "$arm_pid" || fail "Pi /new left the prior arm child alive" +send_prompt "Reply exactly PI_REPLACEMENT_READY." +wait_for_exact_line "PI_REPLACEMENT_READY" 120 || fail "Pi replacement session did not accept a prompt" +pane=$(capture) +arm_tool_result_count=$(printf '%s\n' "$pane" | grep -Ec 'watcher: (started|unchanged|not armed|read-only)' || true) +[ "$arm_tool_result_count" -eq 1 ] || fail "Pi /new required a model-owned watcher re-arm (tool-result count $arm_tool_result_count)" +watcher_pid=$replacement_watcher_pid +arm_pid=$replacement_arm_pid + "$TMUX" -L "$SOCKET" send-keys -t "$SESSION" -l '/quit' sleep 1 "$TMUX" -L "$SOCKET" send-keys -t "$SESSION" Enter @@ -343,4 +372,4 @@ wait_for_text "PI_EXIT=0" 60 || fail "Pi did not exit cleanly" wait_pid_dead "$watcher_pid" || fail "watcher child survived clean Pi exit" wait_pid_dead "$arm_pid" || fail "arm child survived clean Pi exit" -printf 'ok - Pi %s live E2E covered the Calm working ship, Ahoy first/later messages, legacy transcripts, near misses, and watcher continuity\n' "$PI_VERSION" +printf 'ok - Pi %s live E2E covered the Calm working ship, Ahoy first/later messages, legacy transcripts, near misses, watcher continuity, and automatic /new re-arm\n' "$PI_VERSION" diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index 28a4b927818..bc13f37d927 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -662,14 +662,18 @@ import { pathToFileURL } from "node:url"; let tool = null; const prompts = []; +const handlers = new Map(); const pi = { - on() {}, + on(event, handler) { + handlers.set(event, handler); + }, registerCommand() {}, registerTool(candidate) { if (candidate.name === "fm_watch_arm_pi") tool = candidate; }, sendUserMessage: async (message) => { prompts.push(message); + handlers.get("before_agent_start")?.({ prompt: message }, {}); }, }; const rows = () => existsSync(process.env.FM_ARM_LOG) @@ -1048,10 +1052,6 @@ const mod = await import(pathToFileURL(process.env.PLUGIN).href); const startup = makePi(); mod.default(startup.pi); await startup.handlers.get("session_start")?.({ type: "session_start", reason: "startup" }, {}); -const first = await startup.getTool().execute("startup", {}, undefined, undefined, {}); -if (!first.details?.ok || !String(first.details.message).includes("started Pi extension arm child")) { - throw new Error(`startup arm failed: ${JSON.stringify(first.details)}`); -} await waitFor(() => existsSync(process.env.FM_CHILD_PID_FILE), "startup child"); const startupChild = readFileSync(process.env.FM_CHILD_PID_FILE, "utf8").trim(); if (!pidAlive(startupChild)) throw new Error("startup child was not alive"); @@ -1072,13 +1072,6 @@ async function replaceSession(previous, reason) { reason, previousSessionFile: `/tmp/previous-${reason}.jsonl`, }, {}); - const armed = await next.getTool().execute(`arm-${reason}`, {}, undefined, undefined, {}); - if (!armed.details?.ok) { - throw new Error(`${reason} replacement arm failed: ${JSON.stringify(armed.details)}`); - } - if (String(armed.details.message).includes("shutting down")) { - throw new Error(`${reason} replacement still refused with shutting-down latch`); - } await waitFor(() => { if (!existsSync(process.env.FM_CHILD_PID_FILE)) return false; const child = readFileSync(process.env.FM_CHILD_PID_FILE, "utf8").trim(); @@ -1094,15 +1087,12 @@ async function replaceSession(previous, reason) { let current = await replaceSession(startup, "new"); current = await replaceSession(current, "resume"); current = await replaceSession(current, "fork"); +current = await replaceSession(current, "reload"); // Same bound instance: ordinary shutdown then session_start without a fresh factory. const sameInstanceChild = readFileSync(process.env.FM_CHILD_PID_FILE, "utf8").trim(); await current.handlers.get("session_shutdown")?.({ type: "session_shutdown", reason: "new" }, {}); await current.handlers.get("session_start")?.({ type: "session_start", reason: "new" }, {}); -const sameInstanceArm = await current.getTool().execute("same-instance", {}, undefined, undefined, {}); -if (!sameInstanceArm.details?.ok || String(sameInstanceArm.details.message).includes("shutting down")) { - throw new Error(`same-instance replacement arm failed: ${JSON.stringify(sameInstanceArm.details)}`); -} await waitFor(() => { if (!existsSync(process.env.FM_CHILD_PID_FILE)) return false; const child = readFileSync(process.env.FM_CHILD_PID_FILE, "utf8").trim(); @@ -1150,8 +1140,123 @@ EOF status=$? expect_code 0 "$status" "Pi session transitions must rearm through an explicit generation owner (output: $out)" [ -z "$out" ] || fail "Pi session-transition generation owner test printed output: $out" - pass "Pi session transitions use a generation owner across /new /resume /fork, stale callbacks, and quit" + pass "Pi session transitions auto-arm through a generation owner across /new /resume /fork/reload, stale callbacks, and quit" +} + +test_pi_session_replacement_carries_inflight_actionable_close() { + local repo home plugin log trigger stop out status + repo="$TMP_ROOT/pi-session-replacement-handoff-root" + home="$TMP_ROOT/pi-session-replacement-handoff-home" + log="$TMP_ROOT/pi-session-replacement-handoff.log" + trigger="$TMP_ROOT/pi-session-replacement-handoff.trigger" + stop="$TMP_ROOT/pi-session-replacement-handoff.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\n' "$$" >> "${FM_ARM_LOG:?}" +count=$(wc -l < "$FM_ARM_LOG" | tr -d '[:space:]') +printf 'watcher: started pid=%s (beacon fresh)\n' "$$" +trap 'exit 0' TERM INT +if [ "$count" -eq 1 ]; then + while [ ! -e "$FM_TRIGGER_FILE" ]; do sleep 0.02; done + printf 'signal: replacement in-flight actionable close\n' + 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_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 releaseOldDelivery = () => {}; +const oldDeliveryRelease = new Promise((resolve) => { + releaseOldDelivery = resolve; +}); +let oldDeliveryStarted = false; + +function makePi(blockDelivery = false) { + const handlers = new Map(); + let tool = null; + const prompts = []; + const pi = { + on(event, handler) { + handlers.set(event, handler); + }, + registerCommand() {}, + registerTool(candidate) { + if (candidate.name === "fm_watch_arm_pi") tool = candidate; + }, + sendUserMessage: async (message) => { + prompts.push(message); + if (blockDelivery) { + oldDeliveryStarted = true; + await oldDeliveryRelease; + return; + } + handlers.get("before_agent_start")?.({ prompt: message }, {}); + }, + events: { on() {}, emit() {} }, + }; + return { pi, handlers, prompts, getTool: () => tool }; +} + +async function waitFor(pred, label) { + for (let i = 0; i < 500; i += 1) { + if (pred()) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + throw new Error(`timeout waiting for ${label}`); +} + +const armCount = () => existsSync(process.env.FM_ARM_LOG) + ? readFileSync(process.env.FM_ARM_LOG, "utf8").trim().split("\n").filter(Boolean).length + : 0; +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +const originalMod = await import(pathToFileURL(process.env.PLUGIN).href); +const original = makePi(true); +originalMod.default(original.pi); +await original.handlers.get("session_start")?.({ type: "session_start", reason: "startup" }, {}); +await waitFor(() => armCount() === 1, "initial automatic arm"); +writeFileSync(process.env.FM_TRIGGER_FILE, "trigger\n"); +await waitFor(() => oldDeliveryStarted && armCount() >= 2, "blocked old delivery and successor arm"); +await original.handlers.get("session_shutdown")?.({ type: "session_shutdown", reason: "new" }, {}); + +const replacementMod = await import(`${pathToFileURL(process.env.PLUGIN).href}?replacement=inflight`); +const replacement = makePi(false); +replacementMod.default(replacement.pi); +await replacement.handlers.get("session_start")?.({ type: "session_start", reason: "new" }, {}); +await waitFor(() => armCount() >= 3, "replacement automatic arm"); +releaseOldDelivery(); +await waitFor( + () => replacement.prompts.some((message) => message.includes("signal: replacement in-flight actionable close")), + "replacement actionable replay", +); +if (original.prompts.filter((message) => message.includes("signal: replacement in-flight actionable close")).length !== 1) { + throw new Error(`old session did not hold exactly one in-flight delivery: ${original.prompts.join(" | ")}`); +} +if (replacement.prompts.filter((message) => message.includes("signal: replacement in-flight actionable close")).length !== 1) { + throw new Error(`replacement did not consume exactly one replay: ${replacement.prompts.join(" | ")}`); } +const handoff = `${process.env.FM_HOME}/state/extensions/pi-primary-watch/session-replacement-actionable.json`; +await waitFor(() => !existsSync(handoff), "replacement handoff cleanup"); +const redundant = await replacement.getTool().execute("replacement-redundant", {}, undefined, undefined, {}); +if (!redundant.details?.ok || !String(redundant.details.message).includes("unchanged")) { + throw new Error(`replacement lost automatic arm ownership: ${JSON.stringify(redundant.details)}`); +} +writeFileSync(process.env.FM_STOP_FILE, "stop\n"); +await replacement.handlers.get("session_shutdown")?.({ type: "session_shutdown", reason: "quit" }, {}); +EOF +) + status=$? + expect_code 0 "$status" "Pi replacement must auto-arm and carry an in-flight actionable close" + [ -z "$out" ] || fail "Pi session-replacement handoff test printed output: $out" + pass "Pi session replacement auto-arms and carries its in-flight actionable 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" @@ -2405,6 +2510,7 @@ test_pi_established_empty_close_honors_retry_limit 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_rebind_without_shutdown_supersedes_prior_binding test_pi_process_exit_cleanup_listener_lifecycle test_pi_process_exit_cleanup_stops_arm_child From 9440c94072e214e0afdb3033e08928db2c57742f Mon Sep 17 00:00:00 2001 From: dnth Date: Thu, 3 Sep 2026 19:19:16 +0800 Subject: [PATCH 2/6] no-mistakes(review): Reconciled OMP session-switch supervision documentation --- docs/omp-supervision-branch.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/omp-supervision-branch.md b/docs/omp-supervision-branch.md index 024c2ca7e8f..b29233a86f2 100644 --- a/docs/omp-supervision-branch.md +++ b/docs/omp-supervision-branch.md @@ -22,7 +22,7 @@ This feature is OMP-only by construction and changes nothing anywhere else: A fleet-wide heartbeat keeps its own all-or-nothing rule (see "Heartbeat routing" below): it takes every branch-ownable unread row or none of them. - The branch itself: `.omp/extensions/fm-branch-supervision-omp.ts` creates and reopens the persistent branch session, serializes wakes, mirrors dialog, and merges outcomes. It checks the current extension generation and `state/.lock` ownership before each guarded branch side effect, so a lost lock or a cold-start re-arm cannot let an old continuation mutate the fleet. - The branch conversation is persistent across main's own `/new`, `/resume`, and `/fork` navigation: the mirror re-anchors on its own when main's session file changes, so no live re-arm runs, and `session_shutdown` only latches shutdown. + The branch conversation remains resident across ordinary main turns. When main performs `/new`, `/resume`, `/fork`, or reload, OMP emits `session_switch`; the primary watcher retires the prior generation and re-arms its replacement before the next model turn. The mirror re-anchors when main's session file changes, and any actionable close whose delivery overlaps the replacement is carried to the new generation exactly once. `session_shutdown` remains the terminal-process boundary rather than a replacement signal. Every accepted path that cannot reach a working branch rejects its settlement to the shared watcher core, which retains delivery ownership until main consumes the follow-up; a broken branch declines later offers so they take that same watcher-owned path directly. - Branch model and effort selection: the same extension registers `/supervision-model`, which picks the branch's model and then its reasoning effort over OMP's portable select dialog, and saves both as pins applied at the next branch build, never to the branch already running; [configuration.md](configuration.md#omp-supervision-branch-model-and-effort-configsupervision-branch-model-configsupervision-branch-effort) owns the full activation boundary, operator-facing schema, and behavior. Model resolution reads OMP's live `ModelRegistry` (available models and their configured credentials) with pure lookups, so pinning the branch never moves main's own conversation; effort uses OMP's `Effort` catalog and the shim's `clampThinkingLevel`. @@ -37,13 +37,13 @@ This feature is OMP-only by construction and changes nothing anywhere else: ## Transitions and what is out of scope -Ownership transitions happen only at a clean boundary or by killing the process and letting a fresh one re-arm - never by a synchronous handoff from a live or hung branch. +Ownership transitions happen at a clean boundary: a cold `session_start`, an OMP `session_switch` replacement, or killing the process and letting a fresh one re-arm - never by a synchronous handoff from a live or hung branch. A branch generation serializes each wake through a clean completion boundary, and the next wake re-prompts the same resident conversation; only the first wake in a fresh process reopens the durable branch conversation. -`session_start` (a cold start of a fresh process) is the sole clean-boundary arm; the branch is otherwise persistent across main's own session navigation. +`session_start` and `session_switch` are the clean-boundary arm points. A `session_switch` replacement re-arms automatically without a foreground watcher command or a model turn; the branch itself remains resident unless the process is restarted. Three capabilities are deliberately out of scope for this port and are a future iteration: -- Mid-flight branch replacement - displacing a live branch and re-arming a replacement inside the same process. +- Mid-flight branch replacement - displacing a live branch and re-arming a replacement inside the same process. This is distinct from the supported main-session `session_switch` replacement above. - Mid-branch model or effort hot-swap - the branch keeps its model and effort until a fresh process rebuilds it; a `/supervision-model` change is a pin applied at the next branch build, not to the running branch. - Hung-branch live takeover - a branch stuck inside a model call is recovered by killing the process (its leases and wake-grant rows go stale on death and are swept), not by main taking ownership away from it. From 912709add5600ae491981fe26eb2a643795bc552 Mon Sep 17 00:00:00 2001 From: dnth Date: Thu, 3 Sep 2026 19:23:38 +0800 Subject: [PATCH 3/6] no-mistakes(review): Updated OMP lifecycle comment for session-switch re-arm --- .omp/extensions/fm-branch-supervision-omp.ts | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/.omp/extensions/fm-branch-supervision-omp.ts b/.omp/extensions/fm-branch-supervision-omp.ts index 14842fac936..be4c11dbce7 100644 --- a/.omp/extensions/fm-branch-supervision-omp.ts +++ b/.omp/extensions/fm-branch-supervision-omp.ts @@ -939,14 +939,13 @@ ${context.command} enqueueMirrorFlush(); }); - // session_start arms this generation at a cold start (a fresh process). It is - // the sole clean-boundary transition: the branch is persistent across main's - // own /new, /resume, and /fork navigation and is never displaced by a - // synchronous live handoff (that path, and hung-branch takeover, are out of - // scope - see docs/omp-supervision-branch.md). The mirror re-anchors on its - // own when the session file changes (collectMainDialog compares the file), and - // mainModel/mainModelRegistry refresh at every turn_end, so no session_switch - // handling is needed. Terminal quit fires session_shutdown and never a start. + // session_start arms this generation at a cold start (a fresh process). Main + // session replacements are handled by the primary OMP adapter's + // session_switch watcher re-arm; this branch remains resident and the mirror + // re-anchors when the session file changes (collectMainDialog compares the + // file). A synchronous live handoff from a branch, or hung-branch takeover, + // remains out of scope (see docs/omp-supervision-branch.md). Terminal quit + // fires session_shutdown and never a start. pi.on?.("session_start", (_event, ctx) => { rememberMainContext(ctx); shuttingDown = false; From b61605a67bd0949049ba546caf55b07e369b804b Mon Sep 17 00:00:00 2001 From: dnth Date: Thu, 3 Sep 2026 19:27:28 +0800 Subject: [PATCH 4/6] no-mistakes(review): Preserved watcher re-arm after handoff persistence failure --- bin/fm-primary-watch-core.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/bin/fm-primary-watch-core.ts b/bin/fm-primary-watch-core.ts index 4f1587abed9..64073253c28 100644 --- a/bin/fm-primary-watch-core.ts +++ b/bin/fm-primary-watch-core.ts @@ -457,6 +457,7 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar async function stopSessionGeneration(owner: SessionGeneration, replacement: boolean): Promise { owner.replacement = replacement; let persistedTokens = ""; + let persistenceFailed = false; try { if (replacement && owner.pendingActionables.length > 0) { persistReplacementHandoff(owner.pendingActionables); @@ -473,13 +474,13 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar `a replacement-session actionable wake\n${detail}`, }); } - throw error; + persistenceFailed = true; } finally { const child = stopGeneration(owner); await waitForGenerationChildClose(child); } const currentTokens = owner.pendingActionables.map((pending) => pending.token).join("\n"); - if (replacement && currentTokens && currentTokens !== persistedTokens) { + if (replacement && !persistenceFailed && currentTokens && currentTokens !== persistedTokens) { persistReplacementHandoff(owner.pendingActionables); } } From 7673e6dd42b5cc609ed92e17106bb9c8fe0d9b44 Mon Sep 17 00:00:00 2001 From: dnth Date: Thu, 3 Sep 2026 19:30:41 +0800 Subject: [PATCH 5/6] no-mistakes(review): Scheduled retries for undelivered actionable wakes --- bin/fm-primary-watch-core.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bin/fm-primary-watch-core.ts b/bin/fm-primary-watch-core.ts index 64073253c28..87b988b914a 100644 --- a/bin/fm-primary-watch-core.ts +++ b/bin/fm-primary-watch-core.ts @@ -800,7 +800,7 @@ export function createPrimaryWatchCore(options: PrimaryWatchCoreOptions): Primar } finally { if (generationIsLive(owner)) { owner.restoring = false; - if (owner.pendingActionables.some((pending) => pending.delivered)) schedulePendingCleanup(owner); + if (owner.pendingActionables.length > 0) schedulePendingCleanup(owner); if (!owner.child && !owner.retryTimer) startArm(owner); } } From b214e5ad33285a418e2b815301096290586d310e Mon Sep 17 00:00:00 2001 From: dnth Date: Thu, 3 Sep 2026 19:38:09 +0800 Subject: [PATCH 6/6] no-mistakes(document): Clarify compaction remains within active watcher generation --- docs/watcher-continuity.md | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/watcher-continuity.md b/docs/watcher-continuity.md index 686362b46a1..d9808466469 100644 --- a/docs/watcher-continuity.md +++ b/docs/watcher-continuity.md @@ -12,6 +12,7 @@ Editing only the core therefore makes a running session's marker stale, and `bin Each adapter starts the next arm before delivering the wake prompt, checks current session-lock ownership at launch, preserves one child or scheduled retry at a time, and applies bounded exponential retry after an unexpected or failed close. A failed follow-up never cancels continuity restoration. Same-process session replacement on Pi and OMP follows the generation-owner contract stated once in `bin/fm-primary-watch-core.ts`: an owning replacement activation arms without a model turn, and a state-scoped handoff carries every actionable close whose delivery overlaps replacement, including a follow-up not yet consumed and a retiring child that reports during shutdown. +Compaction remains inside the active generation and does not trigger replacement shutdown, re-arm, or actionable handoff. Claude's `.claude/settings.json` Stop `asyncRewake` hook (`bin/fm-claude-stop-autoarm.sh`) owns routine tokenless re-arm. The hook fires on every Stop, and an eligible primary with supervision need admits one home-scoped owner that foregrounds `bin/fm-watch-arm.sh` inside the hook-owned process tree. A numeric session-lock owner that fails the shared `fm_harness_pid_alive` predicate is reclaimed through `bin/fm-lock.sh` before auto-arm state changes, while a live owner, absent lock, or malformed lock keeps the competing hook inert.