From 8494a17dd4fc9b90bdba248132c45136e6de3f6a Mon Sep 17 00:00:00 2001 From: Ivan Li Date: Mon, 31 Aug 2026 20:55:31 +0800 Subject: [PATCH 1/2] fix(pi): recheck deferred stale wake rows --- .pi/extensions/fm-branch-supervision.ts | 34 ++++++- docs/pi-supervision-branch.md | 2 +- tests/fm-pi-branch-extension.test.sh | 122 +++++++++++++++++++++++- 3 files changed, 152 insertions(+), 6 deletions(-) diff --git a/.pi/extensions/fm-branch-supervision.ts b/.pi/extensions/fm-branch-supervision.ts index fa4853e56cd..eef706f4b02 100644 --- a/.pi/extensions/fm-branch-supervision.ts +++ b/.pi/extensions/fm-branch-supervision.ts @@ -533,6 +533,10 @@ export default function (pi: ExtensionAPI) { let pendingWakeGeneration = -1; const pendingWakeMessages: string[] = []; const lastDeliveredStaleByWindow = new Map(); + // One same-text offer can arrive after the active turn's eligible-row + // snapshot. Remember only the newest signal per window and re-run the normal + // durable queue scan once the serialized turn has settled. + const deferredStaleRechecks = new Map(); const pendingMirror: MirrorItem[] = []; const mirrorCollection: MirrorCollectionState = { collectAnchor: null, @@ -1246,7 +1250,29 @@ ${context.command} } deliverPendingActionDeliveries(acceptedGeneration, wakeContext.eligibleSeqs); }) + .then(() => { + if (shuttingDown || acceptedGeneration !== generation) return; + const recheck: string[] = []; + for (const candidate of messages) { + const window = staleWakeWindow(candidate); + if (!window || lastDeliveredStaleByWindow.get(window) !== candidate) continue; + const deferred = deferredStaleRechecks.get(window); + if (deferred?.generation === acceptedGeneration && deferred.message === candidate) { + deferredStaleRechecks.delete(window); + recheck.push(candidate); + } else { + lastDeliveredStaleByWindow.delete(window); + } + } + if (recheck.length > 0) enqueueWake(recheck, acceptedGeneration); + }) .catch(async (error: unknown) => { + for (const candidate of messages) { + const window = staleWakeWindow(candidate); + if (!window || lastDeliveredStaleByWindow.get(window) !== candidate) continue; + lastDeliveredStaleByWindow.delete(window); + deferredStaleRechecks.delete(window); + } releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration)); releaseBranchLeases(acceptedGeneration); try { @@ -1276,7 +1302,11 @@ ${context.command} flushPendingWakes(); } const staleWindow = staleWakeWindow(message); - if (staleWindow && lastDeliveredStaleByWindow.get(staleWindow) === message) return; + if (staleWindow && lastDeliveredStaleByWindow.get(staleWindow) === message) { + deferredStaleRechecks.set(staleWindow, { message, generation: acceptedGeneration }); + return; + } + if (staleWindow) deferredStaleRechecks.delete(staleWindow); pendingWakeGeneration = acceptedGeneration; if (!pendingWakeMessages.includes(message)) pendingWakeMessages.push(message); if (urgentWake(message)) { @@ -1395,6 +1425,7 @@ ${context.command} branchBroken = ""; generation += 1; lastDeliveredStaleByWindow.clear(); + deferredStaleRechecks.clear(); if (actingAsOwner(generation)) activatePendingActionDeliveries(generation); }); @@ -1437,6 +1468,7 @@ ${context.command} pendingActionDeliveries.clear(); rehydratedActionGeneration = -1; lastDeliveredStaleByWindow.clear(); + deferredStaleRechecks.clear(); pendingMirror.length = 0; currentMainSession = null; mirrorCollection.collectAnchor = null; diff --git a/docs/pi-supervision-branch.md b/docs/pi-supervision-branch.md index ed3035b633b..60216a779a8 100644 --- a/docs/pi-supervision-branch.md +++ b/docs/pi-supervision-branch.md @@ -78,7 +78,7 @@ Main can read the durable outcome store on demand through its `fm_branch_outcome Before starting a branch prompt, the extension holds non-urgent offers for a bounded 250 millisecond window and combines their unique wake text into one turn. For a given window, the first `stale:` delivery is urgent and bypasses that delay. -An identical same-text `stale:` repeat for that window is store-only while that text remains the last-delivered stale, so it opens no branch turn and causes no re-prompt. +An identical same-text `stale:` repeat for that window opens no immediate branch turn while that text remains in flight, but it schedules one normal pre-drain recheck after the current handling settles so a row outside the current claim is not stranded. Same-text idle repeats inside the bounded window therefore open at most one branch turn. Urgent status-tail bypass applies when the final nonblank line starts with `done:`, `needs-decision:`, `blocked:`, or `failed:`, or contains `login`, `credential`, `credentials`, `PR ready`, `ready for review`, or `checks green`. The status-tail check reads only a validated direct child of the home `state/` directory whose filename follows the shared task-id grammar, reads at most 4 KiB through a no-follow, nonblocking descriptor, and treats an unreadable or non-regular target as non-urgent. diff --git a/tests/fm-pi-branch-extension.test.sh b/tests/fm-pi-branch-extension.test.sh index fd1d543acb8..355cb6fc779 100644 --- a/tests/fm-pi-branch-extension.test.sh +++ b/tests/fm-pi-branch-extension.test.sh @@ -1183,8 +1183,9 @@ test_branch_coalesces_repeat_wakes_and_bypasses_for_urgent_work() { home="$TMP_ROOT/coalesced-dispatch-home" mkdir -p "$home/state" "$home/config" install_pi_branch_extension_fixture "$repo" - out=$(PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \ - FM_TEST_BRANCH_WAKE_COALESCE_MS=1000 DRIVER_PRELUDE="$DRIVER_PRELUDE" node --input-type=module 2>&1 <<'EOF' + PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \ + FM_TEST_BRANCH_WAKE_COALESCE_MS=1000 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 = { dispatch, settle, fire, mainUserMessages, home }; })()`); const { dispatch, settle, fire, mainUserMessages, home } = globalThis.__t; @@ -1208,6 +1209,13 @@ if ((globalThis.__fmPrompts ?? []).length !== 1) { } const repeated = "stale: assets waiting-for-merge"; +// This presentation-only probe does not run the branch's real drain loop. +// Consume its synthetic row when the prompt starts so the deferred durable-row +// recheck remains an empty no-op; the actor-scoped acknowledgement behavior is +// exercised below through the real wake-drain script. +globalThis.__fmOnBranchPrompt = async () => { + writeFileSync(`${home}/state/.wake-queue`, ""); +}; const staleStarted = Date.now(); for (let index = 0; index < 3; index += 1) { const offer = dispatch(repeated); @@ -1221,6 +1229,7 @@ await new Promise((resolve) => setTimeout(resolve, 450)); if ((globalThis.__fmPrompts ?? []).length !== 2) { throw new Error(`same-text repeats opened ${globalThis.__fmPrompts.length} branch turns`); } +delete globalThis.__fmOnBranchPrompt; const statusPath = `${home}/state/branch-driver.status`; const urgentStatusLines = [ @@ -1318,11 +1327,116 @@ if (!readFileSync(`${home}/state/.wake-queue`, "utf8").includes("branch-driver.s throw new Error("shutdown cleared the durable wake instead of leaving it queued"); } EOF - ) status=$? + out=$(cat "$TMP_ROOT/node-output") expect_code 0 "$status" "Pi branch must coalesce same-text routine floods and bypass the delay for urgent work: $out" [ -z "$out" ] || fail "Pi branch wake-coalescing test printed output: $out" - pass "Pi branch coalesces same-text routine floods and bypasses the delay for urgent work" + + home="$TMP_ROOT/coalesced-routine-ack-home" + mkdir -p "$home/state" "$home/config" + PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \ + FM_TEST_BRANCH_WAKE_COALESCE_MS=1000 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 = { bus, makeOffer, settle, home, realRoot }; })()`); +const { bus, makeOffer, settle, home, realRoot } = globalThis.__t; +import { appendFileSync, readFileSync } from "node:fs"; +import { spawnSync } from "node:child_process"; + +const queue = `${home}/state/.wake-queue`; +const repeated = "stale: assets waiting-for-merge"; +let promptCount = 0; +let ackCount = 0; +let appendAfterSnapshot = false; + +globalThis.__fmExecuteBranchBash = async (context) => { + const actor = spawnSync( + "bash", + ["-c", '. "$1"; fm_lease_actor', "_", `${realRoot}/bin/fm-lease-lib.sh`], + { encoding: "utf8", cwd: context.cwd, env: context.env }, + ); + if (actor.status !== 0 || actor.stdout.trim() !== "branch") { + throw new Error(`branch bash actor resolution failed: ${actor.stdout}${actor.stderr}`); + } + const result = spawnSync("bash", ["-c", context.command], { + encoding: "utf8", + cwd: context.cwd, + env: context.env, + }); + return { + content: [{ type: "text", text: `${result.stdout}${result.stderr}` }], + details: { stdout: result.stdout, stderr: result.stderr, exitCode: result.status }, + isError: result.status !== 0, + }; +}; + +function appendStale(seq) { + appendFileSync(queue, `${seq}\t${seq}\tstale\tbranch-driver\t${repeated}\n`); + const offer = makeOffer(repeated); + bus.emit("fm-branch-supervision:dispatch", offer); + if (!offer.accepted) throw new Error(`stale row ${seq} was not accepted`); +} + +async function runWakeDrain(session, args) { + const bash = session.options.customTools.find((tool) => tool.name === "bash"); + const result = await bash.execute( + `routine-ack-${promptCount}-${ackCount}`, + { command: ["bin/fm-wake-drain.sh", ...args].join(" ") }, + undefined, + undefined, + {}, + ); + if (result.isError) throw new Error(`wake drain failed: ${JSON.stringify(result)}`); + return result.details; +} + +globalThis.__fmOnBranchPrompt = async ({ session }) => { + promptCount += 1; + if (appendAfterSnapshot) { + appendAfterSnapshot = false; + appendStale(5); + } + const drained = await runWakeDrain(session, []); + const ack = drained.stderr.match(/--ack-through ([0-9]+) --recovery-generation ([A-Za-z0-9._-]+)/); + if (!ack) throw new Error(`drain did not return its acknowledgement command: ${drained.stderr}`); + const report = session.options.customTools.find((tool) => tool.name === "fm_branch_report"); + const result = await report.execute( + `routine-${promptCount}`, + { task: "branch-driver", verdict: "routine", summary: "routine stale classification", wake: repeated }, + undefined, + undefined, + {}, + ); + if (result.isError) throw new Error(`routine report failed: ${JSON.stringify(result)}`); + await runWakeDrain(session, ["--ack-through", ack[1], "--recovery-generation", ack[2]]); + ackCount += 1; +}; + +// All three rows exist before the first serialized pre-drain scan. The +// repeated offers must collapse to one prompt whose one acknowledgement owns +// the complete snapshot. +appendStale(1); +appendStale(2); +appendStale(3); +await settle(() => ackCount === 1, "one acknowledgement for the pre-scan duplicates"); +if (promptCount !== 1) throw new Error(`pre-scan duplicates opened ${promptCount} prompts`); +if (readFileSync(queue, "utf8") !== "") throw new Error("pre-scan duplicate rows remained queued after acknowledgement"); +await new Promise((resolve) => setTimeout(resolve, 100)); +if (promptCount !== 1 || ackCount !== 1) throw new Error("an already-consumed pre-scan duplicate opened another turn"); + +// Row 5 arrives only after row 4's eligible-row snapshot has been published. +// It therefore needs one deferred normal recheck after row 4 settles. +appendAfterSnapshot = true; +appendStale(4); +await settle(() => ackCount === 3, "deferred acknowledgement for the post-snapshot duplicate"); +if (promptCount !== 3) throw new Error(`post-snapshot duplicate produced ${promptCount - 1} prompts instead of two`); +if (readFileSync(queue, "utf8") !== "") throw new Error("post-snapshot duplicate remained in the durable wake queue"); +EOF + status=$? + out=$(cat "$TMP_ROOT/node-output") + expect_code 0 "$status" "Pi branch must recheck an accepted identical stale row after the current snapshot settles: $out" + [ -z "$out" ] || fail "Pi branch routine acknowledgement regression printed output: $out" + pass "Pi branch coalesces routine floods, bypasses urgent delay, and rechecks post-snapshot stale rows" } test_requested_healthy_outcome_and_unsolicited_routine_outcome_delivery() { From 8e2e5fba8129a2b028f1b8286292fe1751c4e823 Mon Sep 17 00:00:00 2001 From: Ivan Li Date: Mon, 31 Aug 2026 21:39:36 +0800 Subject: [PATCH 2/2] no-mistakes(document): Document deferred stale-row acknowledgement --- docs/pi-supervision-branch.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/pi-supervision-branch.md b/docs/pi-supervision-branch.md index 60216a779a8..3ec8c42716a 100644 --- a/docs/pi-supervision-branch.md +++ b/docs/pi-supervision-branch.md @@ -78,7 +78,7 @@ Main can read the durable outcome store on demand through its `fm_branch_outcome Before starting a branch prompt, the extension holds non-urgent offers for a bounded 250 millisecond window and combines their unique wake text into one turn. For a given window, the first `stale:` delivery is urgent and bypasses that delay. -An identical same-text `stale:` repeat for that window opens no immediate branch turn while that text remains in flight, but it schedules one normal pre-drain recheck after the current handling settles so a row outside the current claim is not stranded. +An identical same-text `stale:` repeat for that window opens no immediate branch turn while the text remains in flight, but it records one deferred normal durable-queue recheck after the active eligible-row snapshot settles; that recheck claims and acknowledges any remaining branch-owned row, while an already-consumed row is an empty no-op. Same-text idle repeats inside the bounded window therefore open at most one branch turn. Urgent status-tail bypass applies when the final nonblank line starts with `done:`, `needs-decision:`, `blocked:`, or `failed:`, or contains `login`, `credential`, `credentials`, `PR ready`, `ready for review`, or `checks green`. The status-tail check reads only a validated direct child of the home `state/` directory whose filename follows the shared task-id grammar, reads at most 4 KiB through a no-follow, nonblocking descriptor, and treats an unreadable or non-regular target as non-urgent. @@ -113,6 +113,6 @@ What is new is only the attended path: outside away mode, the branch absorbs the ## Verification -Portable regressions: `tests/fm-pi-branch-extension.test.sh` (dispatch, default-on eligibility, main-only classification, three-way outcome classification and delivery, routine store-only delivery, verdict-specific main envelopes, crash-before-ack replay, duplicate-wake handoff and `wake_seq` idempotency, wake-ack and branch-lease ordering, wake coalescing and urgent bypass, same-text stale suppression, pre-turn-end complete-current-request mirroring, fleet-event ownership, main outcome access, eligible-row claim lifecycle, partial pre-drain recheck, fallback, filter, model-visible outcome typing and plain-instruction fallback, cache key, persistence, model pin and searchable picker, effort pin), `tests/fm-branch-supervision.test.sh` (prompt stability, store append-only, leases, guards, non-branch-home invariance), the branch-offer, heartbeat-offer, heartbeat-not-ridden-by-a-check, and main-only-check-class tests in `tests/fm-pi-watch-extension.test.sh`, the recovery test in `tests/fm-session-start.test.sh`, and the per-actor consume regression in `tests/fm-wake-queue.test.sh`). +Portable regressions: `tests/fm-pi-branch-extension.test.sh` (dispatch, default-on eligibility, main-only classification, three-way outcome classification and delivery, routine store-only delivery, verdict-specific main envelopes, crash-before-ack replay, duplicate-wake handoff and `wake_seq` idempotency, wake-ack and branch-lease ordering, wake coalescing and urgent bypass, same-text stale handling before and after the active eligible-row snapshot, including the empty no-op and deferred acknowledgement boundaries, pre-turn-end complete-current-request mirroring, fleet-event ownership, main outcome access, eligible-row claim lifecycle, partial pre-drain recheck, fallback, filter, model-visible outcome typing and plain-instruction fallback, cache key, persistence, model pin and searchable picker, effort pin), `tests/fm-branch-supervision.test.sh` (prompt stability, store append-only, leases, guards, non-branch-home invariance), the branch-offer, heartbeat-offer, heartbeat-not-ridden-by-a-check, and main-only-check-class tests in `tests/fm-pi-watch-extension.test.sh`, the recovery test in `tests/fm-session-start.test.sh`, and the per-actor consume regression in `tests/fm-wake-queue.test.sh`). Live guard: `FM_PI_BRANCH_LIVE_E2E=1 tests/fm-pi-branch-live-e2e.test.sh` exercises the real installed Pi SDK's custom-message conversion and branch-session surfaces; its no-model probe isolates ambient Gemini credentials in the child process, and its version-specific result belongs in [docs/verification/runtime-backends.md](verification/runtime-backends.md). The strict typecheck in `tests/fm-pi-primary-types.test.sh` pins the extension against the installed Pi package.