Skip to content

Commit 9802c61

Browse files
authored
fix(pi): close routine stale wake acknowledgment gap (#98)
* fix(pi): recheck deferred stale wake rows * no-mistakes(document): Document deferred stale-row acknowledgement
1 parent 95b7688 commit 9802c61

3 files changed

Lines changed: 153 additions & 7 deletions

File tree

‎.pi/extensions/fm-branch-supervision.ts‎

Lines changed: 33 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -533,6 +533,10 @@ export default function (pi: ExtensionAPI) {
533533
let pendingWakeGeneration = -1;
534534
const pendingWakeMessages: string[] = [];
535535
const lastDeliveredStaleByWindow = new Map<string, string>();
536+
// One same-text offer can arrive after the active turn's eligible-row
537+
// snapshot. Remember only the newest signal per window and re-run the normal
538+
// durable queue scan once the serialized turn has settled.
539+
const deferredStaleRechecks = new Map<string, { message: string; generation: number }>();
536540
const pendingMirror: MirrorItem[] = [];
537541
const mirrorCollection: MirrorCollectionState = {
538542
collectAnchor: null,
@@ -1246,7 +1250,29 @@ ${context.command}
12461250
}
12471251
deliverPendingActionDeliveries(acceptedGeneration, wakeContext.eligibleSeqs);
12481252
})
1253+
.then(() => {
1254+
if (shuttingDown || acceptedGeneration !== generation) return;
1255+
const recheck: string[] = [];
1256+
for (const candidate of messages) {
1257+
const window = staleWakeWindow(candidate);
1258+
if (!window || lastDeliveredStaleByWindow.get(window) !== candidate) continue;
1259+
const deferred = deferredStaleRechecks.get(window);
1260+
if (deferred?.generation === acceptedGeneration && deferred.message === candidate) {
1261+
deferredStaleRechecks.delete(window);
1262+
recheck.push(candidate);
1263+
} else {
1264+
lastDeliveredStaleByWindow.delete(window);
1265+
}
1266+
}
1267+
if (recheck.length > 0) enqueueWake(recheck, acceptedGeneration);
1268+
})
12491269
.catch(async (error: unknown) => {
1270+
for (const candidate of messages) {
1271+
const window = staleWakeWindow(candidate);
1272+
if (!window || lastDeliveredStaleByWindow.get(window) !== candidate) continue;
1273+
lastDeliveredStaleByWindow.delete(window);
1274+
deferredStaleRechecks.delete(window);
1275+
}
12501276
releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration));
12511277
releaseBranchLeases(acceptedGeneration);
12521278
try {
@@ -1276,7 +1302,11 @@ ${context.command}
12761302
flushPendingWakes();
12771303
}
12781304
const staleWindow = staleWakeWindow(message);
1279-
if (staleWindow && lastDeliveredStaleByWindow.get(staleWindow) === message) return;
1305+
if (staleWindow && lastDeliveredStaleByWindow.get(staleWindow) === message) {
1306+
deferredStaleRechecks.set(staleWindow, { message, generation: acceptedGeneration });
1307+
return;
1308+
}
1309+
if (staleWindow) deferredStaleRechecks.delete(staleWindow);
12801310
pendingWakeGeneration = acceptedGeneration;
12811311
if (!pendingWakeMessages.includes(message)) pendingWakeMessages.push(message);
12821312
if (urgentWake(message)) {
@@ -1395,6 +1425,7 @@ ${context.command}
13951425
branchBroken = "";
13961426
generation += 1;
13971427
lastDeliveredStaleByWindow.clear();
1428+
deferredStaleRechecks.clear();
13981429
if (actingAsOwner(generation)) activatePendingActionDeliveries(generation);
13991430
});
14001431

@@ -1437,6 +1468,7 @@ ${context.command}
14371468
pendingActionDeliveries.clear();
14381469
rehydratedActionGeneration = -1;
14391470
lastDeliveredStaleByWindow.clear();
1471+
deferredStaleRechecks.clear();
14401472
pendingMirror.length = 0;
14411473
currentMainSession = null;
14421474
mirrorCollection.collectAnchor = null;

‎docs/pi-supervision-branch.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ Main can read the durable outcome store on demand through its `fm_branch_outcome
7878

7979
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.
8080
For a given window, the first `stale:` delivery is urgent and bypasses that delay.
81-
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.
81+
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.
8282
Same-text idle repeats inside the bounded window therefore open at most one branch turn.
8383
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`.
8484
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
113113

114114
## Verification
115115

116-
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`).
116+
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`).
117117
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).
118118
The strict typecheck in `tests/fm-pi-primary-types.test.sh` pins the extension against the installed Pi package.

‎tests/fm-pi-branch-extension.test.sh‎

