From 1f3e769616fdf9f31f85f4c3e6a9f71606634238 Mon Sep 17 00:00:00 2001 From: Ian Brown <742554+zestysoft@users.noreply.github.com> Date: Sat, 3 Oct 2026 03:19:54 -0700 Subject: [PATCH 1/6] fix(pi): hide duplicate assistant finals from hidden processing retries (#5863) * fix(pi): silence unacknowledged processing retry replies Suppress autonomous processing prose before persistence and during streaming while retaining tool calls, signed reasoning, and retryable outcomes. Restore ordinary output after acknowledgement or a user message. Fixes #4954 * no-mistakes(review): Silence only processing retries, keep first presentation visible * fix(pi): preserve differing processing retry replies --- .pi/extensions/fm-branch-supervision.ts | 62 +++++++- docs/pi-supervision-branch.md | 10 +- docs/verification/runtime-backends.md | 16 ++ tests/fm-pi-branch-extension.test.sh | 186 +++++++++++++++++++++--- tests/fm-pi-branch-live-e2e.test.sh | 109 ++++++++++++++ 5 files changed, 362 insertions(+), 21 deletions(-) diff --git a/.pi/extensions/fm-branch-supervision.ts b/.pi/extensions/fm-branch-supervision.ts index 25f2e31d924..dc44e5d7c38 100644 --- a/.pi/extensions/fm-branch-supervision.ts +++ b/.pi/extensions/fm-branch-supervision.ts @@ -653,10 +653,16 @@ export default function (pi: ExtensionAPI) { // queued for the captain's next prompt. The durable truth is the store's // processed marker; this only paces re-presentation and resets with the // session generation. - type ProcessingState = { sequences: string; through: number; triggered: number; pending: boolean; nextTurnQueued: boolean }; + type ProcessingState = { sequences: string; through: number; triggered: number; pending: boolean; nextTurnQueued: boolean; visibleFinals: Set }; let processing: ProcessingState | null = null; let queuedProcessingContent: string | null = null; let processingOpenedThisRun = false; + // Compare only replies to the same consumed sequence set. A retry can be + // the first real handling, so only an empty or exact-repeat final is hidden. + // Buffer retry streaming until message_end can make that decision; tool + // messages always keep their prose, and a user message ends this scope. + let activeProcessing: { request: ProcessingState; retry: boolean } | null = null; + let userMessageThisTurn = false; let processedInitializedGeneration = -1; // One revision for BOTH selections: a model or effort change invalidates an // in-flight branch build exactly the same way. @@ -1095,7 +1101,7 @@ export default function (pi: ExtensionAPI) { } if (processing?.pending) return true; if (!processing || processing.sequences !== sequences) { - processing = { sequences, through, triggered: 0, pending: false, nextTurnQueued: false }; + processing = { sequences, through, triggered: 0, pending: false, nextTurnQueued: false, visibleFinals: new Set() }; } // A presentation already sent is consumed by the run it joins or opens; // until that run settles, sending a widened or identical copy would hand @@ -1677,7 +1683,6 @@ ${context.command} // duplicate suppression. Operational extension injections are not dialog. const prompt = event.prompt; processingOpenedThisRun = queuedProcessingContent !== null && prompt === queuedProcessingContent; - if (processingOpenedThisRun) queuedProcessingContent = null; const trimmed = prompt.trim(); if (!trimmed || isOperationalUserText(trimmed)) return; const file = currentMainSession.getSessionFile() ?? ""; @@ -1692,6 +1697,48 @@ ${context.command} // so a fresh copy may be queued again once this run settles unacknowledged. if (processing) processing.nextTurnQueued = false; }); + pi.on?.("turn_start", () => { + userMessageThisTurn = false; + }); + pi.on?.("message_start", (event) => { + if (event.message.role === "user") { + userMessageThisTurn = true; + activeProcessing = null; + } else if ( + event.message.role === "custom" && + isProcessingCustomMessage(event.message) && + queuedProcessingContent !== null && + event.message.content === queuedProcessingContent + ) { + // message_start covers both an idle custom prompt and a follow-up + // consumed inside an existing run; neither needs before_agent_start. + activeProcessing = !userMessageThisTurn && processing ? { request: processing, retry: processing.triggered > 1 } : null; + queuedProcessingContent = null; + } + }); + pi.registerMarkdownTransformer?.((markdown, context) => + activeProcessing?.retry && context.isStreaming && context.messageType !== "user" ? "" : markdown, + ); + pi.on?.("message_end", (event) => { + if (!activeProcessing || event.message.role !== "assistant") return; + // message_end runs before tool execution. Keep the whole message when + // it carries a call, including prose alongside fm_branch_processed. + if (event.message.content.some((part) => part.type === "toolCall")) return; + const text = event.message.content.filter((part) => part.type === "text").map((part) => part.text).join("\n").trim(); + const { request, retry } = activeProcessing; + if (!retry || (text && !request.visibleFinals.has(text))) { + if (text) request.visibleFinals.add(text); + return; + } + // Pi applies the replacement before persistence and transcript rendering. + // Preserve the message envelope, including provider usage accounting. + return { + message: { + ...event.message, + content: [], + }, + }; + }); pi.on?.("context", (event, ctx) => { if (!afkPostureRecordPresent(state)) return; const messages = event.messages ?? []; @@ -1714,6 +1761,7 @@ ${context.command} mainStreaming = false; queuedProcessingContent = null; processingOpenedThisRun = false; + activeProcessing = null; if (processing) processing.pending = false; const settledGeneration = generation; await enqueueDelivery(async () => { @@ -1773,6 +1821,8 @@ ${context.command} consecutiveProviderErrors = 0; providerRecovery = null; generation += 1; + activeProcessing = null; + userMessageThisTurn = false; mirrorCollection.collectAnchor = null; mirrorCollection.pendingCursor = null; mirrorCollection.stagedCaptain = null; @@ -1821,6 +1871,9 @@ ${context.command} shuttingDown = true; generation += 1; processing = null; + queuedProcessingContent = null; + activeProcessing = null; + userMessageThisTurn = false; pendingMirror.length = 0; currentMainSession = null; mirrorCollection.collectAnchor = null; @@ -2339,6 +2392,9 @@ ${context.command} }; } const remaining = await readUnprocessedOutcomes(acknowledgedGeneration); + if (acknowledgedGeneration === generation && activeProcessing && through >= activeProcessing.request.through) { + activeProcessing = null; + } if (remaining !== null && remaining.length === 0) processing = null; const open = remaining === null ? "the remaining outcomes could not be read" diff --git a/docs/pi-supervision-branch.md b/docs/pi-supervision-branch.md index 717be533d37..0c61310e880 100644 --- a/docs/pi-supervision-branch.md +++ b/docs/pi-supervision-branch.md @@ -429,6 +429,14 @@ Nothing else advances that marker. An unrelated reply, an empty reply, or a reply that paraphrases the outcome leaves the sequence unprocessed. The extension presents the current unprocessed sequence set again at the next main run boundary and at every session start. +The first presentation of a sequence set is an ordinary turn whose response stays visible, including prose alongside `fm_branch_processed`. +From the second triggered presentation of that same set on, until its listed outcomes are acknowledged, an assistant final is removed before persistence only when its text is empty after trimming or exactly matches a final already visible for that sequence set after trimming. +Differing replies stay visible, including the first real handling after an empty or unrelated reply. +Messages carrying tool calls always retain their prose, signed reasoning, and usage accounting. +Pi's Markdown transformer API buffers retry prose while streaming on versions that expose it, so the complete reply can be compared before rendering. +A retained reply renders when the message ends. +Successful acknowledgement releases subsequent assistant output, and a real user message restores ordinary output immediately, including when a processing request rides that prompt. + ### Re-presentation pacing A presentation already pending its run boundary is not resent or widened. @@ -633,7 +641,7 @@ At that moment the branch reports any refusal instead of concluding there is "no - The new branch conversation at every main session start with continuation inside one session, and the mirror re-anchor that pairs with it. - Requested-versus-unsolicited delivery, exact visible entry content, and no unkeyed model turn. - The sequence-keyed processing request and its acknowledgement. -- Re-presentation after an empty reply and after an unrelated prior answer, the triggered-then-next-turn pacing, and session-start re-presentation. +- Suppression of empty or exact-repeat retry finals with differing replies preserved, preserved tool calls and user responses, the triggered-then-next-turn pacing, and session-start re-presentation. - Routine outcomes staying turn-free, task-level no-change notes staying hidden, absent-marker re-presentation, and malformed-age reporting without acknowledgement. - Idle and busy main state, and incident-shaped compaction and unrelated-assistant context. - Cold-start post-lock recovery, crash-before-cursor reload recovery, and repeated-reload idempotency. diff --git a/docs/verification/runtime-backends.md b/docs/verification/runtime-backends.md index 0b424269956..82187a3eb93 100644 --- a/docs/verification/runtime-backends.md +++ b/docs/verification/runtime-backends.md @@ -2148,6 +2148,22 @@ FM_HARNESS_LIVENESS_DRIFT=1 bin/fm-test-run.sh tests/fm-harness-liveness-drift-l The supervision-branch extension (`.pi/extensions/fm-branch-supervision.ts`, [docs/pi-supervision-branch.md](../pi-supervision-branch.md)) builds its second session through the Pi SDK surface: `createAgentSession` (including its `model`, `modelRuntime`, and `thinkingLevel` options), `DefaultResourceLoader` with `extensionFactories`, `SessionManager`, `createBashToolDefinition` with a `spawnHook`, `sendCustomMessage` for routine notes, `appendEntry` and `registerEntryRenderer` for captain outcomes, the `before_provider_request` hook, the command context's model registry for picker candidates, a fresh `ModelRuntime` for isolated-branch resolution, and Pi's own `getSupportedThinkingLevels`/`clampThinkingLevel` plus its `getThinkingLevel` and `thinking_level_select` extension surface for effort. In TUI mode, its `/supervision-model` model list is drawn with Pi's own `SelectList`, `Input`, `fuzzyFilter`, and `DynamicBorder` through the extension context's `ui.custom` surface, which is what bounds and searches a long catalog. +Processing-retry visibility was verified on 2026-09-27 against Pi 0.87.1 with a local intercepted provider stream, without credentials or an external provider request: + +```sh +bin/fm-test-run.sh tests/fm-pi-branch-extension.test.sh +FM_PI_BRANCH_LIVE_E2E=1 npm exec --yes --package=typescript@5.9.3 -- bin/fm-test-run.sh tests/fm-pi-branch-live-e2e.test.sh tests/fm-pi-primary-types.test.sh +``` + +```text +ok - real Pi SDK 0.87.1 suppresses only empty or exact-repeat retry finals, retains first and differing replies after reopen, buffers retry streaming, and keeps outcomes retryable +ok - tracked Pi extensions pass strict no-emit typecheck against Pi 0.87.1 +``` + +The guard runs the extension through Pi's actual message event runner, renders its streamed replies with the stock assistant component, and checks both live agent state and a reopened session file. +The portable processing-turn case additionally covers whitespace-only replies, a one-character difference, prose alongside acknowledgment calls, signed reasoning and tool-call preservation, rejected and partial acknowledgements, busy follow-ups, user steering, and both orderings of a user message batched with a processing request. +Other primary harnesses do not load this Pi extension, and these event and persistence boundaries are independent of the runtime session backend. + Evidence produced 2026-08-25 on macOS 26.5.2 arm64, Node v24.13.1: - Historical real-SDK guard: `FM_PI_BRANCH_LIVE_E2E=1 bin/fm-test-run.sh tests/fm-pi-branch-live-e2e.test.sh` against the globally installed `@earendil-works/pi-coding-agent` 0.81.1 printed `ok - real Pi SDK 0.81.1 accepts the branch session construction and preserves an unpromptable wake`. diff --git a/tests/fm-pi-branch-extension.test.sh b/tests/fm-pi-branch-extension.test.sh index fd651cefa01..bdadcffa390 100644 --- a/tests/fm-pi-branch-extension.test.sh +++ b/tests/fm-pi-branch-extension.test.sh @@ -557,6 +557,7 @@ const mainUserMessages = []; const mainTools = []; const renderers = new Map(); const entryRenderers = new Map(); +const markdownTransformers = []; const mainEntries = []; const mainSessionManager = { getSessionFile: () => `${home}/main.jsonl`, @@ -581,6 +582,9 @@ const pi = { registerEntryRenderer(customType, renderer) { entryRenderers.set(customType, renderer); }, + registerMarkdownTransformer(transformer) { + markdownTransformers.push(transformer); + }, appendEntry(customType, data) { activeMainSession.getEntries().push({ type: "custom", customType, data }); }, @@ -1263,14 +1267,47 @@ test_captain_outcome_processing_turn_is_sequence_keyed_and_re_presented() { PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \ DRIVER_PRELUDE="$DRIVER_PRELUDE" node --input-type=module > "$TMP_ROOT/node-output" 2>&1 <<'EOF' const prelude = process.env.DRIVER_PRELUDE; -await eval(`(async () => { ${prelude}; globalThis.__t = { fire, dispatch, settle, sentToMain, mainEntries, mainTools, outcomeScript, defaultSessionCtx, home, bus }; })()`); -const { fire, dispatch, settle, sentToMain, mainEntries, mainTools, outcomeScript, defaultSessionCtx, home, bus } = globalThis.__t; +await eval(`(async () => { ${prelude}; globalThis.__t = { fire, dispatch, settle, sentToMain, mainEntries, mainTools, outcomeScript, defaultSessionCtx, home, bus, markdownTransformers }; })()`); +const { fire, dispatch, settle, sentToMain, mainEntries, mainTools, outcomeScript, defaultSessionCtx, home, bus, markdownTransformers } = globalThis.__t; import { writeFileSync } from "node:fs"; let requestsFloor = 0; const requests = () => sentToMain.filter((sent) => sent.message.customType === "fm-branch-process").slice(requestsFloor); const unprocessedSeqs = () => outcomeScript(["unprocessed"]).split("\n").filter(Boolean).map((line) => JSON.parse(line).seq); -const runOf = async (fn) => { await fire("agent_start", {}); await fn?.(); await fire("agent_end", {}); await fire("agent_settled", {}); }; +let consumedRequests = 0; +const consumeRequest = async (userFirst = false) => { + const pending = requests().at(-1); + if (!pending || consumedRequests === requests().length) return; + consumedRequests = requests().length; + if (userFirst || pending.options.deliverAs === "nextTurn") { + await fire("message_start", { message: { role: "user", content: "A new question" } }); + } + await fire("message_start", { message: { role: "custom", ...pending.message } }); +}; +const runOf = async (fn) => { + await fire("agent_start", {}); + await fire("turn_start", {}); + await consumeRequest(); + await fn?.(); + await fire("agent_end", {}); + await fire("agent_settled", {}); +}; +const render = (text, isStreaming = true, messageType = "assistant") => markdownTransformers.reduce( + (value, transform) => transform(value, { messageType, isStreaming, availableWidth: 80 }), text, +); +const finish = async (text, extra = []) => { + const message = { role: "assistant", content: [...(text ? [{ type: "text", text }] : []), ...extra], usage: { totalTokens: 7 }, stopReason: "stop" }; + await fire("message_start", { message }); + const replacement = await fire("message_end", { message }); + const stored = replacement?.message ?? message; + mainEntries.push({ type: "message", message: stored }); + if (stored.usage !== message.usage) throw new Error("suppression lost usage accounting"); + return stored; +}; +const visibleFinals = () => mainEntries.filter((entry) => entry.type === "message" && entry.message.role === "assistant") + .flatMap((entry) => entry.message.content.filter((part) => part.type === "text").map((part) => part.text)); +const priorResult = "Completed the requested work. The checks passed and the result is ready for review."; +await finish(priorResult); // A home with a delivered captain row and no processed marker (upgraded from // before the marker existed, or switched from the supervision host, whose @@ -1320,20 +1357,41 @@ if (request.options.triggerTurn !== true || request.options.deliverAs !== "follo if (!request.message.content.includes(`[seq ${seq}, recorded 0m ago] task-d: ${decision}`)) throw new Error(`the request lost its key or summary: ${request.message.content}`); if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error(`delivery did not leave seq ${seq} unprocessed: ${unprocessedSeqs()}`); -// Case A (timeline report 2026-08-31): the turn returns an EMPTY assistant -// message. The processed marker must not move, and the same sequence is -// presented again at the run boundary. -await runOf(() => mainEntries.push({ type: "message", message: { role: "assistant", content: [] } })); -if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error("an empty answer advanced the processed marker"); -if (requests().length !== 2) throw new Error(`an empty answer did not re-present the outcome: ${requests().length} requests`); -if (requests()[1].options.triggerTurn !== true) throw new Error("the first re-presentation must open its own turn"); +// The first presentation carries the one visible response for this outcome, +// even when it forgets to acknowledge. +await runOf(async () => { + if (render("Captain, task-d needs your call.") !== "Captain, task-d needs your call.") { + throw new Error("the first processing presentation hid its response while streaming"); + } + await finish("Captain, task-d needs your call."); +}); +if (visibleFinals().at(-1) !== "Captain, task-d needs your call.") throw new Error("the first processing presentation lost its final"); +if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error("an unacknowledged answer advanced the processed marker"); +if (requests().length !== 2 || requests()[1].options.triggerTurn !== true) throw new Error("the first re-presentation must open its own turn"); if (!requests()[1].message.content.includes(`[seq ${seq}, recorded 0m ago] task-d: ${decision}`)) throw new Error("the re-presentation changed the outcome"); - -// Case B: the turn repeats an unrelated prior answer. Same result: the marker -// holds, and the request is presented again - now riding the captain's next -// prompt because the triggered budget for this sequence set is spent. -await runOf(() => mainEntries.push({ type: "message", message: { role: "assistant", content: "The retry safe-stopped; diagnosis is underway." } })); -if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error("an unrelated answer advanced the processed marker"); +// The hidden retry repeats this set's prior reply instead of acknowledging. +// Exercise Pi's public message replacement and Markdown transformer surfaces: +// buffer streaming until the complete reply can be compared, then keep every +// differing reply, even one already visible outside this processing set. +await runOf(async () => { + if (render("Captain, task-d needs your call.") !== "" || render("prior reasoning", true, "assistant-thinking") !== "") { + throw new Error("a processing retry reply leaked while streaming"); + } + if (render(priorResult, false) !== priorResult || render("A question", true, "user") !== "A question") { + throw new Error("silencing a retry hid an earlier final or a user message"); + } + const repeated = await finish(" \nCaptain, task-d needs your call. \n"); + if (repeated.content.length) throw new Error("the exact trimmed repeat retained visible content"); + await finish(priorResult); + await finish("Captain, task-d needs your call!"); + const empty = await finish(""); + const whitespace = await finish(" \n\t"); + if (empty.content.length || whitespace.content.length) throw new Error("an empty retry retained visible content"); +}); +if (JSON.stringify(visibleFinals()) !== JSON.stringify([priorResult, "Captain, task-d needs your call.", priorResult, "Captain, task-d needs your call!"])) { + throw new Error(`retry comparison hid new prose or exposed a repeat: ${JSON.stringify(visibleFinals())}`); +} +if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error("a repeated or empty answer advanced the processed marker"); if (requests().length !== 3) throw new Error(`an unrelated answer did not re-present the outcome: ${requests().length} requests`); if (requests()[2].options.deliverAs !== "nextTurn" || requests()[2].options.triggerTurn) { throw new Error(`after the triggered budget the request must ride the next prompt: ${JSON.stringify(requests()[2].options)}`); @@ -1342,7 +1400,11 @@ if (requests()[2].options.deliverAs !== "nextTurn" || requests()[2].options.trig await fire("agent_settled", {}); if (requests().length !== 3) throw new Error("a duplicate next-turn copy was queued"); // The captain's next prompt consumes that copy; settling unacknowledged queues one more. -await runOf(() => mainEntries.push({ type: "message", message: { role: "assistant", content: "Captain, shipshape." } })); +await runOf(async () => { + if (render("The new answer") !== "The new answer") throw new Error("a nextTurn request hid the user response"); + await finish("The new answer"); +}); +if (visibleFinals().at(-1) !== "The new answer") throw new Error("a nextTurn request removed the user final"); if (requests().length !== 4 || requests()[3].options.deliverAs !== "nextTurn") throw new Error("the outcome stopped being re-presented on later prompts"); if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error("a paraphrase advanced the processed marker"); @@ -1354,6 +1416,32 @@ if (mainEntries.filter((entry) => entry.customType === "fm-branch-visible-outcom throw new Error("re-presentation duplicated the visible entry"); } +const call = { type: "toolCall", id: "ack-call", name: "fm_branch_processed", arguments: { through: seq } }; +const thinking = { type: "thinking", thinking: "reasoning for the tool", thinkingSignature: "provider-signature" }; +// The replacement's first presentation keeps prose sent alongside the +// acknowledgement call; this run's call is left unexecuted. +await runOf(async () => { + const firstWithTool = await finish("Captain, task-d still needs your call.", [thinking, call]); + if (firstWithTool.content.length !== 3 || firstWithTool.content[0].text !== "Captain, task-d still needs your call.") { + throw new Error("the first presentation dropped prose sent alongside its acknowledgement"); + } +}); +if (visibleFinals().at(-1) !== "Captain, task-d still needs your call.") throw new Error("the first presentation after replacement lost its final"); +if (requests().length !== 6 || requests()[5].options.triggerTurn !== true) throw new Error("the replacement did not retry the unacknowledged outcome"); +await fire("agent_start", {}); +await fire("turn_start", {}); +await consumeRequest(); +await finish("stale after reload"); +if (visibleFinals().at(-1) !== "stale after reload") throw new Error("session replacement hid differing retry prose"); +const repeatedAfterReload = await finish("stale after reload"); +if (repeatedAfterReload.content.length) throw new Error("session replacement lost retry comparison"); +const withTool = await finish("stale after reload", [thinking, call]); +if (withTool.content.length !== 3 || withTool.content[0].text !== "stale after reload" || withTool.content[1] !== thinking || withTool.content[2] !== call) { + throw new Error("retry comparison dropped prose, signed reasoning, or the acknowledgement call"); +} +await fire("turn_start", {}); +if (render("still unacknowledged") !== "") throw new Error("a tool continuation released suppression before acknowledgement"); + // Only the sequence-bound acknowledgement closes it. const nativeTools = new Map(); const messageTypes = new Set(); @@ -1376,9 +1464,13 @@ if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error const tooFar = await processed.execute("ack-too-far", { through: seq + 100 }, undefined, undefined, {}); if (!tooFar.isError) throw new Error("an acknowledgement beyond the read cursor was accepted"); if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seq])) throw new Error("a refused acknowledgement moved the marker"); +if (render("still refused") !== "") throw new Error("a refused acknowledgement released suppression"); const ack = await processed.execute("ack", { through: seq }, undefined, undefined, {}); if (ack.isError) throw new Error(`acknowledgement failed: ${JSON.stringify(ack)}`); if (unprocessedSeqs().length !== 0) throw new Error("the acknowledgement did not close the sequence"); +if (render("A newly processed response") !== "A newly processed response") throw new Error("successful acknowledgement did not release the response"); +await finish("A newly processed response"); +if (visibleFinals().at(-1) !== "A newly processed response") throw new Error("the acknowledged outcome lost its response"); const before = requests().length; await runOf(); if (requests().length !== before) throw new Error("an acknowledged outcome was presented again"); @@ -1431,9 +1523,13 @@ await runOf(); if (requests().length !== beforePairRepeat + 1 || requests().at(-1).options.triggerTurn !== true) { throw new Error("the second presentation of the widened sequence set did not open its own turn"); } +await fire("agent_start", {}); +await fire("turn_start", {}); +await consumeRequest(); const partial = await processed.execute("ack-partial", { through: seqE }, undefined, undefined, {}); if (partial.isError) throw new Error(`partial acknowledgement failed: ${JSON.stringify(partial)}`); if (JSON.stringify(unprocessedSeqs()) !== JSON.stringify([seqF])) throw new Error(`a partial acknowledgement did not keep the newer sequence open: ${unprocessedSeqs()}`); +if (render("partial response") !== "") throw new Error("partial acknowledgement released the remaining outcome's retry"); const beforeF = requests().length; await runOf(); if ( @@ -1447,6 +1543,62 @@ if ( const done = await processed.execute("ack-final", { through: seqF }, undefined, undefined, {}); if (done.isError || unprocessedSeqs().length !== 0) throw new Error("the final acknowledgement did not close the newer sequence"); +// A follow-up queued while main is busy must not hide the answer already +// underway. Suppression starts only when Pi consumes the custom message. +await fire("agent_start", {}); +await fire("turn_start", {}); +await fire("message_start", { message: { role: "user", content: "An active user request" } }); +await report2.execute("busy-report", { task: "task-g", verdict: "captain", summary: "A new decision" }, undefined, undefined, {}); +await finish("The busy user answer"); +if (visibleFinals().at(-1) !== "The busy user answer") throw new Error("queueing a retry hid an in-flight user answer"); +await fire("turn_start", {}); +await consumeRequest(); +await finish("Captain, task-g needs your call."); +if (visibleFinals().at(-1) !== "Captain, task-g needs your call.") throw new Error("a consumed busy follow-up hid its first presentation"); +// User steering in that same turn must immediately recover ordinary output. +await fire("message_start", { message: { role: "user", content: "A steering question" } }); +await finish("The steering answer"); +if (visibleFinals().at(-1) !== "The steering answer") throw new Error("processing suppression hid a steering response"); +await fire("agent_end", {}); +await fire("agent_settled", {}); +// Pi can also batch the user before the custom follow-up in one turn. +await fire("agent_start", {}); +await fire("turn_start", {}); +await consumeRequest(true); +await finish("The batched user answer"); +if (visibleFinals().at(-1) !== "The batched user answer") throw new Error("processing suppression hid a user batched before the custom message"); +const latestSeq = unprocessedSeqs().at(-1); +await processed.execute("ack-busy", { through: latestSeq }, undefined, undefined, {}); +await fire("agent_end", {}); +await fire("agent_settled", {}); + +// A first reply can be empty or unrelated: a retry that finally handles the +// outcome must stay visible. Reusing the same response across new sequence +// sets must not make it a duplicate, and only the tool closes each outcome. +for (const initial of ["", "Unrelated prior acknowledgment"]) { + await report2.execute("new-set", { task: "task-new", verdict: "captain", summary: "Another decision" }, undefined, undefined, {}); + const newSeq = unprocessedSeqs().at(-1); + const startCount = visibleFinals().length; + await runOf(async () => { + const firstReply = await finish(initial); + if ((firstReply.content[0]?.text ?? "") !== initial) throw new Error("first presentation changed its output"); + }); + const newAnswer = "Handled the new outcome."; + await runOf(async () => { + await finish(newAnswer); + if (visibleFinals().at(-1) !== newAnswer) throw new Error("the first real handling on a retry was hidden"); + const repeat = await finish(newAnswer); + if (repeat.content.length) throw new Error("a repeated retry final was retained"); + }); + if (visibleFinals().length !== startCount + (initial ? 2 : 1)) throw new Error("sequence comparison lost or duplicated a final"); + if (!unprocessedSeqs().includes(newSeq) || requests().at(-1).options.deliverAs !== "nextTurn") { + throw new Error("an empty, unrelated, or differing reply closed the durable obligation"); + } + await processed.execute("ack-new-set", { through: newSeq }, undefined, undefined, {}); + if (unprocessedSeqs().length) throw new Error("acknowledgement did not close the new set"); + await fire("agent_settled", {}); +} + // A session that does not own the fleet lock cannot acknowledge anything. writeFileSync(`${home}/state/.lock`, "1\n"); const foreign = await processed.execute("ack-foreign", { through: seqF }, undefined, undefined, {}); diff --git a/tests/fm-pi-branch-live-e2e.test.sh b/tests/fm-pi-branch-live-e2e.test.sh index 76a49a4995d..376c8b0ee52 100644 --- a/tests/fm-pi-branch-live-e2e.test.sh +++ b/tests/fm-pi-branch-live-e2e.test.sh @@ -991,3 +991,112 @@ 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" + +# The first processing presentation keeps its visible response, while the +# hidden retry drops only empty or exact-repeat finals through the real event +# runner, stock assistant renderer, and persistence, including after reopen. +for retry_case in repeated differing empty first-empty; do +retryhome="$TMP_ROOT/retry-home-$retry_case" +retrydir="$TMP_ROOT/retry-agent-dir-$retry_case" +mkdir -p "$retryhome/state" "$retryhome/config" "$retrydir" +cp "$streamdir/models.json" "$retrydir/models.json" +BRANCH_PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" \ + FM_HOME="$retryhome" FM_ROOT_OVERRIDE="$ROOT" RETRY_CASE="$retry_case" \ + PI_CODING_AGENT_DIR="$retrydir" PI_PACKAGE_DIR="$PI_PACKAGE_DIR" \ + node --input-type=module > "$TMP_ROOT/retry-output" 2>&1 <<'EOF' +import { readFileSync, writeFileSync } from "node:fs"; +import { spawnSync } from "node:child_process"; +import { resolve } from "node:path"; +import { pathToFileURL } from "node:url"; +const home = process.env.FM_HOME; +const pkg = resolve(process.env.PI_PACKAGE_DIR); +const { DefaultResourceLoader, ModelRegistry, ModelRuntime, SessionManager, SettingsManager, createAgentSession, initTheme } = + await import(pathToFileURL(`${pkg}/dist/index.js`).href); +const { AssistantMessageComponent } = await import(pathToFileURL(`${pkg}/dist/modes/interactive/components/assistant-message.js`).href); +initTheme("dark"); +writeFileSync(`${home}/state/.lock`, `${process.pid}\n`); +const outcome = (...args) => { + const result = spawnSync("bash", [`${process.env.FM_ROOT_OVERRIDE}/bin/fm-branch-outcome.sh`, ...args], { encoding: "utf8" }); + if (result.status !== 0) throw new Error(result.stderr); + return result.stdout.trim(); +}; +outcome("processed-init"); +const seq = Number(outcome("append", "--task", "example", "--verdict", "captain", "--summary", "A decision is needed")); +const original = "The requested result is complete and verified."; +const handled = process.env.RETRY_CASE === "first-empty" ? "" : "The example task needs a decision."; +const retryReply = process.env.RETRY_CASE === "repeated" ? handled : process.env.RETRY_CASE === "empty" ? "" : original; +const expectedFinals = [original, ...(handled ? [handled] : []), ...(retryReply && retryReply !== handled ? [retryReply] : [])]; +let completions = 0; +let settled = 0; +const failures = []; +const chunk = (text, finish = null) => `data: ${JSON.stringify({ + id: "local-retry-probe", object: "chat.completion.chunk", created: 1, model: "fm-live-stream-model", + choices: [{ index: 0, delta: text ? { role: "assistant", content: text } : {}, finish_reason: finish }], +})}\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: ${url}`); + completions += 1; + const text = completions === 1 ? original : completions === 2 ? handled : completions === 3 ? retryReply : "The new user answer."; + return new Response(chunk(text) + chunk(null, "stop") + "data: [DONE]\n\n", { + headers: { "content-type": "text/event-stream" }, + }); +}; +const agentDir = process.env.PI_CODING_AGENT_DIR; +const settings = SettingsManager.create(home, agentDir); +const loader = new DefaultResourceLoader({ + cwd: home, agentDir, settingsManager: settings, + additionalExtensionPaths: [process.env.BRANCH_PLUGIN], + extensionFactories: [{ name: "retry-probe", factory: (pi) => { + pi.on("agent_settled", () => { settled += 1; }); + } }], + 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 manager = SessionManager.create(home, `${home}/sessions`); +const { session } = await createAgentSession({ + cwd: home, sessionManager: manager, settingsManager: settings, resourceLoader: loader, + modelRuntime: runtime, model: registry.find("fm-live-stream", "fm-live-stream-model"), noTools: "builtin", +}); +let streamedRetries = 0; +const unsubscribe = session.subscribe((event) => { + if (event.type !== "message_update" || completions !== 3) return; + streamedRetries += 1; + const component = new AssistantMessageComponent(undefined, false, undefined, undefined, 0, session.extensionRunner.getMarkdownTransformers()); + component.updateContent(event.message, true); + const rendered = component.render(120).join("\n"); + if (retryReply && rendered.includes(retryReply)) failures.push("retry prose leaked from the streaming renderer"); +}); +await session.prompt("Finish the requested work."); +for (let i = 0; i < 600 && settled < 3; i += 1) await new Promise((done) => setTimeout(done, 50)); +if (settled !== 3 || completions !== 3) throw new Error(`retry chain did not settle: ${settled} settlements, ${completions} completions`); +if ((retryReply && streamedRetries === 0) || failures.length) throw new Error(`streaming suppression failed: ${streamedRetries} updates, ${failures}`); +const assistantText = (messages) => messages.filter((message) => message.role === "assistant") + .flatMap((message) => message.content.filter((part) => part.type === "text").map((part) => part.text)); +if (JSON.stringify(assistantText(session.messages)) !== JSON.stringify(expectedFinals)) throw new Error(`agent state lost new prose or retained an exact repeat: ${JSON.stringify(assistantText(session.messages))}`); +const reopened = SessionManager.open(manager.getSessionFile(), `${home}/sessions`); +if (JSON.stringify(assistantText(reopened.buildSessionContext().messages)) !== JSON.stringify(expectedFinals)) throw new Error("reopened session lost new prose or retained an exact repeat"); +if (!outcome("unprocessed").includes(`"seq":${seq}`)) throw new Error("silent retries advanced the processed marker"); +await session.prompt("A new user question."); +if (assistantText(session.messages).at(-1) !== "The new user answer.") throw new Error("a nextTurn retry hid the new user answer"); +const acknowledged = await session.getToolDefinition("fm_branch_processed").execute("ack", { through: seq }, undefined, undefined, {}); +if (acknowledged.isError || outcome("unprocessed")) throw new Error("the outcome could not be acknowledged after silent retries"); +const entries = readFileSync(manager.getSessionFile(), "utf8").split("\n").filter(Boolean).map((line) => JSON.parse(line)); +if (assistantText(entries.filter((entry) => entry.type === "message").map((entry) => entry.message)).length !== expectedFinals.length + 1) { + throw new Error("persisted finals do not match the retained handling outcomes"); +} +unsubscribe(); +session.dispose(); +console.log("RETRY_OK"); +process.exit(0); +EOF +status=$? +out=$(cat "$TMP_ROOT/retry-output") +if [ "$status" -ne 0 ] || [ "$out" != "RETRY_OK" ]; then + fail "real-SDK processing retry visibility guard ($retry_case) failed against pi-coding-agent $PI_VERSION: $out" +fi +done +pass "real Pi SDK $PI_VERSION suppresses only empty or exact-repeat retry finals, retains first and differing replies after reopen, buffers retry streaming, and keeps outcomes retryable" From fede6197556c188dd2f451809d17851760d1ef46 Mon Sep 17 00:00:00 2001 From: Kun Chen <3233006+kunchenguid@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:06:52 -0700 Subject: [PATCH 2/6] test(pi): accept Pi 1.0.1's renamed HTML export renderer lookup (#6530) Pi 1.0.1's createToolHtmlRenderer reads getToolRenderers and ignores getToolDefinition. Calm /export still includes stock grep HTML; the fixture has to pass the lookup key the installed Pi actually reads. --- tests/fm-calm-pi-extension.test.sh | 25 +++++++++++++++---------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/tests/fm-calm-pi-extension.test.sh b/tests/fm-calm-pi-extension.test.sh index 287c5de2b0e..f4baded12f4 100755 --- a/tests/fm-calm-pi-extension.test.sh +++ b/tests/fm-calm-pi-extension.test.sh @@ -1659,14 +1659,23 @@ for (const { name, actual } of rows) { throw new Error(`${name} was not hidden before export rendering`); } } -async function assertStockHtmlRendering(command, submitData) { - editorText = command; - terminalInputHandler(submitData); - const htmlRenderer = createToolHtmlRenderer({ - getToolDefinition: (name) => tools.find((tool) => tool.name === name), +// Pi 1.0.0 resolves export HTML through getToolDefinition. Pi 1.0.1 renamed that +// dependency to getToolRenderers and ignores the old key, so a fixture that +// passes only the old key reports every tool as missing. Supply both; each +// release reads the key it knows and renders the same wrapped definitions. +function createInstalledToolHtmlRenderer() { + const lookup = (name) => tools.find((tool) => tool.name === name); + return createToolHtmlRenderer({ + getToolDefinition: lookup, + getToolRenderers: lookup, theme, cwd: process.cwd(), }); +} +async function assertStockHtmlRendering(command, submitData) { + editorText = command; + terminalInputHandler(submitData); + const htmlRenderer = createInstalledToolHtmlRenderer(); const exportCases = [ ...cases.filter(([toolName]) => toolName === "grep" || toolName === "find"), ["fm_watch_arm_pi", watchArgs, watchResult], @@ -1693,11 +1702,7 @@ await assertStockHtmlRendering("/export calm.html", "\r"); getKeybindings().setUserBindings({ "tui.input.submit": "alt+s" }); editorText = "/export remapped.html"; terminalInputHandler("\r"); -const unmatchedRenderer = createToolHtmlRenderer({ - getToolDefinition: (name) => tools.find((tool) => tool.name === name), - theme, - cwd: process.cwd(), -}); +const unmatchedRenderer = createInstalledToolHtmlRenderer(); if (unmatchedRenderer.renderCall("unmatched-submit", "grep", { pattern: "alpha", path: "." })) { throw new Error("ordinary non-submit input activated HTML export rendering"); } From 2e659ffdc0b008c606152e76a594804306684653 Mon Sep 17 00:00:00 2001 From: Joseph Kim Date: Sun, 4 Oct 2026 03:20:26 -0700 Subject: [PATCH 3/6] fix(bin): tolerate transient quota read failures (#6490) * fix(procevent-quota): tolerate consecutive slow quota-axi reads The quota allowance poll treated any quota_json failure as terminal, so one slow quota-axi --json (measured max ~29s under a 48s derived bound) shut the watch down until someone re-armed it, and the detail always said "missing/incompatible". Tolerate three consecutive failed or timed-out reads before going terminal, reset the streak on any good read, and report a timeout distinctly from a missing or incompatible tool. Each timed poll runs exactly one bounded --version and one bounded --json: validate the captured version text through fm_quota_axi_version_compatible rather than launching a second probe, and describe a mixed failure streak by count plus last cause. * no-mistakes(document): Document quota polling failure tolerance * no-mistakes(ci): Fixed ci-2 and ci-3. Permanent quota read failures (rc 2 missing, rc 3 incompatible) now report on the first poll, while transient rc 1/4 failures retain the existing three-failure retry behavior. The missing-binary test now uses an isolated PATH without quota-axi and asserts both permanent failures stop at condition_polls: 1. Verification passed: tests/fm-procevent-quota.test.sh, canonical fast lint for both changed files, bash syntax checks, and git diff --check * no-mistakes(review): Classify untimed quota version failures as transient * no-mistakes(document): Clarify quota polling failure budget --- bin/fm-procevent-quota.sh | 69 ++++++++-- bin/fm-quota-axi-lib.sh | 33 +++-- tests/fm-procevent-quota.test.sh | 216 ++++++++++++++++++++++++++++++- 3 files changed, 293 insertions(+), 25 deletions(-) diff --git a/bin/fm-procevent-quota.sh b/bin/fm-procevent-quota.sh index 16ce34da2c4..5c41aa1104e 100755 --- a/bin/fm-procevent-quota.sh +++ b/bin/fm-procevent-quota.sh @@ -17,7 +17,10 @@ # registered through `bin/fm-procevent.sh register`. # poll The blocking child the generic runner executes; never run this # directly in a conversational turn. It polls `quota-axi --json` -# until quota drops below the threshold or an error stops the watch. +# until quota drops below the threshold, invalid quota data stops +# the watch, or three consecutive transient command failures stop +# it. Missing or incompatible tools stop it immediately, and a +# successful read resets the command-failure streak. # classify Print the captured outcome class: low, exhausted, error, or unknown. # terminal Every quota poll is terminal because the source fires at most once. # source-id Print the canonical source id. @@ -52,6 +55,9 @@ STATE="${FM_STATE_OVERRIDE:-$FM_HOME/state}" DEFAULT_INTERVAL=60 DEFAULT_THRESHOLD=10 +# Consecutive transient quota-axi read failures before poll goes terminal. +# Missing and incompatible tools bypass this budget. No config knob on purpose. +MAX_CONSECUTIVE_READ_FAILURES=3 SOURCE_ID_BASE=quota @@ -98,16 +104,43 @@ valid_percent() { } # quota_json [timeout] -# Run `quota-axi --json` bounded by the given timeout. A missing or incompatible -# quota-axi is an error condition, not a signal to fire. +# Run `quota-axi --json` bounded by the given timeout. +# Exit status: 0 prints JSON; 1 timed out; 2 missing; 3 incompatible; 4 other failure. +# A missing or incompatible quota-axi is an error condition, not a signal to fire. +# Callers tolerate a bounded streak of 1/4 before going terminal; 2/3 stay distinct. +# Each path probes --version once, validates that captured text through +# fm_quota_axi_version_compatible, then probes --json once. +# A slow or failing probe stays 1/4; "incompatible" is reserved for an actual +# unsupported or unparseable version string. quota_json() { - local timeout=${1:-} output + local timeout=${1:-} output rc=0 + if ! command -v quota-axi >/dev/null 2>&1; then + return 2 + fi if [ -n "$timeout" ]; then - fm_quota_axi_compatible "$timeout" >/dev/null 2>&1 || return 2 - output=$(fm_run_timed "$timeout" quota-axi --json 2>/dev/null /dev/null /dev/null /dev/null 2>&1 || return 2 - output=$(quota-axi --json 2>/dev/null /dev/null /dev/null +# True when the printed `quota-axi --version` text meets FM_QUOTA_AXI_MIN. +# Callers that already captured a bounded --version pass that text here so they +# do not launch a second, unbounded version probe. +fm_quota_axi_version_compatible() { + local output=${1-} parts major minor patch extra local min_major min_minor min_patch min_extra - command -v quota-axi >/dev/null 2>&1 || return 1 - if [ -n "$timeout" ]; then - case "$timeout" in - ''|*[!0-9]*|0) return 1 ;; - esac - [ "$(type -t fm_run_timed)" = function ] || return 1 - output=$(fm_run_timed "$timeout" quota-axi --version 2>/dev/null /dev/null /dev/null 2>&1 || return 1 + if [ -n "$timeout" ]; then + case "$timeout" in + ''|*[!0-9]*|0) return 1 ;; + esac + [ "$(type -t fm_run_timed)" = function ] || return 1 + output=$(fm_run_timed "$timeout" quota-axi --version 2>/dev/null /dev/null "$FAKEBIN/quota-axi" <<'SH' #!/usr/bin/env bash if [ "${1:-}" = "--version" ]; then - printf 'quota-axi 0.1.51\n' + vcount=0 + [ -z "${QUOTA_AXI_VERSION_COUNT:-}" ] || [ ! -f "$QUOTA_AXI_VERSION_COUNT" ] || read -r vcount < "$QUOTA_AXI_VERSION_COUNT" + vcount=$((vcount + 1)) + [ -z "${QUOTA_AXI_VERSION_COUNT:-}" ] || printf '%s\n' "$vcount" > "$QUOTA_AXI_VERSION_COUNT" + if [ -n "${QUOTA_AXI_VERSION_OK_COUNT:-}" ] && [ "$vcount" -gt "$QUOTA_AXI_VERSION_OK_COUNT" ]; then + case "${QUOTA_AXI_VERSION_LATER:-fail}" in + slow) sleep 10 ;; + *) exit 42 ;; + esac + fi + if [ "${QUOTA_AXI_SLOW_VERSION:-0}" = 1 ]; then + sleep 10 + fi + if [ "${QUOTA_AXI_VERSION_FAIL:-0}" = 1 ]; then + exit 42 + fi + printf 'quota-axi %s\n' "${QUOTA_AXI_VERSION:-0.1.51}" exit 0 fi case "${QUOTA_AXI_MALFORMED:-}" in @@ -84,10 +105,40 @@ if [ "${QUOTA_AXI_UNKNOWN_EXHAUSTED:-0}" = 1 ]; then printf '{"schemaVersion":5,"providers":[{"provider":"codex","quotaSemantics":{"status":"known","effectiveAvailability":[{"scope":"all_models","status":"unknown","runway":{"status":"exhausted_now"}}]}}]}\n' exit 0 fi +if [ "${QUOTA_AXI_ALWAYS_SLOW:-0}" = 1 ] && [ "${1:-}" != "--version" ]; then + sleep 10 +fi count=0 [ ! -f "$QUOTA_AXI_COUNT" ] || read -r count < "$QUOTA_AXI_COUNT" count=$((count + 1)) printf '%s\n' "$count" > "$QUOTA_AXI_COUNT" +if [ "${QUOTA_AXI_SLOW_FIRST:-0}" = 1 ] && [ "$count" -eq 1 ] && [ "${1:-}" != "--version" ]; then + sleep 10 +fi +# Two immediate JSON failures, then one timeout: wording must not claim three slow reads. +if [ "${QUOTA_AXI_FAIL_THEN_SLOW:-0}" = 1 ] && [ "${1:-}" != "--version" ]; then + if [ "$count" -le 2 ]; then + exit 42 + fi + sleep 10 +fi +# Reset-streak sequence: timeouts on 1/2/4/5, healthy on 3, exhausted on 6+. +# Without consecutive_failures=0 after a good read, poll 4 would go terminal. +if [ "${QUOTA_AXI_RESET_STREAK:-0}" = 1 ]; then + case "$count" in + 1|2|4|5) + sleep 10 + ;; + 3) + printf '{"schemaVersion":5,"providers":[{"provider":"codex","quotaSemantics":{"status":"known","effectiveAvailability":[{"scope":"all_models","status":"known","effectivePercentRemaining":20,"runway":{"status":"through_reset"}}]}}]}\n' + exit 0 + ;; + *) + printf '{"schemaVersion":5,"providers":[{"provider":"codex","quotaSemantics":{"status":"known","effectiveAvailability":[{"scope":"all_models","status":"known","effectivePercentRemaining":0,"runway":{"status":"exhausted_now"}}]}}]}\n' + exit 0 + ;; + esac +fi if [ "${QUOTA_AXI_UNKNOWN_FIRST:-0}" = 1 ] && [ "$count" -eq 1 ]; then printf '{"schemaVersion":5,"providers":[{"provider":"codex","quotaSemantics":{"status":"unknown","effectiveAvailability":[]}}]}\n' exit 0 @@ -109,6 +160,19 @@ if [ "${QUOTA_AXI_AT_THRESHOLD:-0}" = 1 ]; then printf '{"schemaVersion":5,"providers":[{"provider":"codex","quotaSemantics":{"status":"known","effectiveAvailability":[{"scope":"all_models","status":"known","effectivePercentRemaining":%s,"runway":{"status":"through_reset"}}]}}]}\n' "$remaining" exit 0 fi +# After a first timed-out slow read (count already advanced), the next read is +# healthy and the one after that exhausts so the poll can prove it stayed live. +if [ "${QUOTA_AXI_SLOW_FIRST:-0}" = 1 ]; then + if [ "$count" -eq 2 ]; then + model_remaining=20 + runway=through_reset + else + model_remaining=0 + runway=exhausted_now + fi + printf '{"schemaVersion":5,"providers":[{"provider":"codex","quotaSemantics":{"status":"known","effectiveAvailability":[{"scope":"all_models","status":"known","effectivePercentRemaining":20,"runway":{"status":"through_reset"}},{"scope":"model:codex_bengalfox","status":"known","effectivePercentRemaining":%s,"runway":{"status":"%s"}}]}}]}\n' "$model_remaining" "$runway" + exit 0 +fi if [ "$count" -eq 1 ]; then model_remaining=20 runway=through_reset @@ -284,4 +348,152 @@ printf '%s\n' "$out" | grep -qx 'status: exhausted' || fail "known semantics wit printf '%s\n' "$out" | grep -qx 'condition_polls: 2' || fail "known semantics with unknown headroom stopped early" ok "poll preserves unknown headroom under known semantics" +# One slow (timed-out) read then a good read must keep the watch live. +rm -f "$COUNT" +out=$(QUOTA_AXI_SLOW_FIRST=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: exhausted' \ + || fail "one slow read then a good read did not stay live: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 3' \ + || fail "one slow read then a good read used unexpected poll count: $out" +ok "one slow read then a good read stays live" + +# N consecutive timed-out reads go terminal with distinct slow-read detail. +rm -f "$COUNT" +out=$(QUOTA_AXI_ALWAYS_SLOW=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "consecutive slow reads did not go terminal: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 3' \ + || fail "consecutive slow reads used unexpected poll count: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read timed out' \ + || fail "consecutive slow reads omitted slow-read detail: $out" +printf '%s\n' "$out" | grep -Fq 'missing/incompatible' \ + && fail "consecutive slow reads still used the missing/incompatible detail: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive slow reads' \ + && fail "consecutive slow reads still claimed the whole streak was slow: $out" +ok "N consecutive slow reads go terminal with slow-read detail" + +# A missing quota-axi still reports missing (distinct from a slow read). +rm -f "$COUNT" +out=$(PATH="$NO_QUOTA_BIN" QUOTA_AXI_COUNT="$COUNT" "$BASH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "missing quota-axi did not go terminal: $out" +printf '%s\n' "$out" | grep -qx 'detail: quota-axi is missing' \ + || fail "missing quota-axi did not report missing: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 1' \ + || fail "missing quota-axi was not reported immediately: $out" +printf '%s\n' "$out" | grep -Fq 'timed out' \ + && fail "missing quota-axi was mislabeled as a slow read: $out" +ok "missing quota-axi reports missing immediately" + +# An incompatible quota-axi reports incompatible (distinct from missing and slow). +rm -f "$COUNT" +out=$(QUOTA_AXI_VERSION=0.1.50 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "incompatible quota-axi did not go terminal: $out" +printf '%s\n' "$out" | grep -qx 'detail: quota-axi is incompatible' \ + || fail "incompatible quota-axi did not report incompatible: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 1' \ + || fail "incompatible quota-axi was not reported immediately: $out" +printf '%s\n' "$out" | grep -Fq 'missing' \ + && fail "incompatible quota-axi was mislabeled as missing: $out" +printf '%s\n' "$out" | grep -Fq 'timed out' \ + && fail "incompatible quota-axi was mislabeled as a slow read: $out" +ok "incompatible quota-axi reports incompatible immediately" + +# A healthy read must reset the consecutive-failure streak: two timeouts, one +# healthy, two more timeouts, then exhausted reaches the sixth poll. +rm -f "$COUNT" +out=$(QUOTA_AXI_RESET_STREAK=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: exhausted' \ + || fail "reset-streak sequence did not stay live through six polls: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 6' \ + || fail "reset-streak sequence used unexpected poll count: $out" +ok "a healthy read resets the consecutive-failure streak" + +# A slow --version probe is a timeout, not an incompatible tool. +rm -f "$COUNT" +out=$(QUOTA_AXI_SLOW_VERSION=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "slow version probe did not go terminal: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read timed out' \ + || fail "slow version probe omitted timeout detail: $out" +printf '%s\n' "$out" | grep -Fq 'incompatible' \ + && fail "slow version probe was mislabeled as incompatible: $out" +ok "slow version probe reports timeout not incompatible" + +# A failing --version probe is an execution failure, not an incompatible tool. +rm -f "$COUNT" +out=$(QUOTA_AXI_VERSION_FAIL=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "failing version probe did not go terminal: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read failed' \ + || fail "failing version probe omitted failure detail: $out" +printf '%s\n' "$out" | grep -Fq 'incompatible' \ + && fail "failing version probe was mislabeled as incompatible: $out" +printf '%s\n' "$out" | grep -Fq 'timed out' \ + && fail "failing version probe was mislabeled as a timeout: $out" +ok "failing version probe reports failure not incompatible" + +rm -f "$COUNT" +out=$(QUOTA_AXI_VERSION_FAIL=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "untimed failing version probe did not go terminal: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 3' \ + || fail "untimed failing version probe used unexpected poll count: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read failed' \ + || fail "untimed failing version probe omitted failure detail: $out" +printf '%s\n' "$out" | grep -Fq 'incompatible' \ + && fail "untimed failing version probe was mislabeled as incompatible: $out" +ok "untimed failing version probe reports failure not incompatible" + +# Mixed streak: two immediate JSON failures then one timeout - wording names the +# streak and the last cause, without calling every failure a slow read. +rm -f "$COUNT" +out=$(QUOTA_AXI_FAIL_THEN_SLOW=1 QUOTA_AXI_COUNT="$COUNT" PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "mixed failure streak did not go terminal: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read timed out' \ + || fail "mixed failure streak omitted last-cause detail: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive slow reads' \ + && fail "mixed failure streak claimed three slow reads: $out" +ok "mixed failure streak names the last cause without calling all reads slow" + +# First poll-time --version succeeds (plus a healthy JSON read); later version +# probes fail. Exactly one version launch per poll: a later execution failure +# must stay "failed", never "incompatible". +rm -f "$COUNT" "$VERSION_COUNT" +out=$(QUOTA_AXI_VERSION_OK_COUNT=1 QUOTA_AXI_VERSION_LATER=fail \ + QUOTA_AXI_COUNT="$COUNT" QUOTA_AXI_VERSION_COUNT="$VERSION_COUNT" \ + PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "later version failure did not go terminal: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 4' \ + || fail "later version failure used unexpected poll count: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read failed' \ + || fail "later version failure omitted failure detail: $out" +printf '%s\n' "$out" | grep -Fq 'incompatible' \ + && fail "later version failure was mislabeled as incompatible: $out" +# One successful poll-time version + three failing ones; no second probe per poll. +[ "$(cat "$VERSION_COUNT")" = 4 ] \ + || fail "expected exactly four version launches across the streak, got $(cat "$VERSION_COUNT" 2>/dev/null)" +ok "later version probe failure reports failed not incompatible" + +# Same shape with a later slow version probe: timeout, not incompatible, and +# still exactly one version launch per poll. +rm -f "$COUNT" "$VERSION_COUNT" +out=$(QUOTA_AXI_VERSION_OK_COUNT=1 QUOTA_AXI_VERSION_LATER=slow \ + QUOTA_AXI_COUNT="$COUNT" QUOTA_AXI_VERSION_COUNT="$VERSION_COUNT" \ + PATH="$FAKEBIN:$PATH" "$BIN/fm-procevent-quota.sh" poll --interval 0.01 --threshold 10 --provider codex --timeout 1) +printf '%s\n' "$out" | grep -qx 'status: error' \ + || fail "later slow version probe did not go terminal: $out" +printf '%s\n' "$out" | grep -qx 'condition_polls: 4' \ + || fail "later slow version probe used unexpected poll count: $out" +printf '%s\n' "$out" | grep -Fq '3 consecutive read failures; last quota-axi read timed out' \ + || fail "later slow version probe omitted timeout detail: $out" +printf '%s\n' "$out" | grep -Fq 'incompatible' \ + && fail "later slow version probe was mislabeled as incompatible: $out" +[ "$(cat "$VERSION_COUNT")" = 4 ] \ + || fail "expected exactly four version launches across the slow streak, got $(cat "$VERSION_COUNT" 2>/dev/null)" +ok "later slow version probe reports timeout not incompatible" + printf '# all fm-procevent-quota tests passed\n' From 918a5bf15c1563949842903aca2df8091bcf4f17 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micka=C3=ABl=20R=C3=A9mond?= Date: Sun, 4 Oct 2026 12:20:48 +0200 Subject: [PATCH 4/6] fix(bin): clarify scratch guidance and dirty teardown refusals (#6505) * fix(teardown): clarify scratch guidance and dirty worktree refusals Keep ship proof material outside the task worktree and distinguish untracked-only leftovers from tracked edits without changing cleanup guards. Fixes #6319 * fix(ci): Fixed ci-3 only. Both promotion outputs now replace the scout restriction and require external scratch storage and a clean worktree before done. Verification: 27 delivery tests passed, five mutations caught, restored test passed, pinned ShellCheck and syntax/whitespace checks passed. ci-1, ci-2, and ci-4 remain untouched --- bin/fm-brief.sh | 4 +- bin/fm-promote.sh | 4 ++ bin/fm-teardown.sh | 26 ++++++++++-- docs/architecture.md | 2 +- tests/fm-task-delivery.test.sh | 14 ++++++- tests/fm-teardown.test.sh | 73 ++++++++++++++++++++++++++++++++++ 6 files changed, 116 insertions(+), 7 deletions(-) diff --git a/bin/fm-brief.sh b/bin/fm-brief.sh index 2ff9c727762..871673d56bd 100755 --- a/bin/fm-brief.sh +++ b/bin/fm-brief.sh @@ -661,7 +661,9 @@ If the top-level path is the primary checkout or not the worktree you were launc # Rules $RULE1 -2. Stay inside this worktree; modify nothing outside it. +2. Keep project edits inside this worktree; keep proof and scratch output outside it, under \`$DATA/$ID/\` or a temporary directory. + Outside the worktree, write only that task material and the status and steering-inbox records authorized below. + Leave the worktree clean before reporting done. 3. Use gh-axi for GitHub operations and chrome-devtools-axi for browser operations. 4. Report status by appending one line: \`$STATUS_APPEND\` diff --git a/bin/fm-promote.sh b/bin/fm-promote.sh index 3d53e50cecf..53a63897a26 100755 --- a/bin/fm-promote.sh +++ b/bin/fm-promote.sh @@ -250,6 +250,10 @@ This task is now kind=ship with mode=$MODE$PROMOTE_FORGE_WORDS. This section supersedes every earlier brief instruction about delivery mode. These current ship instructions supersede the scout delivery rules and report-based Definition of done. Any earlier "Never push" or scout-only delivery language in this file is superseded. +This replaces the scout rule limiting outside-worktree writes to the report and status file. +Keep project edits inside this worktree; keep proof and scratch output outside it, under \`$DATA/$ID/\` or a temporary directory. +Outside the worktree, write only that task material and the status and steering-inbox records authorized below. +Leave the worktree clean before reporting done. The mode-specific Definition of done below is the current delivery contract. # Current ship safety rule diff --git a/bin/fm-teardown.sh b/bin/fm-teardown.sh index 24ed4644c76..8b9603d7e2d 100755 --- a/bin/fm-teardown.sh +++ b/bin/fm-teardown.sh @@ -65,7 +65,8 @@ # by itself causes a false refusal of landed work. # A gh lookup error falls back to the content check; if that is also inconclusive, # teardown refuses rather than risk discarding unlanded work. -# Uncommitted changes are never landed. +# Uncommitted changes are never landed; dirty refusals distinguish untracked-only +# leftovers from tracked edits and list at most ten non-exempt untracked paths. # local-only projects additionally accept work merged into the local default # branch (firstmate performs that merge after configured approval) as a fallback # for the common case where there is no remote at all. @@ -1864,6 +1865,23 @@ teardown_treehouse_return() { return 1 } +report_worktree_dirt() { + # Use the same porcelain snapshot and exemptions as the refusal predicate. + printf '%s\n' "$1" | awk ' + /^\?\? / { if (++untracked <= 10) paths = paths " " substr($0, 4) "\n"; next } + NF { tracked = 1 } + END { + if (tracked) print "uncommitted changes present (includes tracked edits)" + else print "uncommitted changes present (untracked-only leftovers)" + if (untracked) { + print "untracked paths (up to 10):" + printf "%s", paths + if (untracked > 10) print " ... additional untracked paths omitted" + } + } + ' >&2 +} + validate_worktree_teardown_safety() { local dirty_raw dirty unpushed_raw unpushed DEFAULT unmerged_raw unmerged branch [ -d "$WT" ] || return 0 @@ -1880,7 +1898,7 @@ validate_worktree_teardown_safety() { echo "Restore the git index state, or get the captain's explicit OK to discard, then --force." >&2 return 1 fi - dirty=$(printf '%s\n' "$dirty_raw" | grep -vE '^\?\? (\.claude/|\.fm-(grok|kimi)-turnend$)' | head -1 || true) + dirty=$(printf '%s\n' "$dirty_raw" | grep -vE '^\?\? (\.claude/|\.fm-(grok|kimi)-turnend$)' || true) if ! unpushed_raw=$(git -C "$WT" log --oneline HEAD --not --remotes -- 2>/dev/null); then if worktree_safety_blocked_by_lock "commits not on a remote"; then @@ -1905,14 +1923,14 @@ validate_worktree_teardown_safety() { unmerged=$(printf '%s\n' "$unmerged_raw" | head -5) if [ -n "$dirty" ] || [ -n "$unmerged" ]; then echo "REFUSED: local-only worktree $WT has work not yet merged into $DEFAULT and not on any remote." >&2 - [ -n "$dirty" ] && echo "uncommitted changes present" >&2 + [ -n "$dirty" ] && report_worktree_dirt "$dirty" [ -n "$unmerged" ] && printf 'commits not yet on %s:\n%s\n' "$DEFAULT" "$unmerged" >&2 echo "Merge the branch into local $DEFAULT first (bin/fm-merge-local.sh after the captain approves), or push to a fork/remote, or get the captain's explicit OK to discard, then --force." >&2 return 1 fi elif [ -n "$dirty" ]; then echo "REFUSED: worktree $WT has uncommitted changes." >&2 - echo "uncommitted changes present" >&2 + report_worktree_dirt "$dirty" echo "Commit them (or get the captain's explicit OK to discard, then --force)." >&2 return 1 elif [ -n "$unpushed" ]; then diff --git a/docs/architecture.md b/docs/architecture.md index 4b2b6f9cbfe..a4ae0a3f06b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -8,7 +8,7 @@ firstmate's supervisor contract and routing index for conditional procedures is ## Event-driven supervision -The declared-wait vocabulary, including the legacy "external wait" label, is owned by [`bin/fm-classify-lib.sh`](../bin/fm-classify-lib.sh); worker declaration instructions are owned by [`bin/fm-brief.sh`](../bin/fm-brief.sh). +The declared-wait vocabulary, including the legacy "external wait" label, is owned by [`bin/fm-classify-lib.sh`](../bin/fm-classify-lib.sh); worker declaration instructions and the ship worker's scratch-location and clean-worktree contract are owned by [`bin/fm-brief.sh`](../bin/fm-brief.sh). A zero-token bash watcher (`bin/fm-watch.sh`) sleeps on the fleet, classifies detected wakes in bash, and wakes the first mate only when something is actionable. Actionable wakes include captain-relevant status signals, no-verb signals without positive evidence that their crew is still executing, authenticated check output such as PR merge polling or a Relay mention, stale panes whose crew is not provably working whether their status log looks terminal or non-terminal, provably-working stale panes that persist past `FM_STALE_ESCALATE_SECS` with no wait their own worker declared, no writes to their own task worktree, and - in a home that armed `config/wedge-defer-parked-gate` - no validation gate of their own awaiting an unanswered supervisor decision, declared external waits and attended captain-held transfers that remain declared past `FM_PAUSE_RESURFACE_SECS`, and heartbeat backstop hits. diff --git a/tests/fm-task-delivery.test.sh b/tests/fm-task-delivery.test.sh index a0ced0a1a6e..f0e8f21a9dc 100755 --- a/tests/fm-task-delivery.test.sh +++ b/tests/fm-task-delivery.test.sh @@ -308,7 +308,7 @@ test_promote_refuses_a_symlinked_task_record() { # prints against a capturing fm-send.sh, and asserts on the message the worker would # actually receive - for every supported mode. test_promotion_delivers_the_real_definition_of_done() { - local home meta out sendroot payload mode id brief_dod delivered_dod + local home meta out sendroot payload mode id brief_dod delivered_dod contract home="$TMP_ROOT/promote-dod/home" sendroot="$TMP_ROOT/promote-dod/sendroot" mkdir -p "$home/state" "$sendroot/bin" @@ -356,6 +356,18 @@ STUB assert_grep "## Firstmate spec" "$payload" \ "$mode: promoted worker did not receive the Firstmate spec subsection" + # Both the delivered prompt and persisted relaunch brief are public outputs. + for contract in "$payload" "$home/data/$id/brief.md"; do + assert_grep "This replaces the scout rule limiting outside-worktree writes to the report and status file." "$contract" \ + "$mode: $contract retained the scout-only write restriction" + assert_grep "Keep project edits inside this worktree; keep proof and scratch output outside it, under \`$home/data/$id/\` or a temporary directory." "$contract" \ + "$mode: $contract omitted the ship scratch-location rule" + assert_grep "Outside the worktree, write only that task material and the status and steering-inbox records authorized below." "$contract" \ + "$mode: $contract omitted the ship outside-worktree write boundary" + assert_grep "Leave the worktree clean before reporting done." "$contract" \ + "$mode: $contract omitted the clean-before-done rule" + done + # Compare the public outputs of both real generation paths. The promoted # payload ends at its Definition of done, as does an ordinary generated # brief, so identical suffixes prove both workers receive the same contract. diff --git a/tests/fm-teardown.test.sh b/tests/fm-teardown.test.sh index a53d66b70b2..40f1ee03da7 100755 --- a/tests/fm-teardown.test.sh +++ b/tests/fm-teardown.test.sh @@ -1212,6 +1212,76 @@ test_dirty_worktree_refuses() { pass "dirty worktree is refused even when its committed work has landed (dirty always wins)" } +assert_dirty_diagnostic() { + local kind=$1 mode=$2 case_dir rc before n + case_dir=$(make_case "dirty-$kind-$mode") + write_meta "$case_dir" "$mode" ship + wt_commit_file "$case_dir" feature.txt hello + # Exercise both dirty refusal sites: remote-reachable work and local-only + # work merged into local main but absent from every remote. + if [ "$mode" = local-only ]; then + git -C "$case_dir/project" merge -q --ff-only fm/task-x1 + else + git -C "$case_dir/wt" push -q origin fm/task-x1 + fi + if [ "$kind" != untracked ]; then + printf '%s\n' 'uncommitted edit' > "$case_dir/wt/feature.txt" + # Cover index edits as well as unstaged edits. + [ "$mode" != local-only ] || git -C "$case_dir/wt" add feature.txt + fi + if [ "$kind" != tracked ]; then + mkdir "$case_dir/wt/00 proof scratch" + printf '%s\n' 'manual server log' > "$case_dir/wt/00 proof scratch/server.log" + for n in 01 02 03 04 05 06 07 08 09 10 11; do + touch "$case_dir/wt/$n-scratch.txt" + done + # Preserve the existing exemptions without counting them as leftovers. + mkdir "$case_dir/wt/.claude" + touch "$case_dir/wt/.claude/settings.local.json" "$case_dir/wt/.fm-grok-turnend" + fi + before=$(git -C "$case_dir/wt" status --porcelain) + rc=0 + run_teardown "$case_dir" > "$case_dir/stdout" 2> "$case_dir/stderr" || rc=$? + expect_code 1 "$rc" "$kind/$mode: dirty teardown must still refuse" + grep -q REFUSED "$case_dir/stderr" || fail "$kind/$mode: no refusal" + if [ "$kind" = untracked ]; then + grep -Fq 'uncommitted changes present (untracked-only leftovers)' "$case_dir/stderr" \ + || fail "$kind/$mode: missing untracked-only classification" + ! grep -q 'includes tracked edits' "$case_dir/stderr" || fail "$kind/$mode: misclassified as tracked" + else + grep -Fq 'uncommitted changes present (includes tracked edits)' "$case_dir/stderr" \ + || fail "$kind/$mode: missing tracked-edit classification" + ! grep -q 'untracked-only' "$case_dir/stderr" || fail "$kind/$mode: misclassified as untracked-only" + fi + if [ "$kind" != tracked ]; then + grep -Fq '00 proof scratch/' "$case_dir/stderr" || fail "$kind/$mode: scratch folder not named" + grep -Fxq ' 09-scratch.txt' "$case_dir/stderr" || fail "$kind/$mode: tenth path missing" + ! grep -q '10-scratch.txt\|11-scratch.txt\|\.claude/\|\.fm-grok-turnend' "$case_dir/stderr" \ + || fail "$kind/$mode: path list exceeded its bound or included exempt files" + grep -Fq 'additional untracked paths omitted' "$case_dir/stderr" || fail "$kind/$mode: no truncation notice" + else + ! grep -q 'untracked paths' "$case_dir/stderr" || fail "$kind/$mode: invented untracked paths" + fi + [ -f "$case_dir/state/task-x1.meta" ] || fail "$kind/$mode: task metadata removed" + [ "$before" = "$(git -C "$case_dir/wt" status --porcelain)" ] || fail "$kind/$mode: worktree changed" + pass "$kind/$mode: dirty refusal classifies leftovers and preserves work" +} + +test_untracked_only_refusal_diagnostic() { + assert_dirty_diagnostic untracked no-mistakes + assert_dirty_diagnostic untracked local-only +} + +test_tracked_edit_refusal_diagnostic() { + assert_dirty_diagnostic tracked no-mistakes + assert_dirty_diagnostic tracked local-only +} + +test_mixed_refusal_diagnostic() { + assert_dirty_diagnostic mixed no-mistakes + assert_dirty_diagnostic mixed local-only +} + test_gh_error_and_content_absent_refuses() { local case_dir rc case_dir=$(make_case gh-error) @@ -4337,6 +4407,9 @@ test_pr_check_records_remote_head_when_local_lags test_content_in_default_fallback_allows test_content_fallback_refreshes_stale_origin_ref test_dirty_worktree_refuses +test_untracked_only_refusal_diagnostic +test_tracked_edit_refusal_diagnostic +test_mixed_refusal_diagnostic test_gh_error_and_content_absent_refuses test_legacy_record_without_the_flag_refuses test_windowless_legacy_record_with_gone_worktree_tears_down From 42dd906d02f672652606057a193a6386b956c3f3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micka=C3=ABl=20R=C3=A9mond?= Date: Sun, 4 Oct 2026 12:21:16 +0200 Subject: [PATCH 5/6] fix(bin): escalate inbox instructions blocked by busy workers (#6518) * Escalate inbox instructions stuck behind a busy worker Count consecutive busy-deferred due doorbells durably and escalate at the configured bound without typing into the worker pane. Fixes #6445 * fix(review): Fix inbox escalation deduplication and busy streak resets * fix(review): Preserve busy inbox escalations through daemon supervision * fix(document): Correct busy-inbox escalation documentation * fix(ci): Fixed SC2034 in tests/fm-task-inbox.test.sh by including the loop counter in the failure diagnostic. Source-aware lint, all 34 inbox tests, and git diff --check pass. Behavior portable serial 6 reproduces identically on base 1f3e7696 and target 78156b86 with Pi 1.0.1: an unrelated renderer API change breaks the unchanged Calm test. No Calm changes made; that failure is addressed separately by https://github.com/kunchenguid/firstmate/pull/6516. Logs retained in scratchpad-ci/ * fix(ci): Fixed ci-2, ci-3, and ci-4: successor failures surface, reset alerts deduplicate, and oversized busy limits fall back to two. Passed 47 inbox tests, 7 focused daemon checks, all 13 mutation checks, lint, documentation checks, and diff checks. Evidence: scratchpad-ci-selected/summary.json. ci-1 remains unchanged and unwaived. Fresh live Herdr proof remains with the outer driver --- .agents/skills/afk/SKILL.md | 6 +- bin/fm-send.sh | 13 +- bin/fm-supervise-daemon.sh | 10 +- bin/fm-task-inbox-lib.sh | 56 +++++- bin/fm-watch.sh | 50 +++-- docs/architecture.md | 5 +- docs/configuration.md | 1 + docs/verification/runtime-backends.md | 17 ++ tests/fm-daemon.test.sh | 91 +++++++++ tests/fm-task-inbox.test.sh | 275 ++++++++++++++++++++++++-- 10 files changed, 456 insertions(+), 68 deletions(-) diff --git a/.agents/skills/afk/SKILL.md b/.agents/skills/afk/SKILL.md index 68b7bf1b7cd..f82de256fb8 100644 --- a/.agents/skills/afk/SKILL.md +++ b/.agents/skills/afk/SKILL.md @@ -163,13 +163,15 @@ The daemon still clears its buffer only on the backend's `empty` success verdict The daemon wraps `fm-watch.sh`, runs the watcher as a child, presents every durable wake after each actionable watcher close, classifies each presented record in bash, and acknowledges the presented generation only after routing completes. It self-handles the routine majority without consuming a firstmate turn. -Captain-relevant events, plus a bounded recheck of a declared external wait that is still declared, escalate to firstmate's context as one pre-read, single-line, batched digest. +Events selected by the routing below escalate to firstmate's context as one pre-read, single-line, batched digest. The digest is byte-bounded so every transport can carry it; when it cuts an event or omits events past its budget, it names a `state/.subsuper-digests/` file that holds every buffered event verbatim, so read that file before acting on a cut event. The captain-relevant verb set, declared-wait vocabulary, status-span classifier, and presentation-marker contract live in shared `bin/fm-classify-lib.sh`, while each supervisor owns its routing and fleet scan as a consumer of that policy. While `state/.afk` exists the daemon owns the watcher, so the watcher reverts to one-shot and lets the daemon do the triage - the two never run their triage at the same time. -Classify each wake this way: +Classify each wake this way, applying the steering-inbox exception before status-based routing: +- `stale` whose detail begins `unread firstmate instruction: stuck-busy ` or `steering-inbox busy bookkeeping unwritable: ` -> buffer the explicit inbox escalation for supervision in away and quiet mode, without consuming worker status or entering transient-stale recovery. + [`bin/fm-task-inbox-lib.sh`](../../../bin/fm-task-inbox-lib.sh) owns the busy budget and bookkeeping contract; `tests/fm-daemon.test.sh` covers this consumer boundary. - `signal` whose newly classified status span contains captain-relevant events -> escalate every event in source order. A nonterminal progress verb remains nonterminal even when its prose contains a legacy free-text token such as `PR ready`, `checks green`, `ready in branch`, or `merged`; only a bare legacy line with such a token escalates. Other signals with no captain-relevant event in the span -> self-handle. diff --git a/bin/fm-send.sh b/bin/fm-send.sh index 19562680313..cdd4ddcb95d 100755 --- a/bin/fm-send.sh +++ b/bin/fm-send.sh @@ -47,15 +47,10 @@ # instruction. There is no delivered-unconfirmed # outcome on this plane: "did the doorbell land" is no longer the question - # "was the message acted on" is, and that is answered asynchronously for an -# ordinary record by the worker's acknowledgement move into handled/. The -# watcher re-rings an unacknowledged message while its endpoint remains -# available, escalates after the bounded ladder, and instead routes a positively -# dead or missing endpoint directly to recovery without typing. An explicit -# fire-and-forget record is excluded from that ladder; when config/wait-no-turns -# is present and its ring here was skipped or failed, the watcher rings it -# exactly once more. -# bin/fm-task-inbox-lib.sh owns the record format, the doorbell line, and the -# re-ring ladder. The composer pre-check before the ring is ADVISORY only: when +# ordinary record by the worker's acknowledgement move into handled/. +# bin/fm-task-inbox-lib.sh owns the record format, doorbell line, and retry and +# escalation policy for ordinary and fire-and-forget records. +# The composer pre-check before the ring is ADVISORY only: when # the composer visibly holds pending text the ring is skipped with a notice and # the watcher re-rings an ordinary record later; no composer verdict is # delivery proof on this plane, and a failed ring never fails the send. diff --git a/bin/fm-supervise-daemon.sh b/bin/fm-supervise-daemon.sh index a2a7664f4fb..1364203efb5 100755 --- a/bin/fm-supervise-daemon.sh +++ b/bin/fm-supervise-daemon.sh @@ -7,9 +7,9 @@ # ESCALATES a batched, distilled digest to the supervisor pane on # captain-relevant events plus bounded declared-wait rechecks. This is the # token-efficient replacement for the prior always-inject daemon: routine -# signal/stale/heartbeat wakes cost zero firstmate context; only done/ -# needs-decision/blocked/failed/persistent-wedge/check-output events and a -# declared-wait recheck reach the LLM, and even then as one pre-read digest per +# signal/stale/heartbeat wakes cost zero firstmate context; routing is owned by +# .agents/skills/afk/SKILL.md (Classification policy). +# Escalated events reach the LLM as one pre-read digest per # batch window. That digest is byte-bounded (see escalate_flush); when it cuts # or omits anything it names a state/.subsuper-digests/ file holding every # buffered event verbatim. @@ -48,7 +48,7 @@ # drain and acknowledges it only after routing completes. # - Fail-safe-to-escalate: any wake the classifier cannot confidently mark # routine is escalated. -# - Bounded wedge latency: a stale pane without a declared wait is escalated +# - Bounded wedge latency: ordinary pane staleness without a declared wait escalates # only after it has been idle for STALE_ESCALATE_SECS # (configurable), rechecked once. A wedged crewmate is therefore detected # within STALE_ESCALATE_SECS + a tick, never lost. A declared wait - either a @@ -1562,6 +1562,8 @@ handle_wake() { # *) arg="${reason#signal: }" ;; esac decision=$(FM_STATUS_SPAN_ENDPOINT_FILE="$capture" classify_signal "$arg" "$state") ;; + stale:*" (unread firstmate instruction: stuck-busy "*|stale:*" (steering-inbox busy bookkeeping unwritable: "*) + decision="escalate|${reason#stale: }" ;; stale:*) kind=stale; arg="${reason#stale: }"; stale_detail="${arg#"$arg"}" case "$arg" in *" ("*) stale_detail="${arg#*" ("}"; arg="${arg%% \(*}" ;; esac task=$(window_to_task "$arg" "$state") diff --git a/bin/fm-task-inbox-lib.sh b/bin/fm-task-inbox-lib.sh index 257da1dd54d..3b87a429a21 100644 --- a/bin/fm-task-inbox-lib.sh +++ b/bin/fm-task-inbox-lib.sh @@ -28,6 +28,7 @@ # .inbox/.seq.lock serializes sequence allocation across writers # (the session and the away daemon) # .inbox/.ring-state watcher re-ring ladder: "\t\t" +# .inbox/.busy-state consecutive busy deferrals: "\t" # .inbox/.escalated oldest-message name already surfaced as stale, # so later polls suppress another escalation # .inbox/.retry-ring name of a fire-and-forget record still owed its @@ -52,10 +53,14 @@ # attempt may ring or be skipped to protect another draft in a proven pending # composer; an unsubmitted copy of this doorbell is retried. After # FM_TASK_INBOX_RING_MAX attempts without an acknowledgement it escalates. The -# caller owns the busy and recovery-grade endpoint checks: a busy pane waits, -# while a positively dead or missing endpoint skips delivery and the ladder and -# escalates directly. This library owns only the schedule and escalation marker. -# If attempt bookkeeping cannot be persisted while the record remains unhandled, +# caller owns the busy and recovery-grade endpoint checks: due actions deferred +# by a busy pane consume a separate durable consecutive-poll budget, +# FM_TASK_INBOX_BUSY_MAX. At that bound the same escalation path surfaces a +# stuck-busy reason without typing. A non-busy due check or acknowledgement resets +# this budget. Fire-and-forget retries remain outside escalation. A positively +# dead or missing endpoint skips delivery and the ladder and escalates directly. +# This library owns the schedule, durable budgets, and escalation marker. +# If delivery-attempt or busy-deferral bookkeeping fails while the record remains unhandled, # the caller surfaces that failure instead of retrying silently; a concurrently # removed inbox is a quiet no-op. Escalation deliberately queues the wake before # writing the deduplication marker: normal polls surface a message once, while a @@ -85,6 +90,7 @@ # Tunables (env): # FM_TASK_INBOX_GRACE_SECS default 90; delivery-attempt grace and spacing # FM_TASK_INBOX_RING_MAX default 3; delivery attempts before escalation +# FM_TASK_INBOX_BUSY_MAX default 2; consecutive busy-deferred due polls before escalation _FM_TASK_INBOX_LIB_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" # Both dependencies are canonical lint roots in their own right. Keep them as @@ -98,6 +104,7 @@ _FM_TASK_INBOX_LIB_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" FM_TASK_INBOX_SCHEMA='fm-task-inbox.v1' FM_TASK_INBOX_GRACE_DEFAULT=90 FM_TASK_INBOX_RING_MAX_DEFAULT=3 +FM_TASK_INBOX_BUSY_MAX_DEFAULT=2 FM_TASK_INBOX_LOCK_WAIT_DEFAULT=5 fm_task_inbox_grace_secs() { @@ -112,6 +119,40 @@ fm_task_inbox_ring_max() { printf '%s' "$m" } +fm_task_inbox_busy_max() { + local m=${FM_TASK_INBOX_BUSY_MAX:-$FM_TASK_INBOX_BUSY_MAX_DEFAULT} + case "$m" in ''|*[!0-9]*) m=$FM_TASK_INBOX_BUSY_MAX_DEFAULT ;; esac + # Check the length before numeric comparison so oversized input cannot overflow. + if [ "${#m}" -gt 9 ] || [ "$m" -eq 0 ]; then + m=$FM_TASK_INBOX_BUSY_MAX_DEFAULT + fi + printf '%s' "$m" +} + +# Persist before returning the new count, so a fresh watcher continues the same +# bounded wait. A removed or acknowledged record is a quiet no-op. +fm_task_inbox_record_busy() { # + local dir base previous count + dir=$(fm_task_inbox_dir "$1" "$2") + base=${3##*/} + { IFS=$(printf '\t') read -r previous count < "$dir/.busy-state"; } 2>/dev/null || true + [ "${previous:-}" = "$base" ] || count=0 + case "${count:-}" in ''|*[!0-9]*) count=0 ;; esac + [ -f "$3" ] || { printf '0'; return 0; } + count=$((count + 1)) + if ! { printf '%s\t%s\n' "$base" "$count" > "$dir/.busy-state"; } 2>/dev/null; then + [ -f "$3" ] || { printf '0'; return 0; } + return 1 + fi + printf '%s' "$count" +} + +fm_task_inbox_clear_busy() { # + local dir + dir=$(fm_task_inbox_dir "$1" "$2") + rm -f "$dir/.busy-state" 2>/dev/null +} + fm_task_inbox_dir() { # printf '%s/%s.inbox' "$1" "$2" } @@ -415,7 +456,7 @@ fm_task_inbox_due_action() { # local dir oldest base now grace max ladder rec_base count last dir=$(fm_task_inbox_dir "$1" "$2") if ! oldest=$(fm_task_inbox_oldest_unhandled "$1" "$2"); then - rm -f "$dir/.ring-state" "$dir/.escalated" 2>/dev/null || true + rm -f "$dir/.ring-state" "$dir/.escalated" "$dir/.busy-state" 2>/dev/null || true # The one retry ring exists only while config/wait-no-turns is present. # Absent, a mark is left untouched and the inbox stays quiet, as before. if [ -e "${FM_CONFIG_OVERRIDE:-${FM_HOME:-}/config}/wait-no-turns" ]; then @@ -443,13 +484,8 @@ fm_task_inbox_due_action() { # $ladder EOF if [ -n "$rec_base" ] && [ "$rec_base" != "$base" ]; then - # A different oldest message: the previous ladder is stale. An absent - # ladder is left alone so a dead-pane escalation, which never rings and so - # never writes one, keeps its marker (the marker check below still ignores - # a marker naming some other message). count=0 last=0 - rm -f "$dir/.escalated" 2>/dev/null || true fi case "$count" in ''|*[!0-9]*) count=0 ;; esac case "$last" in ''|*[!0-9]*) last=0 ;; esac diff --git a/bin/fm-watch.sh b/bin/fm-watch.sh index 6e51f76770a..dcd1f7c4933 100755 --- a/bin/fm-watch.sh +++ b/bin/fm-watch.sh @@ -77,12 +77,10 @@ # interrupt, signal, or restart of the worker or its # tool process. # stale: (unread firstmate instruction: ...) -# the steering-inbox ladder spent its delivery-attempt -# budget on an idle pane without an acknowledgement # stale: (steering-inbox ladder bookkeeping unwritable: ...) -# an unhandled record's ladder cannot advance; quiet -# successful attempts never wake firstmate -# (bin/fm-task-inbox-lib.sh owns the ladder policy) +# stale: (steering-inbox busy bookkeeping unwritable: ...) +# steering-inbox recovery; bin/fm-task-inbox-lib.sh owns +# delivery-attempt, busy-deferral, and unavailable-endpoint policy # check: