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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
278 changes: 160 additions & 118 deletions .omp/extensions/fm-branch-supervision-omp.ts

Large diffs are not rendered by default.

90 changes: 90 additions & 0 deletions .omp/extensions/lib/fm-async-exec.ts
Original file line number Diff line number Diff line change
@@ -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<AsyncExecResult> {
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 ?? "");
}
});
}
57 changes: 30 additions & 27 deletions .omp/extensions/lib/fm-branch-dispatch.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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<number | null> {
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<boolean> {
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<EligibleRowsSnapshotResult> {
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<boolean> {
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<boolean> {
return (await runGrantScript(state, grantScript, ["rollback-activation", String(ownerPid), generation])) === 0;
}

export interface BranchDispatchOffer {
Expand Down
9 changes: 5 additions & 4 deletions bin/fm-primary-watch-version-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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.
#
Expand Down Expand Up @@ -53,11 +53,12 @@ fm_primary_watch_version() { # <extension-file> <fm-root> -> sha256:<hex>
}

fm_omp_branch_extension_version() { # <extension-file> <fm-root> -> sha256:<hex>
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"
}
15 changes: 12 additions & 3 deletions docs/omp-supervision-branch.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -121,9 +121,18 @@ 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.
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).
1 change: 1 addition & 0 deletions tests/fm-omp-branch-live-e2e.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
35 changes: 35 additions & 0 deletions tests/fm-omp-branch-supervision.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
Expand All @@ -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
1 change: 1 addition & 0 deletions tests/fm-omp-branch-types.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
3 changes: 3 additions & 0 deletions tests/fm-omp-primary-live-e2e.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"
Expand Down
Loading
Loading