Lines changed: 118 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1183,8 +1183,9 @@ test_branch_coalesces_repeat_wakes_and_bypasses_for_urgent_work() {
11831183
home="$TMP_ROOT/coalesced-dispatch-home"
11841184
mkdir -p "$home/state" "$home/config"
11851185
install_pi_branch_extension_fixture "$repo"
1186-
out=$(PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \
1187-
FM_TEST_BRANCH_WAKE_COALESCE_MS=1000 DRIVER_PRELUDE="$DRIVER_PRELUDE" node --input-type=module 2>&1 <<'EOF'
1186+
PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \
1187+
FM_TEST_BRANCH_WAKE_COALESCE_MS=1000 DRIVER_PRELUDE="$DRIVER_PRELUDE" \
1188+
node --input-type=module > "$TMP_ROOT/node-output" 2>&1 <<'EOF'
11881189
const prelude = process.env.DRIVER_PRELUDE;
11891190
await eval(`(async () => { ${prelude}; globalThis.__t = { dispatch, settle, fire, mainUserMessages, home }; })()`);
11901191
const { dispatch, settle, fire, mainUserMessages, home } = globalThis.__t;
@@ -1208,6 +1209,13 @@ if ((globalThis.__fmPrompts ?? []).length !== 1) {
12081209
}
12091210
12101211
const repeated = "stale: assets waiting-for-merge";
1212+
// This presentation-only probe does not run the branch's real drain loop.
1213+
// Consume its synthetic row when the prompt starts so the deferred durable-row
1214+
// recheck remains an empty no-op; the actor-scoped acknowledgement behavior is
1215+
// exercised below through the real wake-drain script.
1216+
globalThis.__fmOnBranchPrompt = async () => {
1217+
writeFileSync(`${home}/state/.wake-queue`, "");
1218+
};
12111219
const staleStarted = Date.now();
12121220
for (let index = 0; index < 3; index += 1) {
12131221
const offer = dispatch(repeated);
@@ -1221,6 +1229,7 @@ await new Promise((resolve) => setTimeout(resolve, 450));
12211229
if ((globalThis.__fmPrompts ?? []).length !== 2) {
12221230
throw new Error(`same-text repeats opened ${globalThis.__fmPrompts.length} branch turns`);
12231231
}
1232+
delete globalThis.__fmOnBranchPrompt;
12241233
12251234
const statusPath = `${home}/state/branch-driver.status`;
12261235
const urgentStatusLines = [
@@ -1318,11 +1327,116 @@ if (!readFileSync(`${home}/state/.wake-queue`, "utf8").includes("branch-driver.s
13181327
throw new Error("shutdown cleared the durable wake instead of leaving it queued");
13191328
}
13201329
EOF
1321-
)
13221330
status=$?
1331+
out=$(cat "$TMP_ROOT/node-output")
13231332
expect_code 0 "$status" "Pi branch must coalesce same-text routine floods and bypass the delay for urgent work: $out"
13241333
[ -z "$out" ] || fail "Pi branch wake-coalescing test printed output: $out"
1325-
pass "Pi branch coalesces same-text routine floods and bypasses the delay for urgent work"
1334+
1335+
home="$TMP_ROOT/coalesced-routine-ack-home"
1336+
mkdir -p "$home/state" "$home/config"
1337+
PLUGIN="$repo/.pi/extensions/fm-branch-supervision.ts" FM_HOME="$home" FM_ROOT_OVERRIDE="$ROOT" \
1338+
FM_TEST_BRANCH_WAKE_COALESCE_MS=1000 DRIVER_PRELUDE="$DRIVER_PRELUDE" \
1339+
node --input-type=module > "$TMP_ROOT/node-output" 2>&1 <<'EOF'
1340+
const prelude = process.env.DRIVER_PRELUDE;
1341+
await eval(`(async () => { ${prelude}; globalThis.__t = { bus, makeOffer, settle, home, realRoot }; })()`);
1342+
const { bus, makeOffer, settle, home, realRoot } = globalThis.__t;
1343+
import { appendFileSync, readFileSync } from "node:fs";
1344+
import { spawnSync } from "node:child_process";
1345+
1346+
const queue = `${home}/state/.wake-queue`;
1347+
const repeated = "stale: assets waiting-for-merge";
1348+
let promptCount = 0;
1349+
let ackCount = 0;
1350+
let appendAfterSnapshot = false;
1351+
1352+
globalThis.__fmExecuteBranchBash = async (context) => {
1353+
const actor = spawnSync(
1354+
"bash",
1355+
["-c", '. "$1"; fm_lease_actor', "_", `${realRoot}/bin/fm-lease-lib.sh`],
1356+
{ encoding: "utf8", cwd: context.cwd, env: context.env },
1357+
);
1358+
if (actor.status !== 0 || actor.stdout.trim() !== "branch") {
1359+
throw new Error(`branch bash actor resolution failed: ${actor.stdout}${actor.stderr}`);
1360+
}
1361+
const result = spawnSync("bash", ["-c", context.command], {
1362+
encoding: "utf8",
1363+
cwd: context.cwd,
1364+
env: context.env,
1365+
});
1366+
return {
1367+
content: [{ type: "text", text: `${result.stdout}${result.stderr}` }],
1368+
details: { stdout: result.stdout, stderr: result.stderr, exitCode: result.status },
1369+
isError: result.status !== 0,
1370+
};
1371+
};
1372+
1373+
function appendStale(seq) {
1374+
appendFileSync(queue, `${seq}\t${seq}\tstale\tbranch-driver\t${repeated}\n`);
1375+
const offer = makeOffer(repeated);
1376+
bus.emit("fm-branch-supervision:dispatch", offer);
1377+
if (!offer.accepted) throw new Error(`stale row ${seq} was not accepted`);
1378+
}
1379+
1380+
async function runWakeDrain(session, args) {
1381+
const bash = session.options.customTools.find((tool) => tool.name === "bash");
1382+
const result = await bash.execute(
1383+
`routine-ack-${promptCount}-${ackCount}`,
1384+
{ command: ["bin/fm-wake-drain.sh", ...args].join(" ") },
1385+
undefined,
1386+
undefined,
1387+
{},
1388+
);
1389+
if (result.isError) throw new Error(`wake drain failed: ${JSON.stringify(result)}`);
1390+
return result.details;
1391+
}
1392+
1393+
globalThis.__fmOnBranchPrompt = async ({ session }) => {
1394+
promptCount += 1;
1395+
if (appendAfterSnapshot) {
1396+
appendAfterSnapshot = false;
1397+
appendStale(5);
1398+
}
1399+
const drained = await runWakeDrain(session, []);
1400+
const ack = drained.stderr.match(/--ack-through ([0-9]+) --recovery-generation ([A-Za-z0-9._-]+)/);
1401+
if (!ack) throw new Error(`drain did not return its acknowledgement command: ${drained.stderr}`);
1402+
const report = session.options.customTools.find((tool) => tool.name === "fm_branch_report");
1403+
const result = await report.execute(
1404+
`routine-${promptCount}`,
1405+
{ task: "branch-driver", verdict: "routine", summary: "routine stale classification", wake: repeated },
1406+
undefined,
1407+
undefined,
1408+
{},
1409+
);
1410+
if (result.isError) throw new Error(`routine report failed: ${JSON.stringify(result)}`);
1411+
await runWakeDrain(session, ["--ack-through", ack[1], "--recovery-generation", ack[2]]);
1412+
ackCount += 1;
1413+
};
1414+
1415+
// All three rows exist before the first serialized pre-drain scan. The
1416+
// repeated offers must collapse to one prompt whose one acknowledgement owns
1417+
// the complete snapshot.
1418+
appendStale(1);
1419+
appendStale(2);
1420+
appendStale(3);
1421+
await settle(() => ackCount === 1, "one acknowledgement for the pre-scan duplicates");
1422+
if (promptCount !== 1) throw new Error(`pre-scan duplicates opened ${promptCount} prompts`);
1423+
if (readFileSync(queue, "utf8") !== "") throw new Error("pre-scan duplicate rows remained queued after acknowledgement");
1424+
await new Promise((resolve) => setTimeout(resolve, 100));
1425+
if (promptCount !== 1 || ackCount !== 1) throw new Error("an already-consumed pre-scan duplicate opened another turn");
1426+
1427+
// Row 5 arrives only after row 4's eligible-row snapshot has been published.
1428+
// It therefore needs one deferred normal recheck after row 4 settles.
1429+
appendAfterSnapshot = true;
1430+
appendStale(4);
1431+
await settle(() => ackCount === 3, "deferred acknowledgement for the post-snapshot duplicate");
1432+
if (promptCount !== 3) throw new Error(`post-snapshot duplicate produced ${promptCount - 1} prompts instead of two`);
1433+
if (readFileSync(queue, "utf8") !== "") throw new Error("post-snapshot duplicate remained in the durable wake queue");
1434+
EOF
1435+
status=$?
1436+
out=$(cat "$TMP_ROOT/node-output")
1437+
expect_code 0 "$status" "Pi branch must recheck an accepted identical stale row after the current snapshot settles: $out"
1438+
[ -z "$out" ] || fail "Pi branch routine acknowledgement regression printed output: $out"
1439+
pass "Pi branch coalesces routine floods, bypasses urgent delay, and rechecks post-snapshot stale rows"
13261440
}
13271441

13281442
test_requested_healthy_outcome_and_unsolicited_routine_outcome_delivery() {

0 commit comments

Comments
 (0)