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
76 changes: 59 additions & 17 deletions .pi/extensions/fm-primary-pi-watch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,21 @@
// consumes at the user message_start carrying the exact wake text; either
// event finishes the pending record, and a still-unconsumed record rides the
// replacement handoff.
//
// Restore versus delivery (stated once here):
// delivering is true while the serialized pending-wake pump is in flight.
// A later actionable close restores a successor even when the previous wake's
// branch settlement has not resolved. Failure closes during delivering still
// defer their bounded retry until that delivery settles.
// restoring holds the one in-flight restore promise for the generation, so a
// concurrent caller awaits that restore's real result instead of reading a
// synthetic success: the pump needs the recovery marker to confirm handling
// and the failure text to annotate the wake it delivers.
// Because that side-path restore runs outside the pump's own try/catch, it
// surfaces its typed failure - and any thrown persistence error - on its own
// rather than waiting for the hung delivery to settle. A shared failing
// restore therefore reaches main twice, standalone and appended to the later
// wake; re-arming is idempotent, and silent stopped monitoring is not.
import { spawn, spawnSync, type ChildProcess } from "node:child_process";
import { createHash } from "node:crypto";
import { existsSync, mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync } from "node:fs";
Expand Down Expand Up @@ -91,6 +106,11 @@ type UnconsumedWake = {
pending: PendingActionableClose;
};

type ContinuityRestoration = {
failure: string;
recovery?: { generation: string; watcherPid: string };
};

