Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
54b179f
fix(pi): route needs-decision wakes and mixed batches wholly to main
kunchenguid Sep 5, 2026
883db7e
test(pi): cover distinct-file mixed batches and heartbeat independence
kunchenguid Sep 5, 2026
9ae58ee
no-mistakes(document): Clarify needs-decision and heartbeat routing
kunchenguid Sep 5, 2026
c42f1c1
no-mistakes(ci): Fixed captain-held stale reminders so they bypass su…
kunchenguid Sep 5, 2026
e2d00c0
no-mistakes(ci): Fixed CI lint by narrowly suppressing false-positive…
kunchenguid Sep 5, 2026
5ab4628
no-mistakes(review): Route stale open decisions directly to main
kunchenguid Sep 5, 2026
a0b8985
no-mistakes(review): Honor configured verbs in stale decision routing
kunchenguid Sep 5, 2026
b6def7f
no-mistakes(review): Route second-mate escalations and configured dec…
kunchenguid Sep 5, 2026
bcc94dc
no-mistakes(review): Ignore trailing whitespace after captain holds
kunchenguid Sep 5, 2026
6fe1083
no-mistakes(review): Cache stale decision classification per status file
kunchenguid Sep 5, 2026
51a61e2
no-mistakes(review): Document unread decision precedence for later ta…
kunchenguid Sep 5, 2026
b393817
no-mistakes(review): Cache unchanged stale decisions across scope scans
kunchenguid Sep 5, 2026
9382199
no-mistakes(review): Resolve decision aliases and reject symlinked st…
kunchenguid Sep 5, 2026
30eacf0
no-mistakes(review): Route surfaced captain-held signals directly to …
kunchenguid Sep 5, 2026
7df2931
no-mistakes(document): Document decision-owned main routing
kunchenguid Sep 5, 2026
f3304d0
no-mistakes(ci): Fixed captain-held spans to remain actionable while …
kunchenguid Sep 5, 2026
e24a937
no-mistakes(ci): Fixed the CI regression: captain-held transfers now …
kunchenguid Sep 5, 2026
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
21 changes: 20 additions & 1 deletion .pi/extensions/fm-primary-pi-watch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -606,7 +606,26 @@ export default function (pi: ExtensionAPI) {
// also let a check-kind trigger itself slip past main's delivery.
const isCheckTrigger = /^check:/.test(message);
const scope = scopeForUnreadWake(state, heartbeat);
const eligible = !isCheckTrigger && scope.eligible;
// A signal close containing a needs-decision status file, or a stale close
// for a captain-held task, gets the identical main-only treatment as a
// check-kind trigger. The cross-reference deliberately includes every
// unread decision row: until that row is read, a later signal or stale
// trigger for the same task stays on main. Other tasks and heartbeat
// handling remain independent.
const triggerKeys = /^signal:/.test(message)
? message
.slice("signal:".length)
.split(/\s+/)
.filter(Boolean)
.map((path) => path.split("/").pop() ?? path)
: /^stale:/.test(message)
? [message.slice("stale:".length).trim().split(/\s+/, 1)[0]].filter(Boolean)
: [];
const taskIdentity = (key: string): string =>
scope.taskByWakeKey[key] ?? scope.taskByWakeKey[key.replace(/^fm-/, "")] ?? key;
const needsDecisionTasks = new Set(scope.needsDecisionKeys.map(taskIdentity));
const isNeedsDecisionTrigger = triggerKeys.some((key) => needsDecisionTasks.has(taskIdentity(key)));
const eligible = !isCheckTrigger && !isNeedsDecisionTrigger && scope.eligible;
const offer = createBranchDispatchOffer(message, scope.projects, heartbeat, eligible);
pi.events?.emit?.(FM_BRANCH_DISPATCH_EVENT, offer);
return offer.accepted ? offer.settlement : null;
Expand Down
147 changes: 146 additions & 1 deletion .pi/extensions/lib/fm-branch-dispatch.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { readdirSync, readFileSync } from "node:fs";
import { lstatSync, readdirSync, readFileSync } from "node:fs";
import { runCommandAsync } from "./fm-async-exec.ts";

// Shared wake-dispatch handshake between the Pi watcher extension (the
Expand Down Expand Up @@ -54,6 +54,16 @@ export interface UnreadWakeScope {
* either mode.
*/
corrupted: boolean;
/**
* The exact "key" field of every decision-owned signal or stale row this
* scan excluded. Signal rows are marked by bin/fm-watch.sh; stale rows are
* decision-owned when their task has an open needs-decision or its current
* declaration is captain-held. fm-primary-pi-watch.ts cross-references these
* keys against the current trigger so its entire coalesced batch is forced
* to main.
*/
needsDecisionKeys: string[];
taskByWakeKey: Record<string, string>;
}

const EMPTY_SCOPE: UnreadWakeScope = {
Expand All @@ -63,6 +73,8 @@ const EMPTY_SCOPE: UnreadWakeScope = {
eligibleSeqs: [],
eligibleTasks: [],
corrupted: false,
needsDecisionKeys: [],
taskByWakeKey: {},
};
const UNSAFE_SCOPE: UnreadWakeScope = {
status: "unsafe",
Expand All @@ -71,6 +83,8 @@ const UNSAFE_SCOPE: UnreadWakeScope = {
eligibleSeqs: [],
eligibleTasks: [],
corrupted: true,
needsDecisionKeys: [],
taskByWakeKey: {},
};

// scopeForUnreadWake is the single owner of branch-eligibility classification
Expand All @@ -86,6 +100,12 @@ const UNSAFE_SCOPE: UnreadWakeScope = {
// (fm-primary-pi-watch.ts forces every check-kind TRIGGER to main), so nothing
// starves by being left behind.
//
// A signal row whose payload is "needs-decision:"-prefixed, or a stale row
// for a task with an open needs-decision or a current captain-held declaration,
// gets the identical treatment: excluded from eligibleSeqs, never a scan veto,
// and forced to main on its own triggering close (fm-primary-pi-watch.ts's
// offerWakeToBranch). Heartbeat handling remains independent.
//
// That applies to a heartbeat review too, and it is the whole point: a
// heartbeat used to be deferred to main merely because some unrelated check
// row happened to be sitting unread, which put a routine fleet review in the
Expand All @@ -102,6 +122,71 @@ const UNSAFE_SCOPE: UnreadWakeScope = {
// this repo's fm_wake_append could never have produced (an unknown kind, or a
// line that fails the structural tab-field check) also still vetoes the whole
// scan - that is queue corruption, not an everyday mixed queue.
function statusLineVerb(line: string): string {
const beforeColon = line.split(":", 1)[0].split("[", 1)[0].trim();
const words = beforeColon.split(/\s+/);
if (!words.some((word) => word.startsWith("corr="))) return beforeColon;
return words.filter((word, index) => index === 0 || !/^corr=[0-9a-f]{16}$/i.test(word)).join(" ");
}

function decisionKey(line: string): string | null {
const colon = line.indexOf(":");
const beforeColon = colon < 0 ? line : line.slice(0, colon);
const beforeMatch = beforeColon.match(/\[key=([^\]]*)\]/);
const noteMatch = beforeMatch || colon < 0 ? null : line.slice(colon + 1).trimStart().match(/^\[key=([^\]]*)\]/);
const key = (beforeMatch ?? noteMatch)?.[1] ?? "default";
return /^[A-Za-z0-9._-]+$/.test(key) ? key : null;
}

function statusLineNote(line: string): string {
const colon = line.indexOf(":");
if (colon < 0) return line;
const note = line.slice(colon + 1).trimStart();
if (/\[key=[^\]]*\]/.test(line.slice(0, colon))) return note;
const match = note.match(/^\[key=([A-Za-z0-9._-]+)\]/);
return match ? note.slice(match[0].length).trimStart() : note;
}

interface StaleDecisionCacheEntry {
version: string;
config: string;
decisionOwned: boolean;
}

const staleDecisionCache = new Map<string, StaleDecisionCacheEntry>();

function statusFileVersion(path: string): string | null {
try {
const stat = lstatSync(path);
if (stat.isSymbolicLink()) throw new Error("status path is a symbolic link");
return `${stat.dev}:${stat.ino}:${stat.size}:${stat.mtimeMs}:${stat.ctimeMs}`;
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return null;
throw error;
}
}

function hasOpenNeedsDecision(
lines: readonly string[],
resolveVerb: string,
heldVerb: string,
reservedPrefixes: readonly string[],
): boolean {
const open = new Map<string, "needs-decision" | "blocked">();
for (const line of lines) {
const verb = statusLineVerb(line);
if (!["needs-decision", "blocked", resolveVerb, heldVerb].includes(verb)) continue;
const key = decisionKey(line);
if (!key) continue;
const note = statusLineNote(line);
const reservedPrefix = reservedPrefixes.find((prefix) => key.startsWith(prefix));
if (reservedPrefix && !(note.startsWith(reservedPrefix) && note.slice(reservedPrefix.length).includes(":"))) continue;
if (verb === "needs-decision" || verb === "blocked") open.set(key, verb);
else open.delete(key);
}
return [...open.values()].includes("needs-decision");
}

export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWakeScope {
let queue = "";
try {
Expand All @@ -128,6 +213,8 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak
if (project) {
metadata.set(task, project);
taskByKey.set(task, task);
taskByKey.set(`${task}.status`, task);
taskByKey.set(`${task}.turn-ended`, task);
if (window) {
metadata.set(window, project);
taskByKey.set(window, task);
Expand All @@ -140,6 +227,14 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak

const eligibleSeqs: string[] = [];
const eligibleTasks = new Set<string>();
const needsDecisionKeys: string[] = [];
const staleDecisionOwnership = new Map<string, boolean>();
const resolveVerb = process.env.FM_CLASSIFY_RESOLVE_VERB || "resolved";
const heldVerb = process.env.FM_CLASSIFY_CAPTAIN_HELD_VERB || "captain-held";
const reservedPrefixes = (process.env.FM_CLASSIFY_RESERVED_KEY_PREFIXES || "pending-reply-")
.split(/\s+/)
.filter(Boolean);
const decisionConfig = `${resolveVerb}\0${heldVerb}\0${reservedPrefixes.join("\0")}`;
for (const line of rows) {
const fields = line.split("\t");
if (fields.length < 5 || !/^[0-9]+$/.test(fields[1])) return UNSAFE_SCOPE;
Expand All @@ -159,11 +254,59 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak
let project = "";
let task = "";
if (kind === "signal") {
const payload = fields[4] ?? "";
if (/^needs-decision:/.test(payload)) {
// Main-owned exactly like a check-kind row above: a needs-decision
// status append surfaced through the actionable signal path is
// excluded from what the branch may claim without vetoing the scan
// (docs/pi-supervision-branch.md "Autonomy").
needsDecisionKeys.push(key);
continue;
}
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-/, "")) ?? "";
if (task) {
const statusPath = `${state}/${task}.status`;
if (!staleDecisionOwnership.has(statusPath)) {
let version: string | null;
try {
version = statusFileVersion(statusPath);
} catch {
return UNSAFE_SCOPE;
}
let decisionOwned = false;
if (version) {
const cached = staleDecisionCache.get(statusPath);
if (cached?.version === version && cached.config === decisionConfig) {
decisionOwned = cached.decisionOwned;
} else {
let statusLines: string[];
try {
statusLines = readFileSync(statusPath, "utf8").split(/\r?\n/).filter((line) => /\S/.test(line));
if (statusFileVersion(statusPath) !== version) return UNSAFE_SCOPE;
} catch {
return UNSAFE_SCOPE;
}
decisionOwned = hasOpenNeedsDecision(statusLines, resolveVerb, heldVerb, reservedPrefixes) ||
statusLineVerb(statusLines.at(-1) ?? "") === heldVerb;
staleDecisionCache.set(statusPath, { version, config: decisionConfig, decisionOwned });
if (staleDecisionCache.size > 512) {
staleDecisionCache.delete(staleDecisionCache.keys().next().value!);
}
}
} else {
staleDecisionCache.delete(statusPath);
}
staleDecisionOwnership.set(statusPath, decisionOwned);
}
if (staleDecisionOwnership.get(statusPath)) {
needsDecisionKeys.push(key);
continue;
}
}
} else {
// A kind fm_wake_append never emits: structural corruption, not an
// ordinary main-only row.
Expand All @@ -189,6 +332,8 @@ export function scopeForUnreadWake(state: string, heartbeat: boolean): UnreadWak
eligibleSeqs,
eligibleTasks: [...eligibleTasks],
corrupted: false,
needsDecisionKeys,
taskByWakeKey: Object.fromEntries(taskByKey),
};
}

Expand Down
56 changes: 48 additions & 8 deletions bin/fm-classify-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -418,6 +418,14 @@ _fm_decision_key_transition_allowed() { # <key> <note>
return 0
}

_fm_is_pending_reply_escalation() { # <key> <note>
case "$1" in pending-reply-*) ;; *) return 1 ;; esac
case "$2" in
pending-reply-missed:*|pending-reply-delivery-unknown:*|pending-reply-recovery-delivery-failed:*|pending-reply-recovery-delivery-unknown:*) return 0 ;;
*) return 1 ;;
esac
}

