From c7b46a0b9feb4ee7b3ce0a5491c3ea03c6d09556 Mon Sep 17 00:00:00 2001 From: dnth Date: Sat, 5 Sep 2026 18:00:05 +0800 Subject: [PATCH 1/2] fix(omp): keep supervision outcome delivery responsive --- .omp/extensions/fm-branch-supervision-omp.ts | 278 +++++++++++-------- .omp/extensions/lib/fm-async-exec.ts | 90 ++++++ .omp/extensions/lib/fm-branch-dispatch.ts | 57 ++-- bin/fm-primary-watch-version-lib.sh | 9 +- docs/omp-supervision-branch.md | 9 + tests/fm-omp-branch-live-e2e.test.sh | 1 + tests/fm-omp-branch-supervision.test.sh | 35 +++ tests/fm-omp-branch-types.test.sh | 1 + tests/fm-omp-primary-live-e2e.test.sh | 3 + tests/fm-omp-primary.test.sh | 8 + tests/fm-omp-secondmate.test.sh | 3 +- tests/fm-send-turn-start.test.sh | 2 + tests/fm-session-start.test.sh | 3 + tests/fm-spawn-dispatch-profile.test.sh | 2 + 14 files changed, 351 insertions(+), 150 deletions(-) create mode 100644 .omp/extensions/lib/fm-async-exec.ts diff --git a/.omp/extensions/fm-branch-supervision-omp.ts b/.omp/extensions/fm-branch-supervision-omp.ts index be4c11dbce7..7cf8599b9eb 100644 --- a/.omp/extensions/fm-branch-supervision-omp.ts +++ b/.omp/extensions/fm-branch-supervision-omp.ts @@ -92,6 +92,7 @@ import { type BranchDispatchOffer, } from "./lib/fm-branch-dispatch.ts"; import { buildBranchModelItems, FOLLOW_MAIN_VALUE } from "./lib/fm-branch-model-picker.ts"; +import { runCommandAsync } from "./lib/fm-async-exec.ts"; const extensionFile = fileURLToPath(import.meta.url); const extensionDir = dirname(extensionFile); @@ -106,13 +107,13 @@ const sessionPointer = join(state, ".branch-session"); const mirrorCursorFile = join(state, ".branch-mirror-cursor"); const promptScript = join(fmRoot, "bin", "fm-branch-prompt.sh"); const outcomeScript = join(fmRoot, "bin", "fm-branch-outcome.sh"); -const operationalInputScript = join(fmRoot, "bin", "fm-operational-input.sh"); const leaseScript = join(fmRoot, "bin", "fm-lease.sh"); const wakeGrantScript = join(fmRoot, "bin", "fm-wake-grant.sh"); const loadedMarker = join(state, ".omp-branch-extension-loaded"); const branchExtensionVersion = `sha256:${createHash("sha256") .update(readFileSync(extensionFile)) .update(readFileSync(join(extensionDir, "lib", "fm-branch-dispatch.ts"))) + .update(readFileSync(join(extensionDir, "lib", "fm-async-exec.ts"))) .update(readFileSync(join(extensionDir, "lib", "fm-branch-model-picker.ts"))) .digest("hex")}`; const modelPinFile = join(config, "supervision-branch-model"); @@ -248,7 +249,13 @@ function offerEligible(offer: BranchDispatchOffer): boolean { return offer.eligible === true; } -function parentPid(pid: string): string { +async function parentPid(pid: string): Promise { + const result = await runCommandAsync("ps", ["-o", "ppid=", "-p", pid]); + if (result.status !== 0) return ""; + return result.stdout.trim(); +} + +function parentPidSync(pid: string): string { const result = spawnSync("ps", ["-o", "ppid=", "-p", pid], { encoding: "utf8" }); if (result.status !== 0) return ""; return result.stdout.trim(); @@ -267,43 +274,57 @@ let ownedLockPid = ""; // Same ownership read as the watcher adapter's lockOwnership(): the lock names // the harness pid, and this process owns it when that pid appears in its own -// ancestry. -function lockOwnership(): LockOwnership { +// ancestry. The awaited form keeps delivery off OMP's event-loop thread; the +// synchronous form is only for the offer handshake and shell spawn hook, both +// of which must decide without waiting. +const LOCK_ANCESTRY_DEPTH = 8; + +function readLockPid(): { lockPid: string; verdict: LockOwnership | null } { ownedLockPid = ""; let lockPid = ""; try { lockPid = readFileSync(`${state}/.lock`, "utf8").trim(); } catch { - return "missing"; + return { lockPid: "", verdict: "missing" }; + } + if (!/^[0-9]+$/.test(lockPid) || lockPid === "1") return { lockPid, verdict: "other" }; + return { lockPid, verdict: null }; +} + +function ownershipVerdict(lockPid: string, ancestryMatched: boolean): LockOwnership { + if (ancestryMatched) { + ownedLockPid = lockPid; + return "owned"; } - if (!/^[0-9]+$/.test(lockPid) || lockPid === "1") return "other"; + return pidAlive(lockPid) ? "other" : "missing"; +} + +async function lockOwnership(): Promise { + const { lockPid, verdict } = readLockPid(); + if (verdict) return verdict; let pid = String(process.pid); - for (let i = 0; i < 8; i += 1) { + for (let i = 0; i < LOCK_ANCESTRY_DEPTH; i += 1) { if (pid === lockPid) { - ownedLockPid = lockPid; - return "owned"; + const current = readLockPid(); + if (current.verdict || current.lockPid !== lockPid) return current.verdict ?? "other"; + return ownershipVerdict(lockPid, true); } - pid = parentPid(pid); + pid = await parentPid(pid); if (!pid || pid === "1") break; } - return pidAlive(lockPid) ? "other" : "missing"; + return ownershipVerdict(lockPid, false); } -// Encode a watcher wake as a marked operational injection through the bash -// owner, matching the primary OMP adapter's own fallback delivery. An encoding -// failure must never lose a wake, so the raw body is returned unmarked. -function encodeOperationalInput(content: string): string { - try { - const result = spawnSync(operationalInputScript, ["encode", "watcher"], { - encoding: "utf8", - input: content, - maxBuffer: 1024 * 1024, - }); - if (result.status === 0 && result.stdout) return result.stdout; - } catch { - // fall through to the unmarked body +function lockOwnershipSync(): LockOwnership { + const { lockPid, verdict } = readLockPid(); + if (verdict) return verdict; + let pid = String(process.pid); + for (let i = 0; i < LOCK_ANCESTRY_DEPTH; i += 1) { + if (pid === lockPid) return ownershipVerdict(lockPid, true); + pid = parentPidSync(pid); + if (!pid || pid === "1") break; } - return content; + return ownershipVerdict(lockPid, false); } function textOfContent(content: unknown): string { @@ -404,6 +425,20 @@ export default function (pi: ExtensionAPI) { // dispatch order, one at a time (the branch runs drain -> handle -> ack // serially by design). let branchChain: Promise = Promise.resolve(); + // Store writes, cursor advances, ownership reads, and visible merges are + // one ordered delivery stream. Awaiting subprocesses lets OMP repaint and + // accept input while a delivery runs; this queue preserves the ordering the + // old synchronous calls provided implicitly. + let deliveryChain: Promise = Promise.resolve(); + + function enqueueDelivery(unit: () => Promise): Promise { + const queued = deliveryChain.then(unit); + deliveryChain = queued.then( + () => {}, + () => {}, + ); + return queued; + } const pendingMirror: MirrorItem[] = []; const mirrorCollection: MirrorCollectionState = { collectAnchor: null, pendingCursor: null }; // Main's own current model and its live registry, tracked from the contexts @@ -492,12 +527,19 @@ export default function (pi: ExtensionAPI) { return model ? clampThinkingLevel(model, chosen) : chosen; } - function generationOwnsLock(expectedGeneration: number): boolean { - return !shuttingDown && expectedGeneration === generation && lockOwnership() === "owned"; + async function generationOwnsLock(expectedGeneration: number): Promise { + if (shuttingDown || expectedGeneration !== generation) return false; + const ownership = await lockOwnership(); + return !shuttingDown && expectedGeneration === generation && ownership === "owned"; + } + + function generationOwnsLockSync(expectedGeneration: number): boolean { + if (shuttingDown || expectedGeneration !== generation) return false; + return lockOwnershipSync() === "owned"; } - function markLoaded(): void { - if (lockOwnership() === "other") return; + async function markLoaded(): Promise { + if ((await lockOwnership()) === "other") return; const temporary = `${loadedMarker}.tmp.${process.pid}.${randomUUID()}`; try { mkdirSync(state, { recursive: true }); @@ -512,56 +554,48 @@ export default function (pi: ExtensionAPI) { // Clear any stray branch leases before a cold-start generation begins work. // One bulk release runs per generation, at activation. - function releaseBranchLeases(expectedGeneration: number): boolean { - if (!generationOwnsLock(expectedGeneration)) return false; - try { - const result = spawnSync("bash", [leaseScript, "release-actor", "--actor", "branch"], { - cwd: fmRoot, - encoding: "utf8", - env: { ...scriptEnv, FM_SUPERVISION_ACTOR: "branch" }, - }); - return result.status === 0; - } catch { - return false; - } + async function releaseBranchLeases(expectedGeneration: number): Promise { + if (!(await generationOwnsLock(expectedGeneration))) return false; + const result = await runCommandAsync("bash", [leaseScript, "release-actor", "--actor", "branch"], { + cwd: fmRoot, + env: { ...scriptEnv, FM_SUPERVISION_ACTOR: "branch" }, + }); + return result.status === 0; } // Lazy, per-action ownership evaluation (see the header). Returns true only // when this session owns the fleet lock right now; the first true evaluation // of a generation also writes the diagnostic marker and clears stray branch // leases from a prior generation. - function actingAsOwner(expectedGeneration = generation): boolean { - if (!generationOwnsLock(expectedGeneration)) return false; + async function actingAsOwner(expectedGeneration = generation): Promise { + if (!(await generationOwnsLock(expectedGeneration))) return false; if (activatedGeneration !== expectedGeneration) { - if (!releaseBranchLeases(expectedGeneration)) return false; - if (!generationOwnsLock(expectedGeneration)) return false; - if (!activateEligibleRowsOwner(state, wakeGrantScript, process.pid, String(expectedGeneration))) return false; - if (!generationOwnsLock(expectedGeneration)) { - rollbackEligibleRowsOwnerActivation(state, wakeGrantScript, process.pid, String(expectedGeneration)); + if (!(await releaseBranchLeases(expectedGeneration))) return false; + if (!(await generationOwnsLock(expectedGeneration))) return false; + if (!(await activateEligibleRowsOwner(state, wakeGrantScript, process.pid, String(expectedGeneration)))) { return false; } - markLoaded(); + if (!(await generationOwnsLock(expectedGeneration))) { + await rollbackEligibleRowsOwnerActivation(state, wakeGrantScript, process.pid, String(expectedGeneration)); + return false; + } + await markLoaded(); activatedGeneration = expectedGeneration; } return generationOwnsLock(expectedGeneration); } - function runOutcomeScript(args: string[]): { ok: boolean; stdout: string; detail: string } { - try { - const result = spawnSync("bash", [outcomeScript, ...args], { - cwd: fmRoot, - encoding: "utf8", - env: scriptEnv, - }); - if (result.status === 0) return { ok: true, stdout: (result.stdout || "").trim(), detail: "" }; - return { - ok: false, - stdout: "", - detail: `fm-branch-outcome.sh exited ${result.status ?? "none"}: ${(result.stderr || "").trim()}`, - }; - } catch (error) { - return { ok: false, stdout: "", detail: error instanceof Error ? error.message : String(error) }; - } + async function runOutcomeScript(args: string[]): Promise<{ ok: boolean; stdout: string; detail: string }> { + const result = await runCommandAsync("bash", [outcomeScript, ...args], { + cwd: fmRoot, + env: scriptEnv, + }); + if (result.status === 0) return { ok: true, stdout: (result.stdout || "").trim(), detail: "" }; + return { + ok: false, + stdout: "", + detail: `fm-branch-outcome.sh exited ${result.status ?? "none"}: ${(result.stderr || "").trim()}`, + }; } // Append-only merge into main. The store row is already durable when this @@ -581,17 +615,17 @@ export default function (pi: ExtensionAPI) { // the outcome durable in the store and recoverable through main's // fm_branch_outcomes tool, which is the strictly safer failure. The store row // is already durable when this runs (the report tool appends store-first). - function mergeIntoMain( + async function mergeIntoMain( expectedGeneration: number, seq: string, task: string, verdict: Verdict, summary: string, silent: boolean, - ): boolean { - if (!actingAsOwner(expectedGeneration)) return false; + ): Promise { + if (!(await actingAsOwner(expectedGeneration))) return false; // Advance the cursor first so the outcome can never be delivered twice. - if (/^[0-9]+$/.test(seq) && !runOutcomeScript(["handoff-next", "--seq", seq]).ok) { + if (/^[0-9]+$/.test(seq) && !(await runOutcomeScript(["handoff-next", "--seq", seq])).ok) { return false; } if (verdict === "captain") { @@ -656,32 +690,34 @@ export default function (pi: ExtensionAPI) { const verdict = verdictRaw as Verdict; const appendArgs = ["append", "--task", task, "--verdict", verdict, "--summary", summary, "--silent", String(silent)]; if (wake) appendArgs.push("--wake", wake); - if (!actingAsOwner(toolGeneration)) { - return { - content: [{ type: "text", text: "report refused: supervision session was replaced or lost lock ownership" }], - details: undefined, - isError: true, - }; - } - const appended = runOutcomeScript(appendArgs); - if (!appended.ok) { + return enqueueDelivery(async () => { + if (!(await actingAsOwner(toolGeneration))) { + return { + content: [{ type: "text", text: "report refused: supervision session was replaced or lost lock ownership" }], + details: undefined, + isError: true, + }; + } + const appended = await runOutcomeScript(appendArgs); + if (!appended.ok) { + return { + content: [{ type: "text", text: `outcome store append failed (nothing merged): ${appended.detail}` }], + details: undefined, + isError: true, + }; + } + if (!(await mergeIntoMain(toolGeneration, appended.stdout, task, verdict, summary, silent))) { + return { + content: [{ type: "text", text: `recorded seq ${appended.stdout}, but merge refused after supervision replacement or lock loss` }], + details: undefined, + isError: true, + }; + } return { - content: [{ type: "text", text: `outcome store append failed (nothing merged): ${appended.detail}` }], + content: [{ type: "text", text: `recorded seq ${appended.stdout} and merged [${verdict}] into main` }], details: undefined, - isError: true, }; - } - if (!mergeIntoMain(toolGeneration, appended.stdout, task, verdict, summary, silent)) { - return { - content: [{ type: "text", text: `recorded seq ${appended.stdout}, but merge refused after supervision replacement or lock loss` }], - details: undefined, - isError: true, - }; - } - return { - content: [{ type: "text", text: `recorded seq ${appended.stdout} and merged [${verdict}] into main` }], - details: undefined, - }; + }); }, }; } @@ -709,9 +745,8 @@ export default function (pi: ExtensionAPI) { // authoritative whenever a fresh process creates or reopens the branch. const pinned = branchModelSelection(); const effort = branchEffortSelection(pinned?.model); - const prompt = spawnSync("bash", [promptScript], { + const prompt = await runCommandAsync("bash", [promptScript], { cwd: fmRoot, - encoding: "utf8", env: scriptEnv, maxBuffer: 4 * 1024 * 1024, }); @@ -720,7 +755,7 @@ export default function (pi: ExtensionAPI) { `fm-branch-prompt.sh did not produce a usable branch prompt (status=${prompt.status ?? "none"}): ${(prompt.stderr || "").trim()}`, ); } - if (!actingAsOwner(branchGeneration)) throw new Error("supervision session was replaced or lost lock ownership"); + if (!(await actingAsOwner(branchGeneration))) throw new Error("supervision session was replaced or lost lock ownership"); mkdirSync(sessionsDir, { recursive: true }); let sessionManager: SessionManager | null = null; try { @@ -737,11 +772,11 @@ export default function (pi: ExtensionAPI) { suppressBreadcrumb: true, }); } - if (!actingAsOwner(branchGeneration)) throw new Error("supervision session was replaced or lost lock ownership"); + if (!(await actingAsOwner(branchGeneration))) throw new Error("supervision session was replaced or lost lock ownership"); const leaseHolderPid = ownedLockPid; const bashTool = createBashToolDefinition(fmRoot, { spawnHook: (context) => { - if (!actingAsOwner(branchGeneration)) { + if (activatedGeneration !== branchGeneration || !generationOwnsLockSync(branchGeneration)) { throw new Error("bash refused: supervision session was replaced or lost lock ownership"); } return { @@ -791,7 +826,7 @@ ${context.command} ...(pinned ? { model: pinned.model } : {}), ...(effort === undefined ? {} : { thinkingLevel: effort }), }); - if (!actingAsOwner(branchGeneration)) { + if (!(await actingAsOwner(branchGeneration))) { try { await created.session.dispose(); } catch {} @@ -806,12 +841,12 @@ ${context.command} } async function ensureBranch(expectedGeneration: number): Promise { - if (!actingAsOwner(expectedGeneration)) throw new Error("supervision session was replaced or lost lock ownership"); + if (!(await actingAsOwner(expectedGeneration))) throw new Error("supervision session was replaced or lost lock ownership"); if (branch) return branch; if (branchBroken) throw new Error(branchBroken); try { const created = await createBranch(expectedGeneration); - if (!actingAsOwner(expectedGeneration)) { + if (!(await actingAsOwner(expectedGeneration))) { try { await created.dispose(); } catch {} @@ -828,19 +863,19 @@ ${context.command} } async function flushMirror(session: AgentSession, expectedGeneration: number): Promise { - if (!actingAsOwner(expectedGeneration)) throw new Error("supervision session no longer owns the fleet lock"); + if (!(await actingAsOwner(expectedGeneration))) throw new Error("supervision session no longer owns the fleet lock"); while (pendingMirror.length > 0) { const item = pendingMirror[0]; - if (!actingAsOwner(expectedGeneration)) throw new Error("supervision session no longer owns the fleet lock"); + if (!(await actingAsOwner(expectedGeneration))) throw new Error("supervision session no longer owns the fleet lock"); await session.sendCustomMessage( { customType: "fm-main-mirror", content: `[${item.tag}] ${item.text}`, display: false }, {}, ); - if (!actingAsOwner(expectedGeneration)) throw new Error("supervision session was replaced during mirror delivery"); + if (!(await actingAsOwner(expectedGeneration))) throw new Error("supervision session was replaced during mirror delivery"); pendingMirror.shift(); } if (mirrorCollection.pendingCursor) { - if (!actingAsOwner(expectedGeneration)) throw new Error("supervision session no longer owns the fleet lock"); + if (!(await actingAsOwner(expectedGeneration))) throw new Error("supervision session no longer owns the fleet lock"); writeMirrorCursor(mirrorCollection.pendingCursor); mirrorCollection.pendingCursor = null; } @@ -852,10 +887,14 @@ ${context.command} if (shuttingDown || acceptedGeneration !== generation) { throw new Error("supervision session was replaced before handling the accepted wake"); } - if (!actingAsOwner(acceptedGeneration)) throw new Error("supervision session no longer owns the fleet lock"); + if (!(await enqueueDelivery(() => actingAsOwner(acceptedGeneration)))) { + throw new Error("supervision session no longer owns the fleet lock"); + } const session = await ensureBranch(acceptedGeneration); await flushMirror(session, acceptedGeneration); - if (!actingAsOwner(acceptedGeneration)) throw new Error("supervision session no longer owns the fleet lock"); + if (!(await enqueueDelivery(() => actingAsOwner(acceptedGeneration)))) { + throw new Error("supervision session no longer owns the fleet lock"); + } const heartbeat = /^heartbeat($|:)/.test(message); const scope = scopeForUnreadWake(state, heartbeat); // A newly-arrived main-owned (check-kind) row never bounces this whole @@ -869,7 +908,7 @@ ${context.command} if (scope.corrupted) { throw new Error("the unread wake queue could not be read safely"); } - const grant = writeEligibleRowsSnapshot(state, scope.eligibleSeqs, wakeGrantScript, String(acceptedGeneration)); + const grant = await writeEligibleRowsSnapshot(state, scope.eligibleSeqs, wakeGrantScript, String(acceptedGeneration)); if (grant === "main-owned") throw new Error("the wake rows are already claimed by main"); if (grant !== "published") throw new Error("could not record the branch's eligible row snapshot"); // A row can still arrive between this re-check and the model starting @@ -877,12 +916,12 @@ ${context.command} await session.prompt( `FIRSTMATE SUPERVISION WAKE: ${message}\n\nHandle this per your operating procedure. Do not finish this turn until you have completed, in order: fm_branch_report, the exact WAKE_ACK_REQUIRED command, and release of every task lease you claimed.`, ); - if (!releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration))) { + if (!(await releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration)))) { throw new Error("could not release the branch's settled wake-row grant"); } }) - .catch((error: unknown) => { - releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration)); + .catch(async (error: unknown) => { + await releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration)); throw error; }); branchChain = delivery.catch(() => {}); @@ -895,7 +934,7 @@ ${context.command} const flushSession = branch; branchChain = branchChain .then(async () => { - if (!actingAsOwner(flushGeneration)) return; + if (!(await actingAsOwner(flushGeneration))) return; await flushMirror(flushSession, flushGeneration); }) .catch(() => { @@ -911,7 +950,7 @@ ${context.command} // gets neither branch routing nor branch-owned state/lease cleanup side // effects. if (!offerEligible(offer)) return; - if (!actingAsOwner()) return; // cold start pre-lock, secondary session, or shutdown + if (!generationOwnsLockSync(generation)) 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)); @@ -928,9 +967,11 @@ ${context.command} // the volatile queue, then deliver it through the serialized chain so it // lands before any later wake. The durable cursor advances only in // flushMirror after the complete pending batch reaches the branch. - pi.on?.("turn_end", (_event, ctx) => { + pi.on?.("turn_end", async (_event, ctx) => { rememberMainContext(ctx); - if (!actingAsOwner()) return; + const turnGeneration = generation; + if (!(await enqueueDelivery(() => actingAsOwner(turnGeneration)))) return; + if (turnGeneration !== generation) return; try { pendingMirror.push(...collectMainDialog(ctx.sessionManager, mirrorCollection)); } catch { @@ -946,12 +987,13 @@ ${context.command} // 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) => { + pi.on?.("session_start", async (_event, ctx) => { rememberMainContext(ctx); shuttingDown = false; branchBroken = ""; generation += 1; - actingAsOwner(generation); + const startedGeneration = generation; + await enqueueDelivery(() => actingAsOwner(startedGeneration)); }); // Terminal quit: latch shutdown so no further wake or mirror turn starts. No @@ -1132,7 +1174,7 @@ ${context.command} execute: async (_toolCallId, params) => { const recentRaw = (params as { recent?: unknown }).recent; const recent = typeof recentRaw === "number" && recentRaw >= 1 ? String(Math.floor(recentRaw)) : "20"; - const listed = runOutcomeScript(["list", "--recent", recent]); + const listed = await enqueueDelivery(() => runOutcomeScript(["list", "--recent", recent])); if (!listed.ok) { return { content: [{ type: "text", text: `could not read the outcome store: ${listed.detail}` }], @@ -1146,5 +1188,5 @@ ${context.command} }; }, }); - markLoaded(); + void markLoaded(); } diff --git a/.omp/extensions/lib/fm-async-exec.ts b/.omp/extensions/lib/fm-async-exec.ts new file mode 100644 index 00000000000..521798ddbbb --- /dev/null +++ b/.omp/extensions/lib/fm-async-exec.ts @@ -0,0 +1,90 @@ +import { spawn } from "node:child_process"; + +// OMP extensions, their tools, and their event handlers share the JavaScript +// thread that drives the active session. A synchronous child-process call in +// supervision delivery therefore stalls prompt handling and rendering for the +// child's whole lifetime. +// +// This helper owns the awaited-spawn replacement. It preserves the status, +// UTF-8 stdout, and UTF-8 stderr shape callers consumed from spawnSync, while +// yielding the OMP event loop until the child exits and its streams drain. +// Ordering that must remain serialized belongs to the caller's queue. + +export interface AsyncExecResult { + /** Exit code, or null when the child was signalled or never started. */ + status: number | null; + stdout: string; + stderr: string; +} + +export interface AsyncExecOptions { + cwd?: string; + env?: NodeJS.ProcessEnv; + /** Written to the child's stdin, which is closed either way. */ + input?: string; + /** Upper bound on each captured output stream. */ + maxBuffer?: number; +} + +const DEFAULT_MAX_BUFFER = 1024 * 1024; + +export function runCommandAsync( + command: string, + args: readonly string[], + options: AsyncExecOptions = {}, +): Promise { + return new Promise((resolve) => { + let stdout = ""; + let stderr = ""; + let stdoutBytes = 0; + let stderrBytes = 0; + const maxBuffer = options.maxBuffer ?? DEFAULT_MAX_BUFFER; + let settled = false; + const finish = (status: number | null, detail = ""): void => { + if (settled) return; + settled = true; + resolve({ status, stdout, stderr: detail ? `${stderr}${detail}` : stderr }); + }; + let child; + try { + child = spawn(command, [...args], { + cwd: options.cwd, + env: options.env, + stdio: ["pipe", "pipe", "pipe"], + }); + } catch (error) { + finish(null, error instanceof Error ? error.message : String(error)); + return; + } + child.stdout?.setEncoding("utf8"); + child.stdout?.on("data", (chunk: string) => { + if (settled) return; + const bytes = Buffer.byteLength(chunk, "utf8"); + if (stdoutBytes + bytes > maxBuffer) { + child.kill(); + finish(null, `stdout exceeded ${maxBuffer} bytes`); + return; + } + stdout += chunk; + stdoutBytes += bytes; + }); + child.stderr?.setEncoding("utf8"); + child.stderr?.on("data", (chunk: string) => { + if (settled) return; + const bytes = Buffer.byteLength(chunk, "utf8"); + if (stderrBytes + bytes > maxBuffer) { + child.kill(); + finish(null, `stderr exceeded ${maxBuffer} bytes`); + return; + } + stderr += chunk; + stderrBytes += bytes; + }); + child.on("close", (code) => finish(code)); + child.on("error", (error: Error) => finish(null, error.message)); + if (child.stdin) { + child.stdin.on("error", () => {}); + child.stdin.end(options.input ?? ""); + } + }); +} diff --git a/.omp/extensions/lib/fm-branch-dispatch.ts b/.omp/extensions/lib/fm-branch-dispatch.ts index 3bcdaf9a5a3..4aa42829602 100644 --- a/.omp/extensions/lib/fm-branch-dispatch.ts +++ b/.omp/extensions/lib/fm-branch-dispatch.ts @@ -1,5 +1,5 @@ -import { spawnSync } from "node:child_process"; import { readdirSync, readFileSync } from "node:fs"; +import { runCommandAsync } from "./fm-async-exec.ts"; // Shared wake-dispatch handshake between the OMP watcher adapter (the // dispatcher) and the supervision-branch extension (the handler), carried over @@ -164,56 +164,59 @@ export const BRANCH_ELIGIBLE_ROWS_FILE = ".branch-eligible-rows"; // actor acquired the requested rows. export type EligibleRowsSnapshotResult = "published" | "main-owned" | "error"; -function runGrantScript(state: string, grantScript: string, args: readonly string[]): number | null { - try { - const result = spawnSync("bash", [grantScript, ...args], { - encoding: "utf8", - env: { - ...process.env, - FM_STATE_OVERRIDE: state, - FM_WAKE_QUEUE: `${state}/.wake-queue`, - FM_WAKE_QUEUE_LOCK: `${state}/.wake-queue.lock`, - }, - }); - return result.status; - } catch { - return null; - } +async function runGrantScript( + state: string, + grantScript: string, + args: readonly string[], +): Promise { + const result = await runCommandAsync("bash", [grantScript, ...args], { + env: { + ...process.env, + FM_STATE_OVERRIDE: state, + FM_WAKE_QUEUE: `${state}/.wake-queue`, + FM_WAKE_QUEUE_LOCK: `${state}/.wake-queue.lock`, + }, + }); + return result.status; } -export function activateEligibleRowsOwner( +export async function activateEligibleRowsOwner( state: string, grantScript: string, ownerPid: number, generation: string, -): boolean { - return runGrantScript(state, grantScript, ["activate", String(ownerPid), generation]) === 0; +): Promise { + return (await runGrantScript(state, grantScript, ["activate", String(ownerPid), generation])) === 0; } -export function writeEligibleRowsSnapshot( +export async function writeEligibleRowsSnapshot( state: string, seqs: readonly string[], grantScript: string, generation: string, -): EligibleRowsSnapshotResult { +): Promise { if (seqs.length === 0 || seqs.some((seq) => !/^[0-9]+$/.test(seq))) return "error"; - const status = runGrantScript(state, grantScript, ["publish", generation, ...seqs]); + const status = await runGrantScript(state, grantScript, ["publish", generation, ...seqs]); if (status === 0) return "published"; if (status === 3) return "main-owned"; return "error"; } -export function releaseEligibleRowsSnapshot(state: string, grantScript: string, generation: string): boolean { - return runGrantScript(state, grantScript, ["release", generation]) === 0; +export async function releaseEligibleRowsSnapshot( + state: string, + grantScript: string, + generation: string, +): Promise { + return (await runGrantScript(state, grantScript, ["release", generation])) === 0; } -export function rollbackEligibleRowsOwnerActivation( +export async function rollbackEligibleRowsOwnerActivation( state: string, grantScript: string, ownerPid: number, generation: string, -): boolean { - return runGrantScript(state, grantScript, ["rollback-activation", String(ownerPid), generation]) === 0; +): Promise { + return (await runGrantScript(state, grantScript, ["rollback-activation", String(ownerPid), generation])) === 0; } export interface BranchDispatchOffer { diff --git a/bin/fm-primary-watch-version-lib.sh b/bin/fm-primary-watch-version-lib.sh index 5834697ad03..0fe8bbf4086 100644 --- a/bin/fm-primary-watch-version-lib.sh +++ b/bin/fm-primary-watch-version-lib.sh @@ -17,7 +17,7 @@ # caller already treats as "not loaded" rather than as a match. # # The OMP supervision branch marker hashes its extension followed by the branch -# dispatch and model-picker helpers, matching its in-process producer. +# dispatch, async-exec, and model-picker helpers, matching its in-process producer. # Missing files or a missing SHA-256 utility return nonzero, which callers treat # as "not loaded" rather than as a match. # @@ -53,11 +53,12 @@ fm_primary_watch_version() { # -> sha256: } fm_omp_branch_extension_version() { # -> sha256: - local extension=${1:-} root=${2:-} dispatch picker + local extension=${1:-} root=${2:-} dispatch async_exec picker [ -n "$extension" ] && [ -f "$extension" ] || return 1 [ -n "$root" ] || return 1 dispatch="$root/.omp/extensions/lib/fm-branch-dispatch.ts" + async_exec="$root/.omp/extensions/lib/fm-async-exec.ts" picker="$root/.omp/extensions/lib/fm-branch-model-picker.ts" - [ -f "$dispatch" ] && [ -f "$picker" ] || return 1 - fm_extension_version_hash "$extension" "$dispatch" "$picker" + [ -f "$dispatch" ] && [ -f "$async_exec" ] && [ -f "$picker" ] || return 1 + fm_extension_version_hash "$extension" "$dispatch" "$async_exec" "$picker" } diff --git a/docs/omp-supervision-branch.md b/docs/omp-supervision-branch.md index b29233a86f2..0fd3ea3f36a 100644 --- a/docs/omp-supervision-branch.md +++ b/docs/omp-supervision-branch.md @@ -121,6 +121,15 @@ The branch can also run on a cheaper model and a shallower reasoning effort than Away mode carries over unchanged: while `state/.afk` exists the away daemon owns supervision, and the branch declines every wake offer for the duration. What is new is only the attended path: outside away mode, the branch absorbs the routine majority that previously interrupted the captain's conversation, applying the same escalation etiquette the daemon applies while away. +## Responsiveness port reconciliation + +Upstream Pi change #3767 moves supervision outcome subprocesses off Pi's render thread and serializes delivery so durable append, cursor handoff, and visible delivery remain ordered. +The shared bash contracts, including `bin/fm-supervise-daemon.sh`, did not change in #3767 because the away daemon already runs in its own process and does not block the OMP session event loop. +This fork ports the mechanism into `.omp/extensions/lib/fm-async-exec.ts`, makes the OMP grant-script helpers awaitable, adds the OMP branch's `deliveryChain` around ownership checks, outcome append/handoff, and report delivery, and includes the new helper in `bin/fm-primary-watch-version-lib.sh`'s branch marker hash. +OMP retains synchronous ownership reads only where its offer handshake and shell spawn hook must answer without waiting, while every delivery-side process call rechecks ownership through the awaited path. +The existing watcher-continuity implementation in `bin/fm-primary-watch-core.ts` remains the delivery owner for rejected branch settlements and replacement handoffs; this port does not bypass or duplicate that PR #99 contract. +Pi-only processed-outcome reconciliation and Pi renderer/live-TUI guards have no OMP equivalent and are deliberately omitted; OMP's direct append-and-handoff outcome store is covered by the portable OMP regression and strict OMP typecheck. + ## Verification 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. diff --git a/tests/fm-omp-branch-live-e2e.test.sh b/tests/fm-omp-branch-live-e2e.test.sh index a0299baf926..99d28e0d83f 100755 --- a/tests/fm-omp-branch-live-e2e.test.sh +++ b/tests/fm-omp-branch-live-e2e.test.sh @@ -39,6 +39,7 @@ repo="$TMP_ROOT/repo" mkdir -p "$repo/.omp/extensions/lib" cp "$ROOT/.omp/extensions/fm-branch-supervision-omp.ts" "$repo/.omp/extensions/fm-branch-supervision-omp.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$repo/.omp/extensions/lib/fm-branch-dispatch.ts" +cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$repo/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" "$repo/.omp/extensions/lib/fm-branch-model-picker.ts" ln -s "$OMP_NODE_MODULES" "$repo/node_modules" diff --git a/tests/fm-omp-branch-supervision.test.sh b/tests/fm-omp-branch-supervision.test.sh index e2d4ea2bf01..77ad152915f 100755 --- a/tests/fm-omp-branch-supervision.test.sh +++ b/tests/fm-omp-branch-supervision.test.sh @@ -354,6 +354,7 @@ test_omp_extension_establishes_main_actor_context() { mkdir -p "$fixture/.omp/extensions/lib" "$package_dir" cp "$ROOT/.omp/extensions/fm-branch-supervision-omp.ts" "$fixture/.omp/extensions/fm-branch-supervision-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-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" "$fixture/.omp/extensions/lib/fm-branch-model-picker.ts" printf '%s\n' '{"type":"module","exports":{".":"./index.js","./extensibility/legacy-pi-coding-agent-shim":"./coding-shim.js","./extensibility/legacy-pi-ai-shim":"./ai-shim.js"}}' > "$package_dir/package.json" printf '%s\n' 'export function createAgentSession() {} export class SessionManager {}' > "$package_dir/index.js" @@ -372,6 +373,39 @@ test_omp_extension_establishes_main_actor_context() { pass "OMP extension load establishes the main supervision actor context" } +test_async_outcome_delivery_keeps_event_loop_responsive_and_ordered() { + command -v bun >/dev/null 2>&1 || { echo "skip: bun not found for OMP async delivery regression"; return 0; } + local out ticks deliveries + out=$(bun - "$ROOT/.omp/extensions/lib/fm-async-exec.ts" <<'JS' +const { runCommandAsync } = await import(process.argv[2]); +let deliveryChain = Promise.resolve(); +const deliveries = []; +const enqueue = (label, delay) => { + const queued = deliveryChain.then(async () => { + const result = await runCommandAsync("bash", ["-c", `sleep ${delay}; printf '%s' '${label}'`]); + if (result.status !== 0) throw new Error(`delivery failed for ${label}`); + deliveries.push(result.stdout); + }); + deliveryChain = queued.then(() => {}, () => {}); + return queued; +}; +let ticks = 0; +const timer = setInterval(() => { ticks += 1; }, 5); +await Promise.all([enqueue("routine", "0.12"), enqueue("captain", "0.02")]); +clearInterval(timer); +if (ticks < 10) throw new Error(`event loop stalled during async outcome delivery: only ${ticks} timer ticks`); +if (deliveries.join(",") !== "routine,captain") throw new Error(`delivery order changed: ${deliveries.join(",")}`); +console.log(`TICKS=${ticks}`); +console.log(`DELIVERIES=${deliveries.join(",")}`); +JS + ) || fail "async outcome delivery regression harness failed: $out" + ticks=$(printf '%s\n' "$out" | sed -n 's/^TICKS=//p') + deliveries=$(printf '%s\n' "$out" | sed -n 's/^DELIVERIES=//p') + [ "${ticks:-0}" -ge 10 ] || fail "async outcome delivery stalled the event loop: $out" + [ "$deliveries" = "routine,captain" ] || fail "async outcome delivery reordered outcomes: $out" + pass "OMP outcome delivery yields to the event loop while preserving routine/captain order" +} + test_branch_prompt_is_byte_stable_and_above_cache_floor test_outcome_store_is_append_only_with_cursor_reads test_outcome_startup_replay_preserves_silence @@ -382,3 +416,4 @@ test_send_lease_guard_serializes_a_held_task test_pr_check_lease_guard_serializes_task_metadata test_non_branch_home_is_untouched test_omp_extension_establishes_main_actor_context +test_async_outcome_delivery_keeps_event_loop_responsive_and_ordered diff --git a/tests/fm-omp-branch-types.test.sh b/tests/fm-omp-branch-types.test.sh index 78ab7bebe45..e3c05a2d516 100755 --- a/tests/fm-omp-branch-types.test.sh +++ b/tests/fm-omp-branch-types.test.sh @@ -40,6 +40,7 @@ cp "$ROOT/.omp/extensions/fm-branch-supervision-omp.ts" "$TMP_ROOT/.omp/extensio cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$TMP_ROOT/.omp/extensions/fm-primary-omp.ts" cp "$ROOT/.omp/extensions/fm-fleet-hooks.ts" "$TMP_ROOT/.omp/extensions/fm-fleet-hooks.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$TMP_ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" +cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$TMP_ROOT/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$TMP_ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" "$TMP_ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" cp "$ROOT/bin/fm-primary-watch-core.ts" "$TMP_ROOT/bin/fm-primary-watch-core.ts" diff --git a/tests/fm-omp-primary-live-e2e.test.sh b/tests/fm-omp-primary-live-e2e.test.sh index ca337bc45a4..817813a2289 100755 --- a/tests/fm-omp-primary-live-e2e.test.sh +++ b/tests/fm-omp-primary-live-e2e.test.sh @@ -50,6 +50,7 @@ cp -R "$ROOT/.agents/skills/harness-adapters" "$PROJECT/.agents/skills/harness-a cp -R "$ROOT/.agents/skills/afk" "$PROJECT/.agents/skills/afk" cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$PROJECT/.omp/extensions/fm-primary-omp.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$PROJECT/.omp/extensions/lib/fm-branch-dispatch.ts" +cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$PROJECT/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$PROJECT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" git init -q -b main "$PROJECT" fm_git_identity fmtest fmtest@example.invalid @@ -224,6 +225,8 @@ cp "$ROOT/.omp/extensions/fm-primary-omp.ts" \ "$FALLBACK_PROJECT/.omp/extensions/fm-primary-omp.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" \ "$FALLBACK_PROJECT/.omp/extensions/lib/fm-branch-dispatch.ts" +cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" \ + "$FALLBACK_PROJECT/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" \ "$FALLBACK_PROJECT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" git init -q -b main "$FALLBACK_PROJECT" diff --git a/tests/fm-omp-primary.test.sh b/tests/fm-omp-primary.test.sh index fdc7e26a995..12765900e52 100755 --- a/tests/fm-omp-primary.test.sh +++ b/tests/fm-omp-primary.test.sh @@ -216,6 +216,7 @@ test_primary_scope_requires_canonical_state() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" mkdir -p "$fixture/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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-primary-scope-lib.sh" "$fixture/bin/fm-primary-scope-lib.sh" @@ -261,6 +262,7 @@ test_native_omp_fresh_checkout_nudges_once() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" mkdir -p "$fixture/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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-primary-scope-lib.sh" "$fixture/bin/fm-primary-scope-lib.sh" @@ -349,6 +351,7 @@ test_primary_marker_refuses_whitespace_identity() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" mkdir -p "$fixture/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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-primary-scope-lib.sh" "$fixture/bin/fm-primary-scope-lib.sh" @@ -390,6 +393,7 @@ test_native_primary_extension_contract() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" mkdir -p "$fixture/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$fixture/.omp/extensions/lib/fm-task-inbox-doorbell.ts" chmod +x "$fixture/.omp/extensions/fm-primary-omp.ts" cp "$ROOT/bin/fm-primary-watch-core.ts" "$fixture/bin/fm-primary-watch-core.ts" @@ -790,6 +794,7 @@ test_native_omp_confirms_recovery_handling_delivery() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" mkdir -p "$fixture/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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-primary-scope-lib.sh" "$fixture/bin/fm-primary-scope-lib.sh" @@ -895,6 +900,7 @@ test_native_omp_refused_handling_delivery_is_typed_once() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$fixture/.omp/extensions/fm-primary-omp.ts" mkdir -p "$fixture/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$fixture/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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-primary-scope-lib.sh" "$fixture/bin/fm-primary-scope-lib.sh" @@ -981,6 +987,7 @@ test_native_omp_session_switch_carries_inflight_actionable_close() { 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-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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" @@ -1117,6 +1124,7 @@ test_native_omp_unacknowledged_wake_keeps_successor_chain() { 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-async-exec.ts" "$fixture/.omp/extensions/lib/fm-async-exec.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" diff --git a/tests/fm-omp-secondmate.test.sh b/tests/fm-omp-secondmate.test.sh index 1f2070bc355..7623f225012 100755 --- a/tests/fm-omp-secondmate.test.sh +++ b/tests/fm-omp-secondmate.test.sh @@ -40,6 +40,7 @@ setup_case() { # : > "$HERDR_LOG" cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$HOME_DIR/.omp/extensions/fm-primary-omp.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$HOME_DIR/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$HOME_DIR/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$HOME_DIR/.omp/extensions/lib/fm-task-inbox-doorbell.ts" cp "$ROOT/AGENTS.md" "$HOME_DIR/AGENTS.md" ln -s "$ROOT/bin" "$HOME_DIR/bin" @@ -52,7 +53,7 @@ setup_case() { # git -C "$HOME_DIR" config user.email fmtest@example.com printf 'home\n' > "$HOME_DIR/README.md" git -C "$HOME_DIR" add README.md .omp/extensions/fm-primary-omp.ts \ - .omp/extensions/lib/fm-branch-dispatch.ts .omp/extensions/lib/fm-task-inbox-doorbell.ts + .omp/extensions/lib/fm-branch-dispatch.ts .omp/extensions/lib/fm-async-exec.ts .omp/extensions/lib/fm-task-inbox-doorbell.ts git -C "$HOME_DIR" commit -qm init cat > "$FAKEBIN/omp" <<'JS' diff --git a/tests/fm-send-turn-start.test.sh b/tests/fm-send-turn-start.test.sh index 1728c382484..3f8cf25bf37 100755 --- a/tests/fm-send-turn-start.test.sh +++ b/tests/fm-send-turn-start.test.sh @@ -359,10 +359,12 @@ test_remote_control_uses_task_bound_omp_route() { "$ROOT/bin/fm-primary-watch-core.ts" "$root/bin/" cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$root/.omp/extensions/" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" \ + "$ROOT/.omp/extensions/lib/fm-async-exec.ts" \ "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$root/.omp/extensions/lib/" cp "$root/bin/fm-primary-watch-core.ts" "$home/bin/" cp "$root/.omp/extensions/fm-primary-omp.ts" "$home/.omp/extensions/" cp "$root/.omp/extensions/lib/fm-branch-dispatch.ts" \ + "$root/.omp/extensions/lib/fm-async-exec.ts" \ "$root/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$home/.omp/extensions/lib/" cat > "$root/bin/fm-backend.sh" <<'SH' fm_backend_validate_task_endpoint() { diff --git a/tests/fm-session-start.test.sh b/tests/fm-session-start.test.sh index 4a8a2f82e4f..53645847dce 100755 --- a/tests/fm-session-start.test.sh +++ b/tests/fm-session-start.test.sh @@ -674,6 +674,7 @@ install_omp_primary_extension_fixture() { cp "$ROOT/.omp/extensions/fm-primary-omp.ts" "$root/.omp/extensions/fm-primary-omp.ts" cp "$ROOT/.omp/extensions/fm-branch-supervision-omp.ts" "$root/.omp/extensions/fm-branch-supervision-omp.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$root/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$root/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$root/.omp/extensions/lib/fm-task-inbox-doorbell.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" "$root/.omp/extensions/lib/fm-branch-model-picker.ts" install_primary_watch_core_fixture "$root" @@ -1773,6 +1774,7 @@ JS branch_file="$root/.omp/extensions/fm-branch-supervision-omp.ts" dispatch_file="$root/.omp/extensions/lib/fm-branch-dispatch.ts" + async_exec_file="$root/.omp/extensions/lib/fm-async-exec.ts" picker_file="$root/.omp/extensions/lib/fm-branch-model-picker.ts" printf '\nexport const staleBranchExtensionFixture = true;\n' >> "$branch_file" out=$(run_session_start "$home" "$root" "$fakebin:$BASE_PATH") @@ -1788,6 +1790,7 @@ JS *) failure=${failure:-"session start trusted a branch marker after its dispatch helper changed"} ;; esac cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$dispatch_file" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$async_exec_file" printf '\nexport const staleBranchPickerFixture = true;\n' >> "$picker_file" out=$(run_session_start "$home" "$root" "$fakebin:$BASE_PATH") case "$out" in diff --git a/tests/fm-spawn-dispatch-profile.test.sh b/tests/fm-spawn-dispatch-profile.test.sh index b8befcb6585..4ec2caa661b 100755 --- a/tests/fm-spawn-dispatch-profile.test.sh +++ b/tests/fm-spawn-dispatch-profile.test.sh @@ -2876,6 +2876,7 @@ test_omp_secondmate_inspects_staged_live_extensions() { cp "$ROOT/.omp/extensions/fm-branch-supervision-omp.ts" "$sm/.omp/extensions/fm-branch-supervision-omp.ts" mkdir -p "$sm/.omp/extensions/lib" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$sm/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$sm/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$sm/.omp/extensions/lib/fm-task-inbox-doorbell.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" "$sm/.omp/extensions/lib/fm-branch-model-picker.ts" git -C "$sm" init -q @@ -2914,6 +2915,7 @@ test_omp_secondmate_rejects_modified_branch_helper() { cp "$ROOT/.omp/extensions/fm-fleet-hooks.ts" "$sm/.omp/extensions/fm-fleet-hooks.ts" cp "$ROOT/.omp/extensions/fm-branch-supervision-omp.ts" "$sm/.omp/extensions/fm-branch-supervision-omp.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-dispatch.ts" "$sm/.omp/extensions/lib/fm-branch-dispatch.ts" + cp "$ROOT/.omp/extensions/lib/fm-async-exec.ts" "$sm/.omp/extensions/lib/fm-async-exec.ts" cp "$ROOT/.omp/extensions/lib/fm-task-inbox-doorbell.ts" "$sm/.omp/extensions/lib/fm-task-inbox-doorbell.ts" cp "$ROOT/.omp/extensions/lib/fm-branch-model-picker.ts" "$sm/.omp/extensions/lib/fm-branch-model-picker.ts" git -C "$sm" init -q From 6d6b7cb48cf11ecd8f7b3757f80a029a54590504 Mon Sep 17 00:00:00 2001 From: dnth Date: Sat, 5 Sep 2026 18:22:49 +0800 Subject: [PATCH 2/2] no-mistakes(document): Clarify async supervision delivery documentation --- docs/omp-supervision-branch.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/omp-supervision-branch.md b/docs/omp-supervision-branch.md index 0fd3ea3f36a..dce6f1297f4 100644 --- a/docs/omp-supervision-branch.md +++ b/docs/omp-supervision-branch.md @@ -47,8 +47,8 @@ Three capabilities are deliberately out of scope for this port and are a future - 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. -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. +The reason is structural: OMP coordinates the two in-process actors through shared filesystem locks (the wake-queue lock and the lease-command lock), and the offer handshake and shell-spawn hook still use 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, so it cannot safely perform a live takeover; the delivery path itself uses awaited subprocesses and does not block OMP's event loop. 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 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. @@ -132,7 +132,7 @@ Pi-only processed-outcome reconciliation and Pi renderer/live-TUI guards have no ## Verification -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. +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, non-branch-home invariance, and responsive ordered async outcome delivery; `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).