type SessionGeneration = {
id: number;
stopping: boolean;
Expand All @@ -99,7 +119,8 @@ type SessionGeneration = {
retryTimer: ReturnType<typeof setTimeout> | null;
cleanupTimer: ReturnType<typeof setTimeout> | null;
retryFailures: number;
restoring: boolean;
restoring: Promise<ContinuityRestoration> | null;
delivering: boolean;
seq: number;
awayStandby: boolean;
awayPoll: ReturnType<typeof setInterval> | null;
Expand Down Expand Up @@ -434,7 +455,8 @@ function createGeneration(): SessionGeneration {
retryTimer: null,
cleanupTimer: null,
retryFailures: 0,
restoring: false,
restoring: null,
delivering: false,
seq: 0,
awayStandby: false,
awayPoll: null,
Expand Down Expand Up @@ -723,7 +745,7 @@ export default function (pi: ExtensionAPI) {
if (retiring) retiring.kill("SIGTERM");
return;
}
if (!owner.awayStandby || owner.child || owner.retryTimer || owner.restoring) return;
if (!owner.awayStandby || owner.child || owner.retryTimer || owner.restoring || owner.delivering) return;
owner.awayStandby = false;
if (!startArm(owner).ok) owner.awayStandby = true;
else void processPendingActionables(owner);
Expand Down Expand Up @@ -781,13 +803,35 @@ export default function (pi: ExtensionAPI) {
owner.cleanupTimer = timer;
}

async function restoreContinuity(owner: SessionGeneration, predecessorArmPid: string): Promise<ContinuityRestoration> {
if (!owner.restoring) {
owner.restoring = restoreAfterActionableClose(owner, predecessorArmPid).finally(() => {
owner.restoring = null;
});
}
return await owner.restoring;
}

async function processPendingActionables(owner: SessionGeneration): Promise<void> {
if (!generationIsLive(owner) || owner.restoring || owner.pendingActionables.length === 0) return;
if (!generationIsLive(owner) || owner.pendingActionables.length === 0) return;
if (awayModeActive()) {
owner.awayStandby = true;
return;
}
owner.restoring = true;
if (owner.delivering) {
const newest = owner.pendingActionables[owner.pendingActionables.length - 1];
if (newest && !newest.delivered && !owner.unconsumedWakes.has(newest.token)) {
try {
const restoration = await restoreContinuity(owner, newest.predecessorArmPid);
if (restoration.failure) surfaceFailure(owner, restoration.failure);
} catch (error) {
const detail = error instanceof Error ? error.message : String(error);
surfaceFailure(owner, `watcher: FAILED - Pi extension could not restore watcher continuity\n${detail}`);
}
}
return;
}
owner.delivering = true;
const attemptedCleanup = new Set<string>();
try {
while (generationIsLive(owner) && !awayModeActive() && owner.pendingActionables.length > 0) {
Expand Down Expand Up @@ -832,7 +876,7 @@ export default function (pi: ExtensionAPI) {
// A new restoration supersedes whatever became of the previous
// successor; only a failure during this delivery is retried after it.
owner.deferredClose = null;
const restoration = await restoreAfterActionableClose(owner, pending.predecessorArmPid);
const restoration = await restoreContinuity(owner, pending.predecessorArmPid);
if (!generationIsLive(owner)) {
settleClaim("failed");
releaseClaim();
Expand Down Expand Up @@ -875,8 +919,8 @@ export default function (pi: ExtensionAPI) {
const detail = error instanceof Error ? error.message : String(error);
surfaceFailure(owner, `watcher: FAILED - Pi extension could not deliver an actionable wake\n${detail}`);
} finally {
owner.delivering = false;
if (generationIsLive(owner)) {
owner.restoring = false;
if (owner.pendingActionables.some((pending) => pending.delivered)) schedulePendingCleanup(owner);
// No bare arm is launched here. A generation without a child at this
// point has either delivered a typed restoration failure after its
Expand Down Expand Up @@ -978,10 +1022,7 @@ export default function (pi: ExtensionAPI) {
});
}

async function restoreAfterActionableClose(owner: SessionGeneration, predecessorArmPid: string): Promise<{
failure: string;
recovery?: { generation: string; watcherPid: string };
}> {
async function restoreAfterActionableClose(owner: SessionGeneration, predecessorArmPid: string): Promise<ContinuityRestoration> {
let failure = "";
for (let attempt = 0; attempt <= retryLimit; attempt += 1) {
if (!generationIsLive(owner)) return { failure: "" };
Expand Down Expand Up @@ -1147,11 +1188,12 @@ export default function (pi: ExtensionAPI) {
return;
}
if (!generationIsLive(owner)) return;
if (owner.restoring) {
// The pipeline is still delivering the wake this successor was
// started for. A verified successor that failed on its own keeps its
// bounded retry for the end of that delivery; an unready child closing
// here was retired by the restoration itself.
if (owner.restoring || owner.delivering) {
// Restore still owns this child, or the serialized pump is still
// delivering the wake this successor was started for. A verified
// successor that failed on its own keeps its bounded retry for the
// end of that delivery; an unready child closing here was retired by
// the restoration itself.
if (verified && !armRetired.has(armChild)) {
owner.deferredClose = { message: classification.message, predecessorArmPid: predecessor };
}
Expand All @@ -1166,7 +1208,7 @@ export default function (pi: ExtensionAPI) {
settleReadiness(false);
releaseChild();
if (!generationIsLive(owner)) return;
if (owner.restoring) return;
if (owner.restoring || owner.delivering) return;
scheduleRetry(owner, `watcher: FAILED - Pi extension arm child ${id} failed: ${error.message}`, String(armChild.pid ?? ""));
});
return {
Expand Down
2 changes: 2 additions & 0 deletions bin/fm-test-run.sh
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,7 @@ family_for_basename() {
fm-rovo-signals-live-e2e.test.sh|\
fm-opencode-primary-live-e2e.test.sh|fm-pi-branch-live-e2e.test.sh|\
fm-pi-branch-responsiveness-live-e2e.test.sh|\
fm-pi-hung-delivery-herdr-e2e.test.sh|\
fm-pi-primary-live-e2e.test.sh|fm-omp-primary-live-e2e.test.sh|fm-stow-horizon-live-e2e.test.sh|\
fm-sessionstart-hook-live-e2e.test.sh|fm-sessionstart-instruction-refresh-live-e2e.test.sh|\
fm-quota-array-dispatch-live-e2e.test.sh|fm-send-secondmate-marker-herdr-e2e.test.sh|\
Expand Down Expand Up @@ -692,6 +693,7 @@ tests/fm-pending-reply.test.sh 24679
tests/fm-pi-branch-extension.test.sh 22239
tests/fm-pi-branch-live-e2e.test.sh 56
tests/fm-pi-branch-responsiveness-live-e2e.test.sh 21
tests/fm-pi-hung-delivery-herdr-e2e.test.sh 23
tests/fm-pi-primary-live-e2e.test.sh 20
tests/fm-pi-watch-extension.test.sh 42970
tests/fm-pointer-check.test.sh 1314
Expand Down
24 changes: 24 additions & 0 deletions docs/verification/supervision.md
Original file line number Diff line number Diff line change
Expand Up @@ -469,6 +469,30 @@ grok 0.2.103 (89c3d36fb6f1) [stable]

Pi 0.81.1 repeated the continuity and clean-exit lifecycle on 2026-07-23 after the Calm presentation changes.

Pi later-cycle restore while an earlier branch settlement is still hung is covered by the portable watcher suite plus this isolated-lab live guard:

```sh
tests/fm-pi-watch-extension.test.sh
FM_PI_HUNG_DELIVERY_HERDR_E2E=1 tests/fm-pi-hung-delivery-herdr-e2e.test.sh
```

The live guard launches a real Pi in a named non-default Herdr lab, hangs branch settlement, requires two further worker completions to keep a live successor, and checks the watcher beacon across an unattended interval.
It does not claim turn-end marker identity, branch-session rotation, or ready-to-validate classification.

On 2026-09-11 that live guard was run against Pi 0.85.1 in a named non-default Herdr 0.8.2 lab:

```sh
FM_PI_HUNG_DELIVERY_HERDR_E2E=1 tests/fm-pi-hung-delivery-herdr-e2e.test.sh
```

Observed output:

```text
ok - isolated Pi hung-settlement Herdr lab restored later-cycle successors and kept a fresh beacon

all fm-pi-hung-delivery-herdr-e2e tests passed
```

Pi same-process session-transition ownership was verified on 2026-07-27 against the tracked extension with a faithful in-process factory rebind (module cache retained, real arm children):

```sh
Expand Down
1 change: 1 addition & 0 deletions docs/watcher-continuity.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ The omp adapter uses the bounded readiness and retirement deadlines owned by its
When that retained arm later closes, its actual close is classified as a new supervised event without replaying the earlier fallback.
After the configured retry bound is exhausted, it delivers the original wake with a typed continuity-restoration failure even if every successor arm hung without reporting readiness.
This is deliberate Option B ordering: the fleet is protected before the model handles the wake whenever restoration succeeds, but the model is never left blind when it does not.
Pi's restore-versus-delivery split is owned by the generation contract in `.pi/extensions/fm-primary-pi-watch.ts`: successor restore is not held across wake delivery.

Claude's Stop hook starts the successor arm at the next Stop after the handling turn, rather than before notification as Pi, omp, and OpenCode do.
The durable wake queue preserves actionable events during the residual active-turn window, and the bounded turn-end guard enforces recovery at Stop when no watcher is live and no open generation claim is still deciding, so a finished, hung, or identity-mismatched claim cannot suppress it ([`turnend-guard.md`](turnend-guard.md#harness-integrations) owns that boundary).
Expand Down
164 changes: 164 additions & 0 deletions tests/fm-pi-hung-delivery-herdr-e2e.test.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
#!/usr/bin/env bash
# Opt-in real Pi/Herdr regression for later-cycle restore while branch
# settlement is hung. Every Herdr call is routed through fm-herdr-lab.sh.
# The named lab is never default. Isolated FM_HOME. Abort before any provider
# call. The portable suite pins the same restore/delivery split without Pi.
set -u

# shellcheck source=tests/lib.sh
. "$(dirname "${BASH_SOURCE[0]}")/lib.sh"
# shellcheck source=tests/herdr-test-safety.sh
. "$(dirname "${BASH_SOURCE[0]}")/herdr-test-safety.sh"

if [ "${FM_PI_HUNG_DELIVERY_HERDR_E2E:-0}" != 1 ]; then
echo "skip: set FM_PI_HUNG_DELIVERY_HERDR_E2E=1 to run the isolated Pi hung-settlement Herdr regression"
exit 0
fi

command -v pi >/dev/null 2>&1 || fail "pi not found"
command -v herdr >/dev/null 2>&1 || fail "herdr not found"
command -v jq >/dev/null 2>&1 || fail "jq not found"

herdr_forget_inherited_pane

LAB_HELPER=${HERDR_LAB_HELPER:-$ROOT/bin/fm-herdr-lab.sh}
SESSION=$("$LAB_HELPER" name fm-pi-hung-deliv)
TMP_ROOT=$(fm_test_tmproot fm-pi-hung-delivery-herdr-e2e)
HOME_DIR="$TMP_ROOT/home"
PROJECT="$TMP_ROOT/project"
PI_DIR="$TMP_ROOT/pi-agent"
HANG_EXT="$TMP_ROOT/hang-dispatch.ts"
LAUNCH="$TMP_ROOT/launch-pi.sh"
PANE=
UNATTENDED_SECS=8

cleanup() {
local rc=$?
trap - EXIT
if [ -n "$PANE" ]; then
"$LAB_HELPER" run "$SESSION" pane send-text "$PANE" '/quit' >/dev/null 2>&1 || true
sleep 1
fi
if ! "$LAB_HELPER" teardown "$SESSION"; then
rc=1
fi
rm -rf "$TMP_ROOT"
exit "$rc"
}
trap cleanup EXIT

"$LAB_HELPER" provision "$SESSION"

mkdir -p "$HOME_DIR"/{state,config,data} "$PROJECT" "$PI_DIR"
printf '# Synthetic isolated hung-settlement lab\n' > "$PROJECT/AGENTS.md"

cat > "$HANG_EXT" <<'EOF'
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
export default function (pi: ExtensionAPI) {
pi.on("project_trust", () => ({ trusted: "yes", remember: false }));
pi.events?.on?.("fm-branch-supervision:dispatch", (data: { accept?: (p: Promise<void>) => void }) => {
data.accept?.(new Promise(() => {}));
});
pi.on("before_agent_start", (_event, ctx) => {
ctx.abort();
});
}
EOF

wait_markers() {
local i=0
while [ "$i" -lt 120 ]; do
if [ -f "$HOME_DIR/state/.pi-watch-extension-loaded" ] && [ -f "$HOME_DIR/state/.pi-turnend-extension-loaded" ]; then
return 0
fi
sleep 0.5
i=$((i + 1))
done
return 1
}

cycle_count() { # <needle>
local needle=$1 file=$HOME_DIR/state/.watch-cycle-exits.log count=0
if [ -f "$file" ]; then
count=$(grep -cE "$needle" "$file" || true)
fi
printf '%s\n' "${count:-0}"
}

wait_started() { # <need>
local need=$1 i=0
while [ "$i" -lt 80 ]; do
[ "$(cycle_count 'successor=started:')" -ge "$need" ] && return 0
sleep 0.5
i=$((i + 1))
done
return 1
}

beacon_age_s() {
local beat=$HOME_DIR/state/.last-watcher-beat now mtime
[ -f "$beat" ] || return 1
now=$(date +%s)
mtime=$(stat -c %Y "$beat")
printf '%s\n' "$((now - mtime))"
}

OUT=$("$LAB_HELPER" run "$SESSION" workspace create --cwd "$PROJECT" --label fm-pi-hung-deliv --no-focus) \
|| fail "Herdr lab workspace create failed"
PANE=$(printf '%s' "$OUT" | jq -er '.result.root_pane.pane_id') \
|| fail "Herdr lab workspace create omitted pane id"

cat > "$LAUNCH" <<EOF
#!/usr/bin/env bash
set -u
echo \$\$ > $(printf %q "$HOME_DIR/state/.lock")
exec env \\
FM_HOME=$(printf %q "$HOME_DIR") \\
FM_ROOT_OVERRIDE=$(printf %q "$ROOT") \\
PI_CODING_AGENT_DIR=$(printf %q "$PI_DIR") \\
FM_POLL=1 \\
FM_SIGNAL_GRACE=0 \\
FM_HEARTBEAT=600 \\
FM_CHECK_INTERVAL=999999 \\
pi --approve --no-session --no-context-files --no-extensions \\
-e $(printf %q "$HANG_EXT") \\
-e $(printf %q "$ROOT/.pi/extensions/fm-primary-turnend-guard.ts") \\
-e $(printf %q "$ROOT/.pi/extensions/fm-primary-pi-watch.ts") \\
--model openai-codex/gpt-5.6-sol --thinking low
EOF
chmod +x "$LAUNCH"

"$LAB_HELPER" run "$SESSION" pane run "$PANE" "$LAUNCH" >/dev/null \
|| fail "could not launch isolated Pi in the named Herdr lab"

wait_markers || fail "Pi watch and turn-end extensions did not load in the isolated lab"

cat > "$HOME_DIR/state/hung-cycle.meta" <<EOF
window=$SESSION:synthetic-worker
backend=herdr
kind=ship
mode=direct-PR
worktree=$PROJECT
project=synthetic-hung-settlement
EOF

printf 'done: simulated hung completion 1\n' >> "$HOME_DIR/state/hung-cycle.status"
wait_started 1 || fail "first actionable close did not start a successor while settlement hung"

printf 'done: simulated hung completion 2\n' >> "$HOME_DIR/state/hung-cycle.status"
wait_started 2 || fail "second actionable close did not restore a successor while settlement hung"

printf 'done: simulated hung completion 3\n' >> "$HOME_DIR/state/hung-cycle.status"
wait_started 3 || fail "third actionable close did not restore a successor while settlement hung"

none=$(cycle_count 'successor=none')
[ "$none" -eq 0 ] || fail "hung settlement left successor=none after later-cycle restores (none=$none)"

sleep "$UNATTENDED_SECS"
age=$(beacon_age_s) || fail "watcher beacon missing after the unattended interval"
[ "$age" -le $((UNATTENDED_SECS + 5)) ] || fail "watcher beacon went stale during the unattended interval (age=${age}s)"
[ "$(cycle_count 'successor=started:')" -ge 3 ] || fail "unattended interval lost restored successors"
[ "$(cycle_count 'successor=none')" -eq 0 ] || fail "unattended interval recorded successor=none"

printf 'ok - isolated Pi hung-settlement Herdr lab restored later-cycle successors and kept a fresh beacon\n'
printf '\nall fm-pi-hung-delivery-herdr-e2e tests passed\n'
Loading
Loading