_fm_decision_fold_line() { # <open-set> <status-line> <resolve-verb> <held-verb>
local open=$1 line=$2 resolve=$3 held=$4 verb key note
# Blank-line guard. A `case` glob answers "does this line hold any non-space
Expand Down Expand Up @@ -1555,12 +1563,16 @@ window_to_task() {

# Capture the bytes of an append-only status log at or after <start-offset> under
# one size-and-identity snapshot.
# The record form prints `<endpoint>\t<identity>\t<events>` and returns 0 when
# The record form produces `<endpoint>\t<identity>\t<events>` and returns 0 when
# the span has actionable events, joining every such event in source order with
# ` ; ` so callers report the complete captured span before committing it.
# With optional <record-var>, it assigns that record instead of printing it; with
# optional <needs-decision-var>, it also assigns 1 when the span newly surfaces a
# needs-decision, captain-held declaration, or pending-reply escalation, otherwise
# 0. This side-band classification never changes the event text.
# It returns 1 after a successful classification with no actionable event; an
# existing log still prints its committable endpoint and identity, while an absent
# log is the ordinary empty case and prints no record.
# existing log still produces its committable endpoint and identity, while an absent
# log is the ordinary empty case and produces no record.
# It returns 2 with no committable endpoint when an existing status object cannot
# be classified.
# The simpler wrapper prints only the event field, and the predicate discards the
Expand Down Expand Up @@ -1615,9 +1627,9 @@ _fm_status_open_decision_origins() { # <status-file>
printf '%s' "$origins"
}

status_span_first_actionable_record() { # <status-file> <start-offset>
local f=$1 start=${2:-0} size ident cur_ident scratch chunk_file full_file prefix_file
local line verb key origins='' folded=0 rc=1 failed=0 prefix_lines=0 line_number=0 live_line='' events='' _line _key
status_span_first_actionable_record() { # <status-file> <start-offset> [record-var] [needs-decision-var]
local f=$1 start=${2:-0} output_var=${3-} needs_var=${4-} size ident cur_ident scratch chunk_file full_file prefix_file result
local line verb key origins='' folded=0 rc=1 failed=0 prefix_lines=0 line_number=0 live_line='' events='' _line _key _fm_span_needs_decision=0
[ -e "$f" ] || { [ -L "$f" ] && return 2; return 1; }
[ -f "$f" ] && [ -r "$f" ] && [ ! -L "$f" ] || return 2
ident=$(_fm_open_decisions_file_ident "$f") || return 2
Expand All @@ -1626,7 +1638,16 @@ status_span_first_actionable_record() { # <status-file> <start-offset>
case "$size" in ''|*[!0-9]*) return 2 ;; esac
case "$start" in ''|*[!0-9]*) start=0 ;; esac
[ "$start" -le "$size" ] || start=0
[ "$start" -lt "$size" ] || { printf '%s\t%s' "$size" "$ident"; return 1; }
if [ "$start" -ge "$size" ]; then
result="${size}"$'\t'"${ident}"
if [ -n "$output_var" ]; then
printf -v "$output_var" '%s' "$result"
[ -z "$needs_var" ] || printf -v "$needs_var" '%s' 0
else
printf '%s' "$result"
fi
return 1
fi
scratch=$(_fm_status_span_scratch "$f") || return 2
chunk_file="${scratch}.span"; full_file="${scratch}.full"; prefix_file="${scratch}.prefix"
_fm_status_read_span "$f" "$start" "$((size - start))" > "$chunk_file" 2>/dev/null \
Expand All @@ -1638,19 +1659,28 @@ status_span_first_actionable_record() { # <status-file> <start-offset>
while IFS= read -r line || [ -n "$line" ]; do
line_number=$((line_number + 1))
case "$line" in *[![:space:]]*) ;; *) continue ;; esac
if status_is_captain_held "$line"; then
# A transfer closes the status-log decision and remains non-actionable to
# stale classification. The side-band marker lets signal routing surface
# the captain-owned hold without changing that established stale verdict.
_fm_span_needs_decision=1
continue
fi
Comment thread
greptile-apps[bot] marked this conversation as resolved.
status_is_captain_relevant "$line" || continue
verb=$(status_line_verb "$line")
case "$verb" in
needs-decision|blocked)
key=$(_fm_decision_key "$line") || {
[ -n "$events" ] && events="${events} ; "
events="${events}${line}"
[ "$verb" = needs-decision ] && _fm_span_needs_decision=1
rc=0
continue
}
_fm_decision_key_transition_allowed "$key" "$(status_line_note "$line")" || {
[ -n "$events" ] && events="${events} ; "
events="${events}reconciliation-required: ${line}"
[ "$verb" = needs-decision ] && _fm_span_needs_decision=1
rc=0
continue
}
Expand All @@ -1674,6 +1704,10 @@ EOF
[ -n "$live_line" ] && [ "$((prefix_lines + line_number))" -eq "$live_line" ] || continue
[ -n "$events" ] && events="${events} ; "
events="${events}${line}"
if [ "$verb" = needs-decision ] || { [ "$verb" = blocked ] &&
_fm_is_pending_reply_escalation "$key" "$(status_line_note "$line")"; }; then
_fm_span_needs_decision=1
fi
rc=0
;;
*)
Expand All @@ -1685,7 +1719,13 @@ EOF
done < "$chunk_file"
rm -f "$chunk_file" "$full_file" "$prefix_file"
[ "$failed" -eq 0 ] || return 2
if [ "$rc" -eq 0 ]; then printf '%s\t%s\t%s' "$size" "$ident" "$events"; else printf '%s\t%s' "$size" "$ident"; fi
if [ "$rc" -eq 0 ]; then result="${size}"$'\t'"${ident}"$'\t'"${events}"; else result="${size}"$'\t'"${ident}"; fi
if [ -n "$output_var" ]; then
printf -v "$output_var" '%s' "$result"
[ -z "$needs_var" ] || printf -v "$needs_var" '%s' "$_fm_span_needs_decision"
else
printf '%s' "$result"
fi
return "$rc"
}

Expand Down
Loading
Loading