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
69 changes: 19 additions & 50 deletions .omp/extensions/fm-branch-supervision-omp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//
Expand Down Expand Up @@ -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<void> {
const delivery = branchChain
.then(async () => {
if (shuttingDown || acceptedGeneration !== generation) {
throw new Error("supervision session was replaced before handling the accepted wake");
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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", () => {
Expand All @@ -969,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;
Expand Down
18 changes: 9 additions & 9 deletions .omp/extensions/fm-primary-omp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> | 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({
Expand Down Expand Up @@ -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 = "";
Expand Down Expand Up @@ -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", {
Expand Down
12 changes: 8 additions & 4 deletions .omp/extensions/lib/fm-branch-dispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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<void>;
accept(settlement?: Promise<void>): void;
}

export function createBranchDispatchOffer(
Expand All @@ -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;
Expand Down
9 changes: 7 additions & 2 deletions .pi/extensions/fm-primary-pi-watch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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", {
Expand Down
Loading
Loading