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
30 changes: 27 additions & 3 deletions .pi/extensions/fm-branch-supervision.ts
Original file line number Diff line number Diff line change
Expand Up @@ -512,6 +512,14 @@ export default function (pi: ExtensionAPI) {
// so a prompt can prove that it created a durable outcome after claiming its
// wake rows without relying on provider text or incidental session shape.
let durableReportRevision = 0;
// The task set the wake being handled right now may be reported on, fixed
// deterministically from the eligible rows before a signal or stale prompt
// opens and cleared when it settles: exactly the tasks those rows resolve
// to. fm_branch_report refuses every other task id during such a prompt,
// `fleet` included, so a report typed from memory about a task the wake
// never named is never stored or delivered. Null outside a wake prompt and
// during a heartbeat review, which is not scoped by task.
let wakeTaskScope: { rows: string[]; tasks: Set<string> } | null = null;
let mainStreaming = false;
let shuttingDown = false;
// Bumps at every session replacement so a stale chain continuation from the
Expand Down Expand Up @@ -921,6 +929,13 @@ export default function (pi: ExtensionAPI) {
return presentUnprocessedOutcomes(expectedGeneration);
}

function wakeScopeRefusal(task: string): string {
if (!wakeTaskScope || wakeTaskScope.tasks.has(task)) return "";
const named = [...wakeTaskScope.tasks].sort().join(", ");
const rows = wakeTaskScope.rows.join(", ");
return `report refused: the wake being handled (row ${rows}) names ${named}, not ${task}; report only that task, never fleet or a task from memory`;
}

function createReportTool(toolGeneration: number): ToolDefinition {
return {
name: "fm_branch_report",
Expand Down Expand Up @@ -956,6 +971,10 @@ export default function (pi: ExtensionAPI) {
};
}
const verdict = verdictRaw as Verdict;
const scopeRefusal = wakeScopeRefusal(task);
if (scopeRefusal) {
return { content: [{ type: "text", text: scopeRefusal }], details: undefined, isError: true };
}
const appendArgs = ["append", "--task", task, "--verdict", verdict, "--summary", summary, "--silent", String(silent)];
if (wake) appendArgs.push("--wake", wake);
if (!actingAsOwner(toolGeneration)) {
Expand Down Expand Up @@ -1215,9 +1234,14 @@ ${context.command}
// the drain; that residual is accepted by the confused-agent-grade boundary.
const reportRevisionBeforePrompt = durableReportRevision;
const entryOffset = sessionManager.getEntries().length;
await session.prompt(
`FIRSTMATE SUPERVISION WAKE: ${message}\n\nHandle this per your operating procedure and finish with fm_branch_report.`,
);
wakeTaskScope = heartbeat ? null : { rows: [...scope.eligibleSeqs], tasks: new Set(scope.eligibleTasks) };
try {
await session.prompt(
`FIRSTMATE SUPERVISION WAKE: ${message}\n\nHandle this per your operating procedure and finish with fm_branch_report.`,
);
} finally {
wakeTaskScope = null;
}
const providerError = settledPromptProviderError(sessionManager, entryOffset);
if (providerError) {
const detail = `supervision branch provider failed after construction: ${providerError}`;
Expand Down
53 changes: 47 additions & 6 deletions .pi/extensions/lib/fm-branch-dispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,15 @@ export interface UnreadWakeScope {
* `eligible` is false.
*/
eligibleSeqs: string[];
/**
* The exact task ids the eligible signal/stale rows name (a signal row by
* its status-log key, a stale row through the task metadata recording that
* endpoint). The branch may report only these tasks while it handles the
* wake; `fleet` or a task it merely remembers is refused (docs/
* pi-supervision-branch.md "Components and their owners"). Empty for a
* heartbeat, which is not scoped by task.
*/
eligibleTasks: string[];
/**
* True only when this scan itself is untrustworthy: the queue or its
* metadata could not be read, a line fails the structural tab-field check,
Expand All @@ -47,8 +56,22 @@ export interface UnreadWakeScope {
corrupted: boolean;
}

const EMPTY_SCOPE: UnreadWakeScope = { status: "empty", eligible: false, projects: [], eligibleSeqs: [], corrupted: false };
const UNSAFE_SCOPE: UnreadWakeScope = { status: "unsafe", eligible: false, projects: [], eligibleSeqs: [], corrupted: true };
const EMPTY_SCOPE: UnreadWakeScope = {
status: "empty",
eligible: false,
projects: [],
eligibleSeqs: [],
eligibleTasks: [],
corrupted: false,
};
const UNSAFE_SCOPE: UnreadWakeScope = {
status: "unsafe",
eligible: false,
projects: [],
eligibleSeqs: [],
eligibleTasks: [],
corrupted: true,
};

// scopeForUnreadWake is the single owner of branch-eligibility classification
// (docs/pi-supervision-branch.md "Autonomy"; docs/watcher-continuity.md
Expand Down Expand Up @@ -92,6 +115,9 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak

const projects = new Set<string>();
const metadata = new Map<string, string>();
// The task id behind each key a signal or stale row may carry: the task id
// itself, or the endpoint its metadata records.
const taskByKey = new Map<string, string>();
try {
for (const name of readdirSync(state)) {
if (!name.endsWith(".meta")) continue;
Expand All @@ -101,14 +127,19 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak
const window = fields.find((line) => line.startsWith("window="))?.slice(7) ?? "";
if (project) {
metadata.set(task, project);
if (window) metadata.set(window, project);
taskByKey.set(task, task);
if (window) {
metadata.set(window, project);
taskByKey.set(window, task);
}
}
}
} catch {
return UNSAFE_SCOPE;
}

const eligibleSeqs: string[] = [];
const eligibleTasks = new Set<string>();
for (const line of rows) {
const fields = line.split("\t");
if (fields.length < 5 || !/^[0-9]+$/.test(fields[1])) return UNSAFE_SCOPE;
Expand All @@ -126,18 +157,21 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak
continue;
}
let project = "";
let task = "";
if (kind === "signal") {
const task = key.replace(/\.(?:status|turn-ended)$/, "");
task = key.replace(/\.(?:status|turn-ended)$/, "");
project = metadata.get(task) ?? "";
} else if (kind === "stale") {
task = taskByKey.get(key) ?? taskByKey.get(key.replace(/^fm-/, "")) ?? "";
project = metadata.get(key) ?? metadata.get(key.replace(/^fm-/, "")) ?? "";
} else {
// A kind fm_wake_append never emits: structural corruption, not an
// ordinary main-only row.
return UNSAFE_SCOPE;
}
if (!project) return UNSAFE_SCOPE;
if (!project || !task) return UNSAFE_SCOPE;
projects.add(project);
eligibleTasks.add(task);
eligibleSeqs.push(seq);
}
const eligible = eligibleSeqs.length > 0;
Expand All @@ -148,7 +182,14 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak
// empty eligible set, so reading eligibility off the claim set rather than
// off the heartbeat flag changes no pre-existing outcome and keeps a
// heartbeat from being offered with nothing to hand over.)
return { status: eligible ? "safe" : "unsafe", eligible, projects: [...projects], eligibleSeqs, corrupted: false };
return {
status: eligible ? "safe" : "unsafe",
eligible,
projects: [...projects],
eligibleSeqs,
eligibleTasks: [...eligibleTasks],
corrupted: false,
};
}

// The exact state-relative filename bin/fm-wake-drain.sh reads for a
Expand Down
16 changes: 15 additions & 1 deletion bin/fm-branch-outcome.sh
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@
# before append and published only after the cache update; processed-init
# rebuilds every cache before publishing it, so interruption or upgrade
# fails closed without making each drain scan lifetime history.
# bin/fm-teardown.sh removes a retired task's cache with its other records,
# and append skips the cache for a task that has neither a live meta nor a
# status log (the outcome itself is still stored), so the branch's report
# of a teardown it just performed leaves no index behind.
# Main-actor drain calls processed-init under the outcome lock when that
# ready marker is absent or invalid, on every harness; only a genuine store
# fault keeps the lost-wake backstop skipped.
Expand Down Expand Up @@ -460,7 +464,17 @@ case "$CMD" in
"$SEQ" "$(date +%s)" "$(json_escape "$TASK")" "$(json_escape "$WAKE")" \
"$VERDICT" "$(json_escape "$SUMMARY")" "$SILENT" "$CAPTURED_STATUS_ENDPOINT" \
"$(json_escape "$CAPTURED_STATUS_IDENT")" >> "$STORE"
if ! write_outcome_index "$TASK" "$SEQ" || ! publish_outcome_index_ready "$SEQ"; then
# A task with neither a live meta nor a status log is retired: the branch
# reports the teardown it just performed, and writing the index here would
# recreate the footprint teardown removed. The outcome itself is still
# stored and delivered; only the reader-less cache is skipped.
if { [ -e "$STATE/$TASK.meta" ] || [ -e "$STATE/$TASK.status" ]; } \
&& ! write_outcome_index "$TASK" "$SEQ"; then
fm_lock_release "$LOCK"
echo "error: outcome was stored but its bounded task index could not be updated" >&2
exit 1
fi
if ! publish_outcome_index_ready "$SEQ"; then
fm_lock_release "$LOCK"
echo "error: outcome was stored but its bounded task index could not be updated" >&2
exit 1
Expand Down
2 changes: 2 additions & 0 deletions bin/fm-branch-prompt.sh
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,8 @@ Stay terse: your context is a cost.
Do not re-read files the drain just printed.
Never use shell background operators for supervision; the watcher and extension own continuity.
Never call fm_branch_report speculatively - only after the event is actually handled or a refusal/lease conflict genuinely ended your handling.
The tool refuses a task the wake being handled did not name, fleet included (a heartbeat review is not scoped by task); a refusal means you reached for a task from memory, so report the wake's own task, never retry with another id.
An acknowledgement that consumed nothing says so and names the exact command for the current wake; run that printed command, do not drain again.

# Recovery playbook (verbatim copy of the tracked skill)

Expand Down
19 changes: 17 additions & 2 deletions bin/fm-guard.sh
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,10 @@
# bounded). Independent alarms (queued wakes, worktree tangle) are never
# suppressed by that dedup. Normal wake handling (watcher briefly down between a
# wake and the next supervision resume) stays inside the grace window and stays
# silent. Always exits 0: the guard warns, it never blocks.
# silent. The queued-wakes warning stays silent for the supervision branch
# actor (FM_SUPERVISION_ACTOR=branch), because that actor runs guarded commands
# while handling exactly the queued rows its grant covers and can drain nothing
# else. Always exits 0: the guard warns, it never blocks.
set -u

SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
Expand All @@ -52,6 +55,12 @@ STALE_BANNER_MARKER="$STATE/.guard-watcher-stale-banner"
. "$SCRIPT_DIR/fm-tangle-lib.sh"
# shellcheck source=bin/fm-supervision-lib.sh
. "$SCRIPT_DIR/fm-supervision-lib.sh"
# shellcheck source=bin/fm-lease-lib.sh
. "$SCRIPT_DIR/fm-lease-lib.sh"

# The current actor (fm_lease_actor is the one owner of that identity); a
# malformed value is a wiring bug elsewhere, so the guard just warns as main.
GUARD_ACTOR=$(fm_lease_actor 2>/dev/null) || GUARD_ACTOR=main

# Deterministic episode key from the qualitative down-state (the failing
# condition), NOT the beacon mtime: under the auto-arm model a healthy
Expand Down Expand Up @@ -232,10 +241,16 @@ fi
# Queued wakes are an independent hazard; warn whenever they are pending, even if
# a watcher is alive. Kept after the banner so the no-watcher alarm reads first.
# Dedup of the watcher-down banner never suppresses this warning.
# The supervision branch is the exception: it runs guarded commands (fm-peek,
# fm-crew-state) in the middle of handling the very rows that are queued, and
# "drain them before anything else" mid-handling reads as "an earlier wake is
# still pending", which is what made it re-run a previous acknowledgement in a
# loop. The branch can act on nothing outside its grant anyway, so for that
# actor the guard stays silent about queued rows.
if "$queue_pending"; then
if [ "$READ_ONLY" -eq 1 ]; then
echo "WARNING: queued wakes pending - left untouched because this session lacks verified fleet-lock ownership." >&2
else
elif [ "$GUARD_ACTOR" != branch ]; then
echo "WARNING: queued wakes pending - drain them with bin/fm-wake-drain.sh before anything else." >&2
fi
fi
Expand Down
5 changes: 3 additions & 2 deletions bin/fm-teardown.sh
Original file line number Diff line number Diff line change
Expand Up @@ -2570,7 +2570,8 @@ cleanup_firstmate_home_children() {
"$sub_state/$child_id.pi-ext.ts" \
"$sub_state/$child_id.grok-turnend-token" "$sub_state/$child_id.kimi-turnend-token" \
"$sub_state/$child_id.muse-session" "$sub_state/$child_id.muse-session-current" \
"$sub_state/$child_id.cursor-session" "$sub_state/$child_id.reconcile-nudged"
"$sub_state/$child_id.cursor-session" "$sub_state/$child_id.reconcile-nudged" \
"$sub_state/.$child_id.branch-outcome-index"
done
}

Expand Down Expand Up @@ -2919,7 +2920,7 @@ rm -f "$STATE/$ID.turn-ended" \
"$STATE/$ID.muse-session-current" "$STATE/$ID.cursor-session" \
"$STATE/$ID.control-relaunch" "$STATE/$ID.control-relaunch.meta-prior" \
"$STATE/$ID.control-relaunch.brief-prior" "$STATE/$ID.control-relaunch.note" \
"$STATE/$ID.reconcile-nudged"
"$STATE/$ID.reconcile-nudged" "$STATE/.$ID.branch-outcome-index"
# The steering inbox (bin/fm-task-inbox-lib.sh) is runtime state for the
# retired endpoint; teardown only runs after landing is confirmed, so any
# leftover unhandled steer here is moot rather than unlanded work.
Expand Down
45 changes: 42 additions & 3 deletions bin/fm-wake-drain.sh
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ RECOVERY_ACK_REQUIRED=false
RECOVERY_ACK_MOVED=false
ACK_THROUGH=
ACK_GENERATION=
ACK_REMOVED=0
PRESENTED_MAX=0
ACK_FINGERPRINTS=
ACK_NOTICE_FINGERPRINTS=
PRESENTATION_LOCK_TIMEOUT=${FM_STATUS_PRESENTATION_LOCK_TIMEOUT:-10}
Expand Down Expand Up @@ -144,6 +146,19 @@ require_branch_eligible_rows() {
}
}

# The highest sequence this actor has already been presented: the branch's
# grant is exactly its current prompt's rows, and main's claim file is what its
# last drain printed. Read BEFORE an ack re-claims, so a row that arrived since
# presentation is never named as "the current wake" the caller may acknowledge
# unseen. 0 when nothing is on record.
presented_max_row() { # <rows-file>
if rows_file_valid "$1" 2>/dev/null; then
awk '$1 ~ /^[0-9]+$/ && $1 > max { max=$1 } END { print max + 0 }' "$1"
else
printf '0\n'
fi
}

case "${1:-}" in
'') ;;
--ack-through)
Expand Down Expand Up @@ -594,6 +609,11 @@ reclaim_stale_branch_grant_locked || exit 1
[ "$ACTOR" != branch ] || require_branch_eligible_rows || exit 1

if [ -n "$ACK_THROUGH" ]; then
if [ "$ACTOR" = branch ]; then
PRESENTED_MAX=$(presented_max_row "$ELIGIBLE_ROWS_FILE") || exit 1
else
PRESENTED_MAX=$(presented_max_row "$MAIN_ROWS_FILE") || exit 1
fi
if [ "$ACTOR" = main ]; then
# Preserve main's original whole-cutoff acknowledgement contract: rows may
# arrive after presentation but before the printed ack runs, and a direct
Expand Down Expand Up @@ -649,6 +669,7 @@ if [ -n "$ACK_THROUGH" ]; then
exit 1
}
fi
ACK_REMOVED=$(( $(awk 'END { print NR }' "$FM_WAKE_QUEUE") - $(awk 'END { print NR }' "$DRAIN_TMP") ))
if [ ! -s "$DRAIN_TMP" ]; then
fm_recovery_marker_ack "$RECOVERY_MARKER" "$ACK_GENERATION"
RECOVERY_ACK_STATUS=$?
Expand Down Expand Up @@ -679,9 +700,27 @@ if [ -n "$ACK_THROUGH" ]; then
fi
fm_lock_release "$FM_WAKE_QUEUE_LOCK"
DRAIN_LOCK_HELD=false
if [ "$RECOVERY_ACK_MOVED" = true ]; then
printf 'wake drain: acknowledged wakes through %s, but a newer recovery episode is pending; re-run bin/fm-wake-drain.sh and use the new WAKE_ACK_REQUIRED command\n' \
"$ACK_THROUGH" >&2
if [ "$ACK_REMOVED" -eq 0 ] && [ "$PRESENTED_MAX" -gt "$ACK_THROUGH" ]; then
# Nothing at or below the cutoff was this actor's to consume, while a
# presented row above it is still waiting: the caller acknowledged an
# earlier wake, not the one it is handling. Say so, and name the exact
# command for the current wake, so the remedy is never "drain again" (which
# re-presents the same row and invites the same stale acknowledgement).
# The generation is the marker's current one; only a retired marker cannot
# be named because the next drain opens a fresh generation for it.
case "$RECOVERY_MARKER_TOKEN" in
pending:*|announced:*)
printf 'wake drain: nothing was acknowledged through %s (none of your presented wake rows is at or below it); the current wake is row %s: run bin/fm-wake-drain.sh --ack-through %s --recovery-generation %s after handling it\n' \
"$ACK_THROUGH" "$PRESENTED_MAX" "$PRESENTED_MAX" "${RECOVERY_MARKER_TOKEN##*:}" >&2
;;
*)
printf 'wake drain: nothing was acknowledged through %s (none of your presented wake rows is at or below it); the current wake is row %s: re-run bin/fm-wake-drain.sh and use the WAKE_ACK_REQUIRED command it prints\n' \
"$ACK_THROUGH" "$PRESENTED_MAX" >&2
;;
esac
elif [ "$RECOVERY_ACK_MOVED" = true ]; then
printf 'wake drain: acknowledged wakes through %s (%s row(s) consumed), but a newer recovery episode is pending; re-run bin/fm-wake-drain.sh and use the new WAKE_ACK_REQUIRED command\n' \
"$ACK_THROUGH" "$ACK_REMOVED" >&2
fi
exit 0
fi
Expand Down
5 changes: 3 additions & 2 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,8 +112,9 @@ It suppresses failed-looking closes when the same identity-matched watcher is he
Cursor's `bin/fm-turnend-guard-cursor.sh` hook is the same between-turns shape in one synchronous step: it parks the awaited `stop` hook on the arm wrapper and translates an actionable close into one `followup_message`, with a generation baton that makes an older park still running after the next `stop` claim stand down instead of leaking a stale duplicate wake.
The existing turn-end guard remains the final backstop for every harness-engine protocol, with pi-signed sharing Pi's protocol, the `--claude` mode cooperating with the auto-arm claim, and Cursor's `--cursor` mode rendering a block as one bounded follow-up because its `stop` step cannot be blocked.
Its `--restart` mode signals only the watcher recorded in the current home's `state/.watch.lock`, so restarting one home cannot kill sibling secondmate watchers.
A pull-based guard (`bin/fm-guard.sh`) warns through supervision tool output if the primary checkout is tangled, if work, process-event sources, or Relay polling has an unhealthy model-aware supervision verdict, or if queued wakes are waiting to be drained.
The drain script calls that guard after presenting the queue; records remain durable, and may keep the queued-wakes warning visible, until the exact generation-bound acknowledgement printed by the drain succeeds after handling.
A pull-based guard (`bin/fm-guard.sh`) warns through supervision tool output if the primary checkout is tangled or if work, process-event sources, or Relay polling has an unhealthy model-aware supervision verdict; on main it also warns when queued wakes are waiting to be drained.
The drain script calls that guard after presenting the queue; records remain durable until the exact generation-bound acknowledgement printed by the drain succeeds after handling, and main may keep the queued-wakes warning visible until then.
The Pi supervision branch's deliberate queued-wake warning exception is owned by [`pi-supervision-branch.md`](pi-supervision-branch.md#components-and-their-owners).
It leads with a prominent bordered tangle banner, while `bin/fm-guard.sh` owns the watcher-down banner and reminder policy so repeated guarded commands stay noisy without reprinting the full banner in the same episode.
On every verified primary harness, tracked hook integration gives the primary session a push-based backstop: when work, a process-event source, or Relay polling needs supervision and no supervision owner provably holds this home with a fresh beacon, blocking-capable Stop hooks block and nonblocking turn-end integrations force one bounded follow-up.
The guard covers the main primary and genuinely marked secondmate homes, exempts child crewmate/scout worktrees, is loop-safe per harness, and is documented in [turnend-guard.md](turnend-guard.md).
Expand Down
Loading
Loading