From 88df9480099cc91bea19e212ffc9ad296a543727 Mon Sep 17 00:00:00 2001 From: kunchenguid Date: Wed, 2 Sep 2026 02:12:17 -0700 Subject: [PATCH 1/2] fix(pi): settle watcher delivery on Pi accepting the follow-up A follow-up queued while main is streaming joins the running run without ever raising before_agent_start, so waiting on that event before clearing the successor pipeline (#3498) stalled every later actionable close: no successor started, no wake was delivered or offered to the branch, and the turn-end guard woke main to re-arm by hand after every close. The pipeline now settles once Pi accepts the follow-up. Consumption is observed at before_agent_start for an idle main and at the user message_start for a streaming main, and decides only what a replacement session (/new, /resume, /fork, reload) replays. An exhausted restoration delivers its typed failure without launching an arm past the retry bound, which the stall had hidden. The replacement-coordinator map is typed so the strict no-emit typecheck passes again. Tests: the doubles no longer raise before_agent_start for a streaming send, a portable regression drives two actionable closes while main streams and proves the successor chain plus consumption-scoped replay, and a credential-free real-SDK probe pins Pi's event contract for both the streaming and the idle follow-up. Claude-Session: https://claude.ai/code/session_01QJjTsUvKkWAwLGNoncaZ3a --- .pi/extensions/fm-primary-pi-watch.ts | 147 +++++++++++++----- docs/pi-supervision-branch.md | 2 +- docs/verification/runtime-backends.md | 24 ++- docs/verification/supervision.md | 3 + docs/watcher-continuity.md | 5 +- tests/fm-pi-branch-live-e2e.test.sh | 214 ++++++++++++++++++++++++++ tests/fm-pi-watch-extension.test.sh | 116 +++++++++++++- 7 files changed, 467 insertions(+), 44 deletions(-) diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index b10fbd82d5a..aa2aba3a01a 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -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"; @@ -66,6 +77,11 @@ type WatchToolRenderContext = { isPartial: boolean; }; +type UnconsumedWake = { + content: string; + pending: PendingActionableClose; +}; + type SessionGeneration = { id: number; stopping: boolean; @@ -78,7 +94,11 @@ type SessionGeneration = { seq: number; pendingActionables: PendingActionableClose[]; cleanupFailure: string; - wakeAcknowledgements: Map 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; }; function refreshWatchToolShell( @@ -145,17 +165,20 @@ type ReplacementCoordinatorGlobal = typeof globalThis & { __firstmatePiWatchReplacements?: Map; }; 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(); +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>(); const armClose = new WeakMap>(); const armRecovery = new WeakMap(); @@ -215,6 +238,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 ?? "") @@ -372,7 +413,7 @@ function createGeneration(): SessionGeneration { seq: 0, pendingActionables: [], cleanupFailure: "", - wakeAcknowledgements: new Map(), + unconsumedWakes: new Map(), }; } @@ -465,30 +506,41 @@ export default function (pi: ExtensionAPI) { async function sendWake( owner: SessionGeneration, message: string, - token?: string, + pending?: PendingActionableClose, ): Promise { 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((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 }): { @@ -556,7 +608,7 @@ export default function (pi: ExtensionAPI) { owner: SessionGeneration, message: string, repairFailed: boolean, - token: string, + pending: PendingActionableClose, recovery?: { generation: string; watcherPid: string }, ): Promise { if (!generationIsLive(owner)) return false; @@ -567,7 +619,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) { @@ -579,7 +631,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 { @@ -654,7 +706,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) { @@ -687,18 +743,30 @@ export default function (pi: ExtensionAPI) { 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) { @@ -714,7 +782,11 @@ 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 arm is launched here. A generation without a child at this point + // has delivered a typed restoration failure after its bounded retries, + // and that message 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". } } } @@ -967,12 +1039,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 () => { @@ -984,8 +1055,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); }); diff --git a/docs/pi-supervision-branch.md b/docs/pi-supervision-branch.md index 44120bdf76a..e03dc694236 100644 --- a/docs/pi-supervision-branch.md +++ b/docs/pi-supervision-branch.md @@ -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. diff --git a/docs/verification/runtime-backends.md b/docs/verification/runtime-backends.md index 024ca44c018..a4e63fd586e 100644 --- a/docs/verification/runtime-backends.md +++ b/docs/verification/runtime-backends.md @@ -1086,9 +1086,31 @@ ok - real Pi SDK 0.84.4 returns a post-construction 429 wake to main without los ``` The current portable regression proves that only consecutive provider errors count toward the two-error broken-branch latch: a durable report between errors resets the streak, the error that reaches the threshold rejects to watcher-owned fallback, and the next wake remains on main without another branch prompt. -`tests/fm-pi-watch-extension.test.sh` owns the provider-free integration evidence that watcher fallback remains pending until main consumption or successful branch settlement. +`tests/fm-pi-watch-extension.test.sh` owns the provider-free integration evidence that watcher fallback remains pending until Pi accepts the main follow-up or the branch settles successfully, and that a follow-up accepted while main is streaming neither stalls the successor chain nor escapes replacement replay until Pi consumes it. [`pi-supervision-branch.md`](../pi-supervision-branch.md) owns the current cooldown, recovery, and re-latch contract and points to the regression that now covers it. Scope of the earlier evidence: the installed signed `pi` CLI (0.82.0 at verification time) is a compiled binary whose bundled SDK is not importable from Node, so the importable npm package is the only surface the guard and the typecheck can pin. The extension executes inside the signed CLI's own runtime, so a CLI upgrade can drift ahead of the pinned npm surface; refresh the SDK construction, picker, renderer, and type evidence after every Pi upgrade by rerunning the applicable live guard probes, picker regression, and strict typecheck above (point `FM_PI_PACKAGE_DIR` at a matching npm install when one exists). The live guard now drives both extensions through the watcher-owned settlement handshake, requires rejected branch settlement before main delivery, and verifies successor-delivery confirmation; rerun it against the matching importable Pi package to refresh end-to-end fallback evidence. + +### 2026-09-02 streaming-time watcher delivery + +The focused watcher suite, strict typecheck, and credential-free live guard were run against the npm `@earendil-works/pi-coding-agent` 0.84.4 package selected with `FM_PI_PACKAGE_DIR`, on macOS 26.6.2 arm64, Node v24.14.1, after the watcher extension stopped waiting for `before_agent_start` before settling a main delivery. +No credential was read, no request left the machine, and the active Pi session was not changed. + +```sh +bin/fm-test-run.sh tests/fm-pi-watch-extension.test.sh +FM_PI_PACKAGE_DIR= npm exec --yes --package=typescript@5.9.3 -- bash tests/fm-pi-primary-types.test.sh +FM_PI_BRANCH_LIVE_E2E=1 FM_PI_PACKAGE_DIR= bin/fm-test-run.sh tests/fm-pi-branch-live-e2e.test.sh +``` + +```text +ok - Pi hung successor falls back to one typed actionable wake +ok - Pi streaming-time wake delivery keeps the successor chain and replays only unconsumed wakes +ok - tracked Pi extensions pass strict no-emit typecheck against Pi 0.84.4 +ok - real Pi SDK 0.84.4 queues a streaming-time watcher wake without before_agent_start, keeps the successor chain, and surfaces consumption of both follow-ups +``` + +The live probe loads the tracked watcher extension through Pi's real resource loader into a real AgentSession whose only provider is a local fake with its fetch intercepted in-process and held open mid-stream. +It proved that a follow-up the extension sends while main is streaming raises no `before_agent_start` at queue time or when the run reaches it, joins the run as a user `message_start` carrying the exact wake text in its own model turn, and is followed by a verified successor and delivery of the next close; a follow-up sent to the idle main raises `before_agent_start` with the exact text before its user `message_start`. +The portable regression drives the same shape with a fake main that never raises `before_agent_start` while streaming, then proves a replacement replays only the follow-up Pi had not consumed and that an exhausted restoration delivers its typed failure without launching a further arm. diff --git a/docs/verification/supervision.md b/docs/verification/supervision.md index 3e3002ad096..d45e123479b 100644 --- a/docs/verification/supervision.md +++ b/docs/verification/supervision.md @@ -473,6 +473,9 @@ Stale prior-generation tool callbacks could not mutate the active child, repeate The strict no-emit check used the installed Pi SDK declarations to hold the lifecycle event contract. Plain Pi and pi-signed share the same tracked `.pi/extensions/fm-primary-pi-watch.ts` path, so both inherit the generation owner; other primary harnesses are not applicable because they do not use this Pi extension lifecycle. +On 2026-09-02 the same suite, the strict typecheck, and the credential-free real-SDK guard were rerun against `@earendil-works/pi-coding-agent` 0.84.4 after the extension stopped waiting for `before_agent_start` before settling a main delivery; [`runtime-backends.md`](runtime-backends.md#2026-09-02-streaming-time-watcher-delivery) owns the exact commands and output. +Observed guarantee: a wake delivered while main was streaming was followed by a verified successor and by delivery of the next actionable close, a replacement replayed only the follow-up Pi had not consumed, and an exhausted restoration delivered its typed failure without launching an arm past the retry bound. + The once-per-generation recovery bound and immediate handling-successor poll were verified on 2026-08-21 with the tracked Pi extension, real watcher processes, and an isolated home. The regression forced handling confirmation to fail, observed one recovery follow-up across the former repeat window, confirmed the successor remained live, and then proved a separate handling successor durably queued a crew event within the bounded poll window. diff --git a/docs/watcher-continuity.md b/docs/watcher-continuity.md index d12a77152b4..537a236d842 100644 --- a/docs/watcher-continuity.md +++ b/docs/watcher-continuity.md @@ -8,7 +8,8 @@ Must-work continuity now lives above that process boundary instead of depending Pi's `.pi/extensions/fm-primary-pi-watch.ts` and OpenCode's `.opencode/plugins/fm-primary-watch-arm.js` own continuous re-arm after an actionable child close. Each adapter starts the next arm before delivering the wake prompt, checks current session-lock ownership at launch, preserves one child or scheduled retry at a time, and applies bounded exponential retry after an unexpected or failed close. A failed follow-up never cancels continuity restoration. -Pi same-process session replacement follows the generation-owner contract in `.pi/extensions/fm-primary-pi-watch.ts`: an owning `session_start` arms the replacement generation without waiting for a model turn, and a state-scoped replacement handoff carries every actionable close whose delivery overlapped `session_shutdown`, including a main follow-up not yet consumed by `before_agent_start`, branch handling, and a retiring child that reports after the bounded shutdown wait. +Pi same-process session replacement follows the generation-owner contract in `.pi/extensions/fm-primary-pi-watch.ts`: an owning `session_start` arms the replacement generation without waiting for a model turn, and a state-scoped replacement handoff carries every actionable close whose delivery overlapped `session_shutdown`, including a main follow-up Pi accepted but had not yet consumed, branch handling, and a retiring child that reports after the bounded shutdown wait. +A main follow-up counts as delivered once Pi accepts it, never once the model reads it, because a follow-up queued while main is streaming joins the running run without a `before_agent_start`; the extension header owns how consumption is observed and why it only decides what a replacement replays. Cursor's `.cursor/hooks.json` `stop` hook (`bin/fm-turnend-guard-cursor.sh`) owns routine tokenless re-arm for a Cursor primary by parking that awaited hook on `bin/fm-watch-arm.sh` and returning an actionable close as one follow-up; [`turnend-guard.md`](turnend-guard.md#harness-integrations) owns its Pi-host stand-down, loop bounds, and supersession baton. Claude's `.claude/settings.json` Stop `asyncRewake` hook (`bin/fm-claude-stop-autoarm.sh`) owns routine tokenless re-arm. The hook fires on every Stop, and an eligible primary with supervision need admits one home-scoped owner that foregrounds `bin/fm-watch-arm.sh` inside the hook-owned process tree. @@ -72,7 +73,7 @@ A main drain validates that owner evidence under the queue lock and reclaims the A main drain claims every currently unclaimed row and excludes an active branch grant from both presentation and acknowledgement. Its `--ack-through ` deletes only claimed main rows at or below the cutoff, while a branch acknowledgement deletes only claimed branch rows at or below its cutoff. Every settled branch prompt releases any residual grant, so an omitted or failed acknowledgement leaves the durable row available to a later main drain; a successful acknowledgement has already removed it. -If a branch offer loses the claim race to main, it rejects its settlement so the watcher retains the actionable close until its consumption-acknowledged main follow-up begins. +If a branch offer loses the claim race to main, it rejects its settlement so the watcher retains the actionable close until Pi accepts its main follow-up. [`pi-supervision-branch.md`](pi-supervision-branch.md#components-and-their-owners) owns branch eligibility, mixed-queue dispatch, the pre-drain recheck, and heartbeat's all-or-nothing rule. A check-kind row is main-owned in every mode, including a heartbeat review, so it is never part of a branch claim and never defers one; main is woken for it on that check's own triggering close. `fm-wake-drain.sh` never reclassifies a row itself: it filters the queue to the current actor's opaque claim before same-key deduplication, then presents and acknowledges only that actor-local view. diff --git a/tests/fm-pi-branch-live-e2e.test.sh b/tests/fm-pi-branch-live-e2e.test.sh index 6eeec79556c..98d64070149 100644 --- a/tests/fm-pi-branch-live-e2e.test.sh +++ b/tests/fm-pi-branch-live-e2e.test.sh @@ -129,11 +129,18 @@ const pi = { registerCommand() {}, registerMessageRenderer() {}, sendMessage() {}, + // Main is idle throughout this probe, so a send starts a run: Pi raises + // before_agent_start with the exact text and then the user message_start. + // A send while main streams raises neither at queue time; the sixth probe + // below proves that against the real AgentSession. async sendUserMessage(content, options) { mainUserMessages.push({ content, options: options ?? {} }); for (const handler of piHandlers.get("before_agent_start") ?? []) { await handler({ prompt: content }, sessionCtx); } + for (const handler of piHandlers.get("message_start") ?? []) { + await handler({ message: { role: "user", content: [{ type: "text", text: content }] } }, sessionCtx); + } }, }; process.env.FM_ROOT_OVERRIDE = process.env.FM_REAL_ROOT; @@ -331,11 +338,16 @@ const pi = { registerCommand() {}, registerMessageRenderer() {}, sendMessage() {}, + // Idle main, as in the first probe: a send starts a run and Pi raises + // before_agent_start, then the user message_start. async sendUserMessage(content, options) { mainUserMessages.push({ content, options: options ?? {} }); for (const handler of piHandlers.get("before_agent_start") ?? []) { await handler({ prompt: content }, sessionCtx); } + for (const handler of piHandlers.get("message_start") ?? []) { + await handler({ message: { role: "user", content: [{ type: "text", text: content }] } }, sessionCtx); + } }, getThinkingLevel() { return "off"; @@ -778,3 +790,205 @@ if [ "$status" -ne 0 ] || [ "$out" != "DELIVERY_OK" ]; then fail "real-SDK visible outcome delivery guard failed against pi-coding-agent $PI_VERSION: $out" fi pass "real Pi SDK $PI_VERSION immediately renders appendEntry in the active transcript, persists it across reopen, and excludes it from model context" + +# Sixth probe: the vendor event contract watcher continuity rests on, against +# the real AgentSession and ExtensionRunner with the tracked watcher extension +# loaded through Pi's own resource loader. A wake the extension delivers while +# main is streaming must join the running run without ever raising +# before_agent_start, the extension must still start the successor and deliver +# the next close, and Pi must surface consumption of both the streaming-time +# and the idle follow-up through the events the extension reads (the user +# message_start, and before_agent_start for the idle one). The provider is a +# local fake whose only fetch is intercepted in-process and held open until +# the follow-up is queued, so no request leaves the machine and no credential +# is read. +streamdir="$TMP_ROOT/stream-agent-dir" +streamhome="$TMP_ROOT/stream-home" +mkdir -p "$streamdir" "$streamhome/state" "$streamhome/config" "$TMP_ROOT/stream-sessions" +cat > "$streamdir/models.json" <<'JSON' +{ + "providers": { + "fm-live-stream": { + "baseUrl": "https://fm-live-stream.invalid/v1", + "api": "openai-completions", + "apiKey": "fm-live-placeholder", + "models": [ + { "id": "fm-live-stream-model", "name": "fm live stream", "contextWindow": 8192, "maxTokens": 512 } + ] + } + } +} +JSON +WATCH_PLUGIN="$repo/.pi/extensions/fm-primary-pi-watch.ts" \ + FM_HOME="$streamhome" FM_ROOT_OVERRIDE="$repo" \ + FM_LIVE_WATCH_LOG="$TMP_ROOT/stream-watch.log" FM_LIVE_WATCH_TRIGGER="$TMP_ROOT/stream-watch.trigger" \ + FM_LIVE_SESSIONS="$TMP_ROOT/stream-sessions" \ + PI_CODING_AGENT_DIR="$streamdir" PI_PACKAGE_DIR="$PI_PACKAGE_DIR" \ + node --input-type=module > "$TMP_ROOT/stream-output" 2>&1 <<'EOF' +import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { resolve } from "node:path"; +import { pathToFileURL } from "node:url"; + +const home = resolve(process.env.FM_HOME); +writeFileSync(`${home}/state/.lock`, `${process.pid}\n`); +const pkg = resolve(process.env.PI_PACKAGE_DIR); +const { DefaultResourceLoader, ModelRegistry, ModelRuntime, SessionManager, SettingsManager, createAgentSession } = + await import(pathToFileURL(`${pkg}/dist/index.js`).href); + +// The local fake provider: the first completion streams one token and then +// holds its stream open until the probe releases it; later ones finish at once. +let completions = 0; +let releaseStream = () => {}; +const streamHeld = new Promise((release) => { + releaseStream = release; +}); +const chunk = (delta, finish) => `data: ${JSON.stringify({ + id: "fm-live-stream", + object: "chat.completion.chunk", + created: 1, + model: "fm-live-stream-model", + choices: [{ index: 0, delta, finish_reason: finish }], + ...(finish ? { usage: { prompt_tokens: 1, completion_tokens: 1, total_tokens: 2 } } : {}), +})}\n\n`; +globalThis.fetch = async (input) => { + const url = typeof input === "string" ? input : input instanceof URL ? input.href : input.url; + if (!url.startsWith("https://fm-live-stream.invalid/")) { + throw new Error(`unexpected network request in provider-free guard: ${url}`); + } + completions += 1; + const hold = completions === 1 ? streamHeld : Promise.resolve(); + const encoder = new TextEncoder(); + const body = new ReadableStream({ + async start(controller) { + controller.enqueue(encoder.encode(chunk({ role: "assistant", content: "OK" }, null))); + await hold; + controller.enqueue(encoder.encode(chunk({}, "stop"))); + controller.enqueue(encoder.encode("data: [DONE]\n\n")); + controller.close(); + }, + }); + return new Response(body, { status: 200, headers: { "content-type": "text/event-stream" } }); +}; + +const events = []; +const userText = (content) => typeof content === "string" + ? content + : content.filter((part) => part.type === "text").map((part) => part.text).join("\n"); +const agentDir = resolve(process.env.PI_CODING_AGENT_DIR); +const settings = SettingsManager.create(process.cwd(), agentDir); +const loader = new DefaultResourceLoader({ + cwd: process.cwd(), + agentDir, + settingsManager: settings, + additionalExtensionPaths: [process.env.WATCH_PLUGIN], + extensionFactories: [{ + name: "fm-event-contract-probe", + factory: (pi) => { + pi.on("before_agent_start", (event) => { + events.push({ type: "before_agent_start", text: event.prompt }); + }); + pi.on("message_start", (event) => { + if (event.message.role !== "user") return; + events.push({ type: "user_message_start", text: userText(event.message.content) }); + }); + pi.on("input", (event) => { + events.push({ type: "input", text: event.text, source: event.source, streamingBehavior: event.streamingBehavior }); + }); + pi.on("agent_settled", () => { + events.push({ type: "agent_settled" }); + }); + }, + }], + noSkills: true, + noPromptTemplates: true, + noThemes: true, + noContextFiles: true, +}); +await loader.reload(); +const runtime = await ModelRuntime.create({ + authPath: `${agentDir}/auth.json`, + modelsPath: `${agentDir}/models.json`, +}); +const registry = new ModelRegistry(runtime); +await registry.refresh(); +const model = registry.find("fm-live-stream", "fm-live-stream-model"); +if (!model) throw new Error("the real registry did not resolve the local streaming model"); +const { session } = await createAgentSession({ + cwd: process.cwd(), + sessionManager: SessionManager.create(process.cwd(), resolve(process.env.FM_LIVE_SESSIONS)), + settingsManager: settings, + resourceLoader: loader, + modelRuntime: runtime, + model, + noTools: "builtin", +}); + +const armLog = process.env.FM_LIVE_WATCH_LOG; +const armRows = (prefix) => existsSync(armLog) + ? readFileSync(armLog, "utf8").split(/\n/).filter((line) => line.startsWith(prefix)).length + : 0; +const has = (type, text) => events.some((event) => event.type === type && (text === undefined || event.text.includes(text))); +const settledRuns = () => events.filter((event) => event.type === "agent_settled").length; +const waitFor = async (predicate, label) => { + for (let i = 0; i < 600; i += 1) { + const failure = events.find((event) => event.type === "prompt_error"); + if (failure) throw new Error(`the real session rejected its prompt: ${failure.text}`); + if (predicate()) return; + await new Promise((tick) => setTimeout(tick, 50)); + } + throw new Error(`timeout waiting for ${label}; events=${JSON.stringify(events)}`); +}; + +const armTool = session.getToolDefinition("fm_watch_arm_pi"); +if (!armTool) throw new Error("the real tool registry did not expose fm_watch_arm_pi"); +const armed = await armTool.execute("live-stream-arm", {}, undefined, undefined, {}); +if (!armed.details?.ok) throw new Error(`watcher did not arm: ${JSON.stringify(armed.details)}`); +await waitFor(() => armRows("arm ") === 1, "initial watcher arm"); + +// Turn 1: main streams against the held provider stream; the watcher closes +// mid-turn and the extension delivers its wake while main is busy. +session.prompt("Reply with exactly the word OK.").catch((error) => { + events.push({ type: "prompt_error", text: error instanceof Error ? error.message : String(error) }); +}); +await waitFor(() => completions === 1 && session.isStreaming, "main streaming on the held completion"); +writeFileSync(process.env.FM_LIVE_WATCH_TRIGGER, "signal: live streaming probe\n"); +await waitFor( + () => events.some((event) => event.type === "input" && event.source === "extension" && event.streamingBehavior === "followUp" && event.text.includes("signal: live streaming probe")), + "the watcher follow-up queued while main streams", +); +await waitFor(() => armRows("arm ") === 2, "successor started while main streams"); +if (has("before_agent_start", "signal: live streaming probe")) { + throw new Error("Pi raised before_agent_start for a follow-up queued while streaming; the extension must never wait for that"); +} +releaseStream(); +await waitFor(() => settledRuns() === 1, "the first run to settle"); +if (!has("user_message_start", "signal: live streaming probe")) { + throw new Error(`the queued follow-up never joined the run as a user message: ${JSON.stringify(events)}`); +} +if (has("before_agent_start", "signal: live streaming probe")) { + throw new Error("Pi raised before_agent_start for a queued follow-up when the run reached it"); +} +if (completions !== 2) throw new Error(`the queued follow-up did not open its own model turn: ${completions} completions`); + +// Idle: the successor closes while main is idle, so the next wake starts a run +// and Pi raises before_agent_start with the exact text, then the user message. +writeFileSync(process.env.FM_LIVE_WATCH_TRIGGER, "signal: live idle probe\n"); +await waitFor(() => has("before_agent_start", "signal: live idle probe"), "the idle follow-up to raise before_agent_start"); +await waitFor(() => armRows("arm ") === 3, "successor after the idle delivery"); +await waitFor(() => settledRuns() === 2, "the idle run to settle"); +if (!has("user_message_start", "signal: live idle probe")) { + throw new Error(`the idle follow-up never reached the run as a user message: ${JSON.stringify(events)}`); +} +if (armRows("confirmed ") !== 2) { + throw new Error(`the watcher did not confirm both successor deliveries: ${readFileSync(armLog, "utf8")}`); +} +session.dispose(); +console.log("STREAM_OK"); +process.exit(0); +EOF +status=$? +out=$(cat "$TMP_ROOT/stream-output") +if [ "$status" -ne 0 ] || [ "$out" != "STREAM_OK" ]; then + fail "real-SDK streaming-time watcher delivery guard failed against pi-coding-agent $PI_VERSION: $out" +fi +pass "real Pi SDK $PI_VERSION queues a streaming-time watcher wake without before_agent_start, keeps the successor chain, and surfaces consumption of both follow-ups" diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index 6ce2de423b8..a7da27b2a76 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -1217,9 +1217,9 @@ const pi = { registerTool(candidate) { if (candidate.name === "fm_watch_arm_pi") tool = candidate; }, + // Nothing here consumes the follow-up: continuity must not depend on it. sendUserMessage: async (message) => { prompts.push(message); - queueMicrotask(() => handlers.get("before_agent_start")?.({ prompt: message }, {})); }, }; const rows = () => existsSync(process.env.FM_ARM_LOG) @@ -2014,6 +2014,119 @@ EOF pass "Pi replacement replays a streaming follow-up before consumption" } +# The 2026-09-02 incident: a wake delivered while main was mid-turn never raised +# before_agent_start, the extension waited for it, and every later actionable +# close was dropped. Continuity must settle on Pi accepting the follow-up, while +# consumption still decides what a replacement replays. +test_pi_streaming_time_delivery_keeps_the_successor_chain() { + local repo home plugin log trigger out status + repo="$TMP_ROOT/pi-streaming-chain-root" + home="$TMP_ROOT/pi-streaming-chain-home" + log="$TMP_ROOT/pi-streaming-chain.log" + trigger="$TMP_ROOT/pi-streaming-chain.trigger" + mkdir -p "$repo/bin" "$home/state" "$home/config" + install_pi_watch_extension_fixture "$repo" + plugin="$repo/.pi/extensions/fm-primary-pi-watch.ts" + cat > "$repo/bin/fm-watch-arm.sh" <<'SH' +#!/usr/bin/env bash +if [ "${1:-}" = --handling-delivered ]; then + printf 'confirmed=%s\n' "$2" >> "${FM_ARM_LOG:?}" + exit 0 +fi +printf 'arm=%s\n' "$$" >> "${FM_ARM_LOG:?}" +count=$(grep -c '^arm=' "$FM_ARM_LOG") +printf 'watcher: started pid=%s (beacon fresh) recovery-generation=chain-%s\n' "$$" "$count" +trap 'exit 0' TERM INT +while [ ! -e "$FM_TRIGGER_FILE.$count" ]; do sleep 0.02; done +printf 'signal: streaming chain wake %s\n' "$count" +exit 0 +SH + chmod +x "$repo/bin/fm-watch-arm.sh" + out=$(PLUGIN="$plugin" FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_ARM_LOG="$log" FM_TRIGGER_FILE="$trigger" node --input-type=module 2>&1 <<'EOF' +import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { pathToFileURL } from "node:url"; + +const handlers = new Map(); +const prompts = []; +let streaming = false; +let beforeAgentStarts = 0; +const pi = { + on(event, handler) { + handlers.set(event, handler); + }, + registerCommand() {}, + registerTool() {}, + // The real Pi prompt path: a follow-up sent while the agent is streaming is + // queued for the running run and raises no before_agent_start; only a send + // to an idle agent starts a run and raises it with the exact text. + sendUserMessage: async (message) => { + prompts.push(message); + if (streaming) return; + beforeAgentStarts += 1; + handlers.get("before_agent_start")?.({ prompt: message }, {}); + }, + events: { on() {}, emit() {} }, +}; +// The running run reaching a queued follow-up: Pi emits the user message. +const consumeQueued = (message) => + handlers.get("message_start")?.({ message: { role: "user", content: [{ type: "text", text: message }] } }, {}); +const arms = () => existsSync(process.env.FM_ARM_LOG) + ? readFileSync(process.env.FM_ARM_LOG, "utf8").split("\n").filter((row) => row.startsWith("arm=")).length + : 0; +async function waitFor(pred, label) { + for (let i = 0; i < 500; i += 1) { + if (pred()) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + throw new Error(`timeout waiting for ${label}`); +} +const wakes = (text) => prompts.filter((message) => message.includes(text)).length; + +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +const mod = await import(pathToFileURL(process.env.PLUGIN).href); +mod.default(pi); +await handlers.get("session_start")?.({ type: "session_start", reason: "startup" }, {}); +await waitFor(() => arms() === 1, "first arm"); +streaming = true; +writeFileSync(`${process.env.FM_TRIGGER_FILE}.1`, "close\n"); +await waitFor(() => prompts.length === 1, "first wake delivered while main streams"); +if (wakes("signal: streaming chain wake 1") !== 1) throw new Error(`wrong first wake: ${prompts.join(" | ")}`); +await waitFor(() => arms() === 2, "successor after the streaming-time delivery"); +writeFileSync(`${process.env.FM_TRIGGER_FILE}.2`, "close\n"); +await waitFor(() => prompts.length === 2, "second wake delivered while main still streams"); +if (wakes("signal: streaming chain wake 2") !== 1) throw new Error(`wrong second wake: ${prompts.join(" | ")}`); +await waitFor(() => arms() === 3, "successor after the second streaming-time delivery"); +if (beforeAgentStarts !== 0) throw new Error(`streaming follow-ups raised before_agent_start ${beforeAgentStarts} times`); + +// The run reaches the first queued follow-up; the second is still queued when +// the captain replaces the session, so only the second rides the handoff. +consumeQueued(prompts[0]); +await handlers.get("session_shutdown")?.({ type: "session_shutdown", reason: "new" }, {}); +const handoffPath = `${process.env.FM_HOME}/state/extensions/pi-primary-watch/session-replacement-actionable.json`; +const handoff = JSON.parse(readFileSync(handoffPath, "utf8")); +if (handoff.pending.length !== 1 || handoff.pending[0].delivered || !handoff.pending[0].message.includes("signal: streaming chain wake 2")) { + throw new Error(`replacement handoff did not carry exactly the unconsumed wake: ${JSON.stringify(handoff)}`); +} +streaming = false; +const replacementMod = await import(`${pathToFileURL(process.env.PLUGIN).href}?replacement=streaming-chain`); +replacementMod.default(pi); +await handlers.get("session_start")?.({ type: "session_start", reason: "new" }, {}); +await waitFor(() => prompts.length === 3, "replacement replay of the unconsumed wake"); +if (wakes("signal: streaming chain wake 2") !== 2 || wakes("signal: streaming chain wake 1") !== 1) { + throw new Error(`replacement replayed the wrong wakes: ${prompts.join(" | ")}`); +} +if (beforeAgentStarts !== 1) throw new Error(`idle replay raised before_agent_start ${beforeAgentStarts} times`); +await waitFor(() => arms() === 4, "replacement arm"); +await waitFor(() => !existsSync(handoffPath), "consumed replay clears its handoff record"); +process.exit(0); +EOF +) + status=$? + expect_code 0 "$status" "Pi streaming-time wake delivery must keep the successor chain and replay only unconsumed wakes" + [ -z "$out" ] || fail "Pi streaming-time delivery chain test printed output: $out" + pass "Pi streaming-time wake delivery keeps the successor chain and replays only unconsumed wakes" +} + test_pi_late_retiring_actionable_reaches_replacement() { local repo home plugin count out status repo="$TMP_ROOT/pi-late-retiring-actionable-root" @@ -3413,6 +3526,7 @@ test_pi_arm_distinguishes_session_lock_ownership test_pi_session_transition_generation_owner test_pi_session_replacement_carries_inflight_actionable_close test_pi_streaming_followup_is_replayed_after_replacement +test_pi_streaming_time_delivery_keeps_the_successor_chain test_pi_late_retiring_actionable_reaches_replacement test_pi_replacement_tokens_are_process_unique test_pi_replacement_persistence_failure_stops_arm_child From b955a9687cdc7d43acdc664bd274858636480795 Mon Sep 17 00:00:00 2001 From: kunchenguid Date: Wed, 2 Sep 2026 02:33:29 -0700 Subject: [PATCH 2/2] fix(pi): retry a verified successor that fails during wake delivery A verified successor can exit while the wake it was started for is still being delivered, most plausibly during a branch turn that holds the settlement for minutes. Its failure close arrived while the pipeline's single-flight guard was set, so the close handler skipped the retry, and the pipeline's end no longer launched an arm, which left the live generation with no watcher and no retry timer. The close handler now records that failure when the child had reported readiness and was not retired by the restoration itself, and the pipeline runs the ordinary bounded, lock-checked retry for it once the delivery settles. A restoration started for a later pending supersedes it, and an exhausted restoration still hands repair to main without a further arm. The regression holds a branch settlement open while the verified successor exits with a failure and proves one retry watcher starts after the settlement releases, none while it is held. Claude-Session: https://claude.ai/code/session_01QJjTsUvKkWAwLGNoncaZ3a --- .pi/extensions/fm-primary-pi-watch.ts | 44 ++++++++++-- docs/verification/runtime-backends.md | 2 + docs/verification/supervision.md | 2 +- tests/fm-pi-watch-extension.test.sh | 99 +++++++++++++++++++++++++++ 4 files changed, 140 insertions(+), 7 deletions(-) diff --git a/.pi/extensions/fm-primary-pi-watch.ts b/.pi/extensions/fm-primary-pi-watch.ts index aa2aba3a01a..31d08615f6d 100644 --- a/.pi/extensions/fm-primary-pi-watch.ts +++ b/.pi/extensions/fm-primary-pi-watch.ts @@ -99,6 +99,10 @@ type SessionGeneration = { // replacement began reads it to tell a main-queued wake (replayed) from a // branch-handled one (finished). unconsumedWakes: Map; + // 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( @@ -181,6 +185,9 @@ function replacementCoordinatorFor(handoff: string): ReplacementCoordinator { const replacementCoordinator = replacementCoordinatorFor(actionableHandoff); const armReadiness = new WeakMap>(); const armClose = new WeakMap>(); +// 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(); const armRecovery = new WeakMap(); const armPendingActionable = new WeakMap(); @@ -414,6 +421,7 @@ function createGeneration(): SessionGeneration { pendingActionables: [], cleanupFailure: "", unconsumedWakes: new Map(), + deferredClose: null, }; } @@ -736,6 +744,9 @@ 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"); @@ -782,11 +793,19 @@ export default function (pi: ExtensionAPI) { if (generationIsLive(owner)) { owner.restoring = false; if (owner.pendingActionables.some((pending) => pending.delivered)) schedulePendingCleanup(owner); - // No arm is launched here. A generation without a child at this point - // has delivered a typed restoration failure after its bounded retries, - // and that message 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". + // 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); + } } } } @@ -823,6 +842,7 @@ export default function (pi: ExtensionAPI) { async function retireArm(armChild: ChildProcess | null): Promise { if (!armChild) return true; + armRetired.add(armChild); armChild.kill("SIGTERM"); const closed = armClose.get(armChild); if (!closed) return false; @@ -933,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((resolveReady) => { @@ -946,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 => { @@ -989,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) => { diff --git a/docs/verification/runtime-backends.md b/docs/verification/runtime-backends.md index a4e63fd586e..c00af4256c2 100644 --- a/docs/verification/runtime-backends.md +++ b/docs/verification/runtime-backends.md @@ -1107,6 +1107,7 @@ FM_PI_BRANCH_LIVE_E2E=1 FM_PI_PACKAGE_DIR= bin/fm-test-run.sh ```text ok - Pi hung successor falls back to one typed actionable wake ok - Pi streaming-time wake delivery keeps the successor chain and replays only unconsumed wakes +ok - Pi retries a verified successor that failed during wake delivery once that delivery settles ok - tracked Pi extensions pass strict no-emit typecheck against Pi 0.84.4 ok - real Pi SDK 0.84.4 queues a streaming-time watcher wake without before_agent_start, keeps the successor chain, and surfaces consumption of both follow-ups ``` @@ -1114,3 +1115,4 @@ ok - real Pi SDK 0.84.4 queues a streaming-time watcher wake without before_agen The live probe loads the tracked watcher extension through Pi's real resource loader into a real AgentSession whose only provider is a local fake with its fetch intercepted in-process and held open mid-stream. It proved that a follow-up the extension sends while main is streaming raises no `before_agent_start` at queue time or when the run reaches it, joins the run as a user `message_start` carrying the exact wake text in its own model turn, and is followed by a verified successor and delivery of the next close; a follow-up sent to the idle main raises `before_agent_start` with the exact text before its user `message_start`. The portable regression drives the same shape with a fake main that never raises `before_agent_start` while streaming, then proves a replacement replays only the follow-up Pi had not consumed and that an exhausted restoration delivers its typed failure without launching a further arm. +A second regression holds a branch settlement open while the verified successor exits with a failure, and proves that failure takes the ordinary bounded retry once the delivery settles rather than leaving the generation with no watcher and no retry. diff --git a/docs/verification/supervision.md b/docs/verification/supervision.md index d45e123479b..41ed77fad30 100644 --- a/docs/verification/supervision.md +++ b/docs/verification/supervision.md @@ -474,7 +474,7 @@ The strict no-emit check used the installed Pi SDK declarations to hold the life Plain Pi and pi-signed share the same tracked `.pi/extensions/fm-primary-pi-watch.ts` path, so both inherit the generation owner; other primary harnesses are not applicable because they do not use this Pi extension lifecycle. On 2026-09-02 the same suite, the strict typecheck, and the credential-free real-SDK guard were rerun against `@earendil-works/pi-coding-agent` 0.84.4 after the extension stopped waiting for `before_agent_start` before settling a main delivery; [`runtime-backends.md`](runtime-backends.md#2026-09-02-streaming-time-watcher-delivery) owns the exact commands and output. -Observed guarantee: a wake delivered while main was streaming was followed by a verified successor and by delivery of the next actionable close, a replacement replayed only the follow-up Pi had not consumed, and an exhausted restoration delivered its typed failure without launching an arm past the retry bound. +Observed guarantee: a wake delivered while main was streaming was followed by a verified successor and by delivery of the next actionable close, a replacement replayed only the follow-up Pi had not consumed, an exhausted restoration delivered its typed failure without launching an arm past the retry bound, and a verified successor that failed while a branch settlement still held its wake took the ordinary bounded retry once that delivery settled. The once-per-generation recovery bound and immediate handling-successor poll were verified on 2026-08-21 with the tracked Pi extension, real watcher processes, and an isolated home. The regression forced handling confirmation to fail, observed one recovery follow-up across the former repeat window, confirmed the successor remained live, and then proved a separate handling successor durably queued a crew event within the bounded poll window. diff --git a/tests/fm-pi-watch-extension.test.sh b/tests/fm-pi-watch-extension.test.sh index a7da27b2a76..da38a3604a5 100755 --- a/tests/fm-pi-watch-extension.test.sh +++ b/tests/fm-pi-watch-extension.test.sh @@ -2127,6 +2127,104 @@ EOF pass "Pi streaming-time wake delivery keeps the successor chain and replays only unconsumed wakes" } +# A verified successor can die while the wake it was started for is still +# being delivered (a branch turn can take minutes). Its failure close arrives +# while the pipeline is busy, so the ordinary retry path must be deferred to +# the end of that delivery rather than skipped, or the live generation is left +# with no watcher and no retry. +test_pi_successor_failure_during_delivery_is_retried_after_delivery() { + local repo home plugin log stop out status + repo="$TMP_ROOT/pi-successor-dies-mid-delivery-root" + home="$TMP_ROOT/pi-successor-dies-mid-delivery-home" + log="$TMP_ROOT/pi-successor-dies-mid-delivery.log" + stop="$TMP_ROOT/pi-successor-dies-mid-delivery.stop" + mkdir -p "$repo/bin" "$home/state" "$home/config" + install_pi_watch_extension_fixture "$repo" + plugin="$repo/.pi/extensions/fm-primary-pi-watch.ts" + cat > "$repo/bin/fm-watch-arm.sh" <<'SH' +#!/usr/bin/env bash +printf 'arm=%s\n' "$$" >> "${FM_ARM_LOG:?}" +count=$(grep -c '^arm=' "$FM_ARM_LOG") +printf 'watcher: started pid=%s (beacon fresh)\n' "$$" +if [ "$count" -eq 1 ]; then + printf 'signal: wake before the successor dies\n' + exit 0 +fi +if [ "$count" -eq 2 ]; then + sleep 0.1 + printf 'watcher: FAILED - successor lost its beacon\n' + exit 3 +fi +trap 'exit 0' TERM INT +while [ ! -e "$FM_STOP_FILE" ]; do sleep 0.02; done +SH + chmod +x "$repo/bin/fm-watch-arm.sh" + out=$(PLUGIN="$plugin" FM_HOME="$home" FM_ROOT_OVERRIDE="$repo" FM_ARM_LOG="$log" FM_STOP_FILE="$stop" FM_WATCH_REARM_RETRY_BASE_MS=5 FM_WATCH_REARM_RETRY_MAX_MS=10 FM_WATCH_REARM_RETRY_LIMIT=2 node --input-type=module 2>&1 <<'EOF' +import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { pathToFileURL } from "node:url"; + +let releaseBranch = () => {}; +const branchSettlement = new Promise((resolve) => { + releaseBranch = resolve; +}); +let branchAccepted = false; +let tool = null; +const prompts = []; +const pi = { + on() {}, + registerCommand() {}, + registerTool(candidate) { + if (candidate.name === "fm_watch_arm_pi") tool = candidate; + }, + sendUserMessage: async (message) => { + prompts.push(message); + }, + events: { + on() {}, + emit(event, data) { + if (event !== "fm-branch-supervision:dispatch") return; + branchAccepted = true; + data.accept(branchSettlement); + }, + }, +}; +const arms = () => existsSync(process.env.FM_ARM_LOG) + ? readFileSync(process.env.FM_ARM_LOG, "utf8").split("\n").filter((row) => row.startsWith("arm=")).length + : 0; +async function waitFor(pred, label) { + for (let i = 0; i < 500; i += 1) { + if (pred()) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + throw new Error(`timeout waiting for ${label}`); +} + +writeFileSync(`${process.env.FM_HOME}/state/.lock`, `${process.pid}\n`); +writeFileSync(`${process.env.FM_HOME}/state/mid-delivery.meta`, "project=/projects/mid-delivery\nwindow=fm-mid-delivery\n"); +writeFileSync(`${process.env.FM_HOME}/state/.wake-queue`, "1\t1\tsignal\tmid-delivery.status\tsignal: wake before the successor dies\n"); +const mod = await import(pathToFileURL(process.env.PLUGIN).href); +mod.default(pi); +await tool.execute("initial-arm", {}, undefined, undefined, {}); +await waitFor(() => branchAccepted, "branch accepted the wake behind a verified successor"); +if (arms() !== 2) throw new Error(`expected the verified successor before delivery, got ${arms()} arms`); +// The successor dies while the branch still holds the delivery. +await new Promise((resolve) => setTimeout(resolve, 300)); +if (arms() !== 2) throw new Error(`a retry launched while the delivery was still in flight: ${arms()} arms`); +releaseBranch(); +await waitFor(() => arms() === 3, "a retry watcher after the delivery settled"); +await new Promise((resolve) => setTimeout(resolve, 150)); +if (arms() !== 3) throw new Error(`the deferred retry was not single-flight: ${arms()} arms`); +if (prompts.length !== 0) throw new Error(`a bounded retry surfaced a failure prompt: ${prompts.join(" | ")}`); +writeFileSync(process.env.FM_STOP_FILE, "stop\n"); +process.exit(0); +EOF +) + status=$? + expect_code 0 "$status" "Pi must retry a verified successor that failed during wake delivery" + [ -z "$out" ] || fail "Pi successor-dies-mid-delivery test printed output: $out" + pass "Pi retries a verified successor that failed during wake delivery once that delivery settles" +} + test_pi_late_retiring_actionable_reaches_replacement() { local repo home plugin count out status repo="$TMP_ROOT/pi-late-retiring-actionable-root" @@ -3527,6 +3625,7 @@ test_pi_session_transition_generation_owner test_pi_session_replacement_carries_inflight_actionable_close test_pi_streaming_followup_is_replayed_after_replacement test_pi_streaming_time_delivery_keeps_the_successor_chain +test_pi_successor_failure_during_delivery_is_retried_after_delivery test_pi_late_retiring_actionable_reaches_replacement test_pi_replacement_tokens_are_process_unique test_pi_replacement_persistence_failure_stops_arm_child