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
25 changes: 22 additions & 3 deletions .omp/extensions/fm-branch-supervision-omp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -411,6 +411,9 @@ function collectMainDialog(sessionManager: ReadonlyEntries, collection: MirrorCo
export default function (pi: ExtensionAPI) {
let branch: AgentSession | null = null;
let branchBroken = "";
// During a signal or stale prompt, reports are restricted to the tasks named
// by that prompt's granted rows. Heartbeat reviews remain unscoped.
let wakeTaskScope: { rows: string[]; tasks: Set<string> } | null = null;
let mainStreaming = false;
let shuttingDown = false;
// Advanced once per cold-start arm (session_start). There is no live-handoff
Expand Down Expand Up @@ -449,6 +452,13 @@ export default function (pi: ExtensionAPI) {
let mainModel: { provider: string; id: string } | null = null;
let mainModelRegistry: ModelRegistry | null = null;

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`;
}

// Main's own current effort needs no such tracking: OMP answers it directly
// on demand, including at wake time. A value that is not one of the branch's
// pinnable levels ("inherit", the auto sentinel, or undefined) means "follow
Expand Down Expand Up @@ -688,6 +698,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);
return enqueueDelivery(async () => {
Expand Down Expand Up @@ -913,9 +927,14 @@ ${context.command}
if (grant !== "published") throw new Error("could not record the branch's eligible row snapshot");
// A row can still arrive between this re-check and the model starting
// the drain; that residual is accepted by the confused-agent-grade boundary.
await session.prompt(
`FIRSTMATE SUPERVISION WAKE: ${message}\n\nHandle this per your operating procedure. Do not finish this turn until you have completed, in order: fm_branch_report, the exact WAKE_ACK_REQUIRED command, and release of every task lease you claimed.`,
);
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. Do not finish this turn until you have completed, in order: fm_branch_report, the exact WAKE_ACK_REQUIRED command, and release of every task lease you claimed.`,
);
} finally {
wakeTaskScope = null;
}
if (!(await releaseEligibleRowsSnapshot(state, wakeGrantScript, String(acceptedGeneration)))) {
throw new Error("could not release the branch's settled wake-row grant");
}
Expand Down
23 changes: 17 additions & 6 deletions .omp/extensions/lib/fm-branch-dispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ export interface UnreadWakeScope {
* `eligible` is false.
*/
eligibleSeqs: string[];
/** Exact task ids named by the signal or stale rows granted to the branch. */
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 +49,8 @@ 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/omp-supervision-branch.md "Components and their owners" and
Expand Down Expand Up @@ -92,6 +94,7 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak

const projects = new Set<string>();
const metadata = new Map<string, string>();
const taskByKey = new Map<string, string>();
try {
for (const name of readdirSync(state)) {
if (!name.endsWith(".meta")) continue;
Expand All @@ -101,14 +104,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 +134,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 +159,7 @@ 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
2 changes: 2 additions & 0 deletions bin/fm-branch-prompt.sh
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,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 @@ -25,7 +25,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 @@ -50,6 +53,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 @@ -230,10 +239,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
49 changes: 46 additions & 3 deletions bin/fm-wake-drain.sh
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ RECOVERY_ACK_REQUIRED=false
RECOVERY_ACK_MOVED=false
ACK_THROUGH=
ACK_GENERATION=
ACK_REMOVED=0
PRESENTED_MAX=0

# --- per-actor consume (docs/omp-supervision-branch.md "Per-actor acknowledgement") --
# main (FM_SUPERVISION_ACTOR unset or "main", via fm-lease-lib.sh's fm_lease_actor
Expand Down Expand Up @@ -155,6 +157,19 @@ decide_scoped_locked() {
fi
}

# 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 @@ -329,6 +344,13 @@ if [ "$SCOPED" = true ]; then
fi

if [ -n "$ACK_THROUGH" ]; then
if [ "$SCOPED" = true ]; 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
fi
# Row consumption is bound to the monotonic sequence and always happens.
# Only retiring the episode is bound to the generation, so a generation that
# moved on names its own remedy instead of refusing and consuming nothing.
Expand Down Expand Up @@ -361,7 +383,10 @@ if [ -n "$ACK_THROUGH" ]; then
NF < 5 || $2 !~ /^[0-9]+$/ || $2 > cutoff || !($2 in owned) { print }
' "$FM_WAKE_QUEUE" > "$DRAIN_TMP" || 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_snapshot "$RECOVERY_MARKER" || exit 1
RECOVERY_MARKER_TOKEN=$FM_RECOVERY_MARKER_TOKEN
fm_recovery_marker_ack "$RECOVERY_MARKER" "$ACK_GENERATION"
RECOVERY_ACK_STATUS=$?
case "$RECOVERY_ACK_STATUS" in
Expand Down Expand Up @@ -393,9 +418,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
Loading
Loading