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
181 changes: 141 additions & 40 deletions .pi/extensions/fm-primary-pi-watch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,17 @@
// state/extensions/pi-primary-watch/session-replacement-actionable.json.
// Terminal quit leaves the final generation stopped so late callbacks cannot rearm.
// Stale callbacks from a prior generation are no-ops against the active replacement.
//
// Delivery versus consumption (stated once here):
// A main follow-up is delivered once Pi accepts it (sendUserMessage resolves).
// The successor pipeline never waits for the model to read it: a follow-up
// queued while main is streaming joins the running run without ever raising
// before_agent_start, so waiting on that event stalls every later close.
// Consumption is tracked only so a replacement can replay a follow-up Pi had
// not consumed. An idle main consumes at before_agent_start; a streaming main
// consumes at the user message_start carrying the exact wake text; either
// event finishes the pending record, and a still-unconsumed record rides the
// replacement handoff.
import { spawn, spawnSync, type ChildProcess } from "node:child_process";
import { createHash } from "node:crypto";
import { mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs";
Expand Down Expand Up @@ -66,6 +77,11 @@ type WatchToolRenderContext = {
isPartial: boolean;
};

type UnconsumedWake = {
content: string;
pending: PendingActionableClose;
};

type SessionGeneration = {
id: number;
stopping: boolean;
Expand All @@ -78,7 +94,15 @@ type SessionGeneration = {
seq: number;
pendingActionables: PendingActionableClose[];
cleanupFailure: string;
wakeAcknowledgements: Map<string, { content: string; settle: (consumed: boolean) => void }>;
// Main follow-ups Pi has accepted but not yet consumed, by pending token.
// Never cleared at shutdown: a delivery continuation that runs after the
// replacement began reads it to tell a main-queued wake (replayed) from a
// branch-handled one (finished).
unconsumedWakes: Map<string, UnconsumedWake>;
// A verified successor's failure close that arrived while the pipeline was
// still delivering the wake it was started for; its bounded retry runs once
// that delivery settles instead of being skipped by the single-flight guard.
deferredClose: { message: string; predecessorArmPid: string } | null;
};

function refreshWatchToolShell(
Expand Down Expand Up @@ -145,19 +169,25 @@ type ReplacementCoordinatorGlobal = typeof globalThis & {
__firstmatePiWatchReplacements?: Map<string, ReplacementCoordinator>;
};
const replacementCoordinatorGlobal = globalThis as ReplacementCoordinatorGlobal;
const replacementCoordinators = replacementCoordinatorGlobal.__firstmatePiWatchReplacements ??= new Map();
let replacementCoordinator = replacementCoordinators.get(actionableHandoff);
if (!replacementCoordinator) {
replacementCoordinator = {
const replacementCoordinators = replacementCoordinatorGlobal.__firstmatePiWatchReplacements ??= new Map<string, ReplacementCoordinator>();
function replacementCoordinatorFor(handoff: string): ReplacementCoordinator {
const existing = replacementCoordinators.get(handoff);
if (existing) return existing;
const created: ReplacementCoordinator = {
receiver: null,
pending: [],
nextTokenId: 0,
deliveries: new Map(),
};
replacementCoordinators.set(actionableHandoff, replacementCoordinator);
replacementCoordinators.set(handoff, created);
return created;
}
const replacementCoordinator = replacementCoordinatorFor(actionableHandoff);
const armReadiness = new WeakMap<ChildProcess, Promise<boolean>>();
const armClose = new WeakMap<ChildProcess, Promise<void>>();
// Children the extension itself asked to exit; their close is not a failure
// of the successor and never earns a deferred retry.
const armRetired = new WeakSet<ChildProcess>();
const armRecovery = new WeakMap<ChildProcess, { generation: string; watcherPid: string }>();
const armPendingActionable = new WeakMap<ChildProcess, PendingActionableClose>();

Expand Down Expand Up @@ -215,6 +245,24 @@ function completedActionableLine(output: string): string {
return newline < 0 ? "" : actionableLine(output.slice(0, newline + 1));
}

// The text Pi carries in a user message_start: sendUserMessage wraps a string
// as one text part, so the joined text parts equal the sent content.
function userMessageText(content: unknown): string {
if (typeof content === "string") return content;
if (!Array.isArray(content)) return "";
const parts: string[] = [];
for (const part of content) {
if (
typeof part === "object" && part !== null &&
(part as { type?: unknown }).type === "text" &&
typeof (part as { text?: unknown }).text === "string"
) {
parts.push((part as { text: string }).text);
}
}
return parts.join("\n");
}

function nodeErrorCode(error: unknown): string {
return typeof error === "object" && error !== null && "code" in error
? String((error as { code?: unknown }).code ?? "")
Expand Down Expand Up @@ -372,7 +420,8 @@ function createGeneration(): SessionGeneration {
seq: 0,
pendingActionables: [],
cleanupFailure: "",
wakeAcknowledgements: new Map(),
unconsumedWakes: new Map(),
deferredClose: null,
};
}

Expand Down Expand Up @@ -465,30 +514,41 @@ export default function (pi: ExtensionAPI) {
async function sendWake(
owner: SessionGeneration,
message: string,
token?: string,
pending?: PendingActionableClose,
): Promise<boolean> {
if (!generationIsLive(owner)) return false;
const content = encodeFirstmateOperationalInput(
"watcher",
`FIRSTMATE WATCHER WAKE: ${message}\n\nRun bin/fm-wake-drain.sh first and handle the queued wake. Watcher continuity is extension-owned.`,
);
if (!token) {
await pi.sendUserMessage(content, { deliverAs: "followUp" });
return generationIsLive(owner);
}
let settleConsumption: (consumed: boolean) => void = () => {};
const consumption = new Promise<boolean>((resolveConsumption) => {
settleConsumption = resolveConsumption;
});
owner.wakeAcknowledgements.set(token, { content, settle: settleConsumption });
if (pending) owner.unconsumedWakes.set(pending.token, { content, pending });
try {
await pi.sendUserMessage(content, { deliverAs: "followUp" });
return await consumption;
} catch (error) {
owner.wakeAcknowledgements.delete(token);
settleConsumption(false);
if (pending) owner.unconsumedWakes.delete(pending.token);
throw error;
}
// Accepted by Pi. A generation replaced while Pi was accepting it may
// have lost the follow-up with the old session, so report it undelivered
// and let the replacement replay the still-pending record.
return generationIsLive(owner);
}

// Pi consumed a main follow-up: an idle main at before_agent_start, a
// streaming main at the user message_start that joins the running run.
function consumeWake(owner: SessionGeneration, text: string): void {
for (const [token, wake] of owner.unconsumedWakes) {
if (wake.content !== text) continue;
owner.unconsumedWakes.delete(token);
wake.pending.delivered = true;
try {
finishPendingActionable(owner, wake.pending);
} catch (error) {
surfaceCleanupFailure(owner, error);
schedulePendingCleanup(owner);
}
return;
}
}

function confirmHandlingDelivery(recovery: { generation: string; watcherPid: string }): {
Expand Down Expand Up @@ -556,7 +616,7 @@ export default function (pi: ExtensionAPI) {
owner: SessionGeneration,
message: string,
repairFailed: boolean,
token: string,
pending: PendingActionableClose,
recovery?: { generation: string; watcherPid: string },
): Promise<boolean> {
if (!generationIsLive(owner)) return false;
Expand All @@ -567,7 +627,7 @@ export default function (pi: ExtensionAPI) {
if (!pidAlive(watcherPid)) {
await retireArm(owner.child);
}
return await sendWake(owner, `${message}\n\n${confirmed.detail}`, token);
return await sendWake(owner, `${message}\n\n${confirmed.detail}`, pending);
}
}
if (!repairFailed) {
Expand All @@ -579,7 +639,7 @@ export default function (pi: ExtensionAPI) {
} catch {}
}
}
return await sendWake(owner, message, token);
return await sendWake(owner, message, pending);
}

function surfaceFailure(owner: SessionGeneration, message: string): void {
Expand Down Expand Up @@ -654,7 +714,11 @@ export default function (pi: ExtensionAPI) {
surfaceCleanupFailure(owner, error);
}
}
const pending = owner.pendingActionables.find((item) => !item.delivered);
// A record Pi has accepted but not consumed is neither redelivered
// nor finished here: consumption finishes it, replacement replays it.
const pending = owner.pendingActionables.find(
(item) => !item.delivered && !owner.unconsumedWakes.has(item.token),
);
if (!pending) break;
const existingClaim = replacementCoordinator.deliveries.get(pending.token);
if (existingClaim && existingClaim.owner !== owner) {
Expand All @@ -680,25 +744,40 @@ export default function (pi: ExtensionAPI) {
}
};
try {
// A new restoration supersedes whatever became of the previous
// successor; only a failure during this delivery is retried after it.
owner.deferredClose = null;
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);
const delivered = await deliverActionableWake(owner, message, Boolean(restoration.failure), pending, restoration.recovery);
if (!delivered) {
settleClaim("failed");
releaseClaim();
return;
}
pending.delivered = true;
const awaitingConsumption = owner.unconsumedWakes.has(pending.token);
if (awaitingConsumption && !generationIsLive(owner)) {
// Pi accepted the follow-up, then the session was replaced before
// this continuation ran: the shutdown persisted the still-pending
// record, so a replacement waiting on this claim must replay it.
settleClaim("failed");
releaseClaim();
return;
}
settleClaim("delivered");
try {
finishPendingActionable(owner, pending);
} catch (error) {
surfaceCleanupFailure(owner, error);
if (!awaitingConsumption) {
// The branch handled it, or Pi consumed it before this ran.
pending.delivered = true;
try {
finishPendingActionable(owner, pending);
} catch (error) {
surfaceCleanupFailure(owner, error);
}
}
releaseClaim();
} catch (error) {
Expand All @@ -714,7 +793,19 @@ export default function (pi: ExtensionAPI) {
if (generationIsLive(owner)) {
owner.restoring = false;
if (owner.pendingActionables.some((pending) => pending.delivered)) schedulePendingCleanup(owner);
if (!owner.child && !owner.retryTimer) startArm(owner);
// No bare arm is launched here. A generation without a child at this
// point has either delivered a typed restoration failure after its
// bounded retries, which hands repair to main through fm_watch_arm_pi
// (one more silent launch past the bound could hold a hung child that
// the repair call would then report as "unchanged"), or lost a
// verified successor during the delivery, which takes the ordinary
// bounded, lock-checked retry it would have taken had the pipeline
// been idle.
const deferred = owner.deferredClose;
owner.deferredClose = null;
if (deferred && !owner.child && !owner.retryTimer) {
scheduleRetry(owner, deferred.message, deferred.predecessorArmPid);
}
}
}
}
Expand Down Expand Up @@ -751,6 +842,7 @@ export default function (pi: ExtensionAPI) {

async function retireArm(armChild: ChildProcess | null): Promise<boolean> {
if (!armChild) return true;
armRetired.add(armChild);
armChild.kill("SIGTERM");
const closed = armClose.get(armChild);
if (!closed) return false;
Expand Down Expand Up @@ -861,6 +953,7 @@ export default function (pi: ExtensionAPI) {
let stderr = "";
let settled = false;
let readinessSettled = false;
let verified = false;
let resolveReadiness: (ready: boolean) => void = () => {};
let resolveClosed: () => void = () => {};
const readiness = new Promise<boolean>((resolveReady) => {
Expand All @@ -874,6 +967,7 @@ export default function (pi: ExtensionAPI) {
const settleReadiness = (ready: boolean): void => {
if (readinessSettled) return;
readinessSettled = true;
verified = ready;
resolveReadiness(ready);
};
const observeEstablishedArm = (): void => {
Expand Down Expand Up @@ -917,7 +1011,17 @@ export default function (pi: ExtensionAPI) {
void processPendingActionables(owner);
return;
}
if (!generationIsLive(owner) || owner.restoring) return;
if (!generationIsLive(owner)) return;
if (owner.restoring) {
// The pipeline is still delivering the wake this successor was
// started for. A verified successor that failed on its own keeps its
// bounded retry for the end of that delivery; an unready child closing
// here was retired by the restoration itself.
if (verified && !armRetired.has(armChild)) {
owner.deferredClose = { message: classification.message, predecessorArmPid: predecessor };
}
return;
}
scheduleRetry(owner, classification.message, predecessor);
});
armChild.on("error", (error: Error) => {
Expand Down Expand Up @@ -967,12 +1071,11 @@ export default function (pi: ExtensionAPI) {
}

pi.on?.("before_agent_start", (event) => {
for (const [token, acknowledgement] of generation.wakeAcknowledgements) {
if (acknowledgement.content !== event.prompt) continue;
generation.wakeAcknowledgements.delete(token);
acknowledgement.settle(true);
break;
}
consumeWake(generation, event.prompt);
});
pi.on?.("message_start", (event) => {
if (event.message.role !== "user") return;
consumeWake(generation, userMessageText(event.message.content));
});

pi.on?.("session_start", async () => {
Expand All @@ -984,8 +1087,6 @@ export default function (pi: ExtensionAPI) {
});
pi.on?.("session_shutdown", async (event) => {
const replacement = event.reason === "reload" || event.reason === "new" || event.reason === "resume" || event.reason === "fork";
for (const acknowledgement of generation.wakeAcknowledgements.values()) acknowledgement.settle(false);
generation.wakeAcknowledgements.clear();
if (replacementCoordinator.receiver === receiveReplacementActionable) replacementCoordinator.receiver = null;
await stopSessionGeneration(generation, replacement);
});
Expand Down
2 changes: 1 addition & 1 deletion docs/pi-supervision-branch.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ This feature is Pi-only by construction and changes nothing anywhere else:
A co-present main-owned check row no longer defers that review to main, because it is not fleet context the branch is missing and main is woken for it on its own triggering close.
- The branch itself: `.pi/extensions/fm-branch-supervision.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 replacement or lock loss cannot let an old continuation mutate the new session.
Every accepted path that cannot reach a working branch rejects its settlement to the watcher, which retains delivery ownership and routes the wake through its consumption-acknowledged main path; a broken branch declines later offers so they take that path directly.
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 as a follow-up that counts as delivered once Pi accepts it; a broken branch declines later offers so they take that path directly.
After wake rows are claimed, a branch prompt counts as handled only when `fm_branch_report` appends a durable outcome before that prompt settles; a settled provider error or a settled prompt with no report releases the grant and rejects delivery ownership back to the watcher.
Two consecutive settled provider errors latch the branch broken and surface a one-line health note only on that initial trip.
Main keeps every wake during a five-minute cooldown, after which one wake may probe the branch while concurrent wakes still stay on main; each probe that settles with another provider error doubles the next cooldown up to one hour.
Expand Down
Loading
Loading