Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 33 additions & 1 deletion .pi/extensions/fm-branch-supervision.ts
Original file line number Diff line number Diff line change
Expand Up @@ -533,6 +533,10 @@ export default function (pi: ExtensionAPI) {
let pendingWakeGeneration = -1;
const pendingWakeMessages: string[] = [];
const lastDeliveredStaleByWindow = new Map<string, string>();
// 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<string, { message: string; generation: number }>();
const pendingMirror: MirrorItem[] = [];
const mirrorCollection: MirrorCollectionState = {
collectAnchor: null,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)) {
Expand Down Expand Up @@ -1395,6 +1425,7 @@ ${context.command}
branchBroken = "";
generation += 1;
lastDeliveredStaleByWindow.clear();
deferredStaleRechecks.clear();
if (actingAsOwner(generation)) activatePendingActionDeliveries(generation);
});

Expand Down Expand Up @@ -1437,6 +1468,7 @@ ${context.command}
pendingActionDeliveries.clear();
rehydratedActionGeneration = -1;
lastDeliveredStaleByWindow.clear();
deferredStaleRechecks.clear();
pendingMirror.length = 0;
currentMainSession = null;
mirrorCollection.collectAnchor = null;
Expand Down
4 changes: 2 additions & 2 deletions docs/pi-supervision-branch.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 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.
Expand Down Expand Up @@ -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.
122 changes: 118 additions & 4 deletions tests/fm-pi-branch-extension.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
Expand All @@ -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 = [
Expand Down Expand Up @@ -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() {
Expand Down
Loading