mirror of
https://github.com/wassname/pi-plan.git
synced 2026-10-03 12:42:00 +08:00
Replace nested workers with visible supervisor session
Co-Authored-By: Pi <288921227+claudypoo@users.noreply.github.com>
This commit is contained in:
1 parent
4c6a7716b1
commit
d56fc55242
16 files changed
+603
-1647
No files matched your search
@@ -0,0 +1,77 @@
|
||||
import { execFile } from "node:child_process";
|
||||
import { promisify } from "node:util";
|
||||
|
||||
const execFileAsync = promisify(execFile);
|
||||
|
||||
interface LaunchSupervisorInput {
|
||||
cwd: string;
|
||||
sourceSessionFile: string;
|
||||
workerSessionId: string;
|
||||
planPath: string;
|
||||
approvalId: string;
|
||||
extensionPath: string;
|
||||
model: string | null;
|
||||
}
|
||||
|
||||
function shellQuote(value: string): string {
|
||||
return `'${value.replaceAll("'", "'\\''")}'`;
|
||||
}
|
||||
|
||||
function findPaneId(value: unknown): string | null {
|
||||
if (!value || typeof value !== "object") return null;
|
||||
const record = value as Record<string, unknown>;
|
||||
for (const key of ["pane_id", "paneId"]) {
|
||||
if (typeof record[key] === "string") return record[key];
|
||||
}
|
||||
for (const child of Object.values(record)) {
|
||||
const found = findPaneId(child);
|
||||
if (found) return found;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
async function herdr(args: string[]): Promise<unknown> {
|
||||
const bin = process.env.HERDR_BIN_PATH ?? "herdr";
|
||||
const { stdout } = await execFileAsync(bin, args, { encoding: "utf8", timeout: 15_000 });
|
||||
return stdout.trim() ? JSON.parse(stdout) : {};
|
||||
}
|
||||
|
||||
export function supervisorCommand(input: LaunchSupervisorInput): string {
|
||||
const env = [
|
||||
"PI_GOALS_ROLE=supervisor",
|
||||
`PI_GOALS_WORKER_ID=${input.workerSessionId}`,
|
||||
`PI_GOALS_PLAN_PATH=${input.planPath}`,
|
||||
`PI_GOALS_APPROVAL_ID=${input.approvalId}`,
|
||||
`PI_GOALS_OWNER_SESSION_ID=${input.workerSessionId}`,
|
||||
];
|
||||
const args = [
|
||||
"pi",
|
||||
"--no-extensions",
|
||||
"-e", input.extensionPath,
|
||||
"-e", "npm:pi-intercom",
|
||||
"-e", "npm:@wassname2/pi-supervise",
|
||||
"--fork", input.sourceSessionFile,
|
||||
"--name", `goals-supervisor-${input.workerSessionId.slice(0, 8)}`,
|
||||
];
|
||||
if (input.model) args.push("--model", input.model);
|
||||
return ["env", ...env, ...args].map(shellQuote).join(" ");
|
||||
}
|
||||
|
||||
export async function openSupervisorPane(input: LaunchSupervisorInput): Promise<string> {
|
||||
if (process.env.HERDR_ENV !== "1") throw new Error("Ready needs a Herdr session so pi-goals can open the supervisor session.");
|
||||
await herdr(["--version"]);
|
||||
const split = await herdr(["pane", "split", "--current", "--direction", "right", "--cwd", input.cwd, "--no-focus"]);
|
||||
const paneId = findPaneId(split);
|
||||
if (!paneId) throw new Error("Herdr did not return the new supervisor pane ID.");
|
||||
try {
|
||||
await herdr(["pane", "run", paneId, supervisorCommand(input)]);
|
||||
return paneId;
|
||||
} catch (error) {
|
||||
await herdr(["pane", "close", paneId]);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function closeSupervisorPane(paneId: string): Promise<void> {
|
||||
await herdr(["pane", "close", paneId]);
|
||||
}
|
||||
+78
-309
@@ -1,36 +1,29 @@
|
||||
/**
|
||||
* PI: pi-goals owns one versioned plan per session. The main agent is a thin coordinator for a
|
||||
* retained pi-subagents supervisor, which owns one foreground implementation worker at a time and approval.
|
||||
* PI: pi-goals owns one versioned plan per session. After Ready, the main session implements the
|
||||
* plan while a compacted, visible fork supervises it through pi-supervise.
|
||||
*
|
||||
* Each /goals call makes `.pi/plan/<session_id>-vN.md`. The selected version survives resume and
|
||||
* compaction. Old plans stay on disk but inactive. A session with no selected plan has no widget,
|
||||
* supervision, worker, or CompleteGoal sign-off.
|
||||
* supervision, or CompleteGoal sign-off.
|
||||
*
|
||||
* TypeScript reads only goal checkbox lines for the widget. Models read the plan as prose. The
|
||||
* worker edits the project and records evidence. The supervisor inspects it and writes a private
|
||||
* approval checkpoint. pi-subagents owns the supervisor and worker sessions, forks, resume, events,
|
||||
* and Fleet controls.
|
||||
* approval checkpoint.
|
||||
*
|
||||
* -- Pi/Codex
|
||||
*/
|
||||
|
||||
import { execFileSync } from "node:child_process";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { existsSync, mkdirSync, readdirSync, readFileSync, rmSync, writeFileSync } from "node:fs";
|
||||
import { join, resolve } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent";
|
||||
import { Type } from "typebox";
|
||||
import { approvalMatches, approvalPath, goalBlock, hashGoalBlock, readApproval, repositoryState } from "./approval.js";
|
||||
import { closeSupervisorPane, openSupervisorPane } from "./herdr.js";
|
||||
import { completeGoalDescription, completeGoalParamDescription, planDrafting, planningState, resync } from "./prompts.js";
|
||||
import {
|
||||
processWorkState,
|
||||
registerGoalSupervisor,
|
||||
resumeGoalSupervisor,
|
||||
startGoalSupervisor,
|
||||
steerGoalSupervisor,
|
||||
stopGoalSupervisor,
|
||||
subagentWorkState,
|
||||
terminalSteerError,
|
||||
} from "./worker.js";
|
||||
import { isVisibleSupervisor, registerVisibleSupervisor } from "./supervisor-session.js";
|
||||
|
||||
const STATE = "pi-goals-state";
|
||||
const STATUS_KEY = "pi-goals";
|
||||
@@ -41,7 +34,6 @@ const PLAN_DIR = ".pi/plan";
|
||||
const PLAN_SHAPE = `${PLAN_DIR}/<session_id>-vN.md`;
|
||||
// Plan mode blocks edit/write except for its plan file. bash remains available for read-only inspection. -- Pi/Codex
|
||||
const PLAN_MODE_BLOCKED_TOOLS = ["edit", "write"];
|
||||
const AUTO_DEFAULT_INTERVAL_MS = 60 * 60 * 1_000;
|
||||
|
||||
// A checkbox line beginning "goal:", used by the widget and supervisor scheduling.
|
||||
// Everything else reads the file as prose.
|
||||
@@ -93,43 +85,32 @@ export function nextPlanVersion(planNames: string[], sessionId: string): number
|
||||
|
||||
type Phase = "planning" | "working" | null;
|
||||
|
||||
/** Goal workers run in child Pi sessions, so they must not receive the main coordinator's tool gate. */
|
||||
export function isSupervisorProcess(isSubagentChild = process.env.PI_SUBAGENT_CHILD === "1"): boolean {
|
||||
return !isSubagentChild;
|
||||
export function isMainSession(isSubagentChild = process.env.PI_SUBAGENT_CHILD === "1"): boolean {
|
||||
return !isSubagentChild && !isVisibleSupervisor();
|
||||
}
|
||||
|
||||
interface PlanState {
|
||||
phase: Phase;
|
||||
supervisorModel: string | null;
|
||||
workerModel: string | null;
|
||||
workerRunId: string | null;
|
||||
workerPending: boolean;
|
||||
supervisorPaneId: string | null;
|
||||
approvalId: string | null;
|
||||
planVersion: number | null;
|
||||
autoIntervalMs: number | null;
|
||||
}
|
||||
|
||||
export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
if (!isSupervisorProcess()) return;
|
||||
if (isVisibleSupervisor()) {
|
||||
registerVisibleSupervisor(pi);
|
||||
return;
|
||||
}
|
||||
if (!isMainSession()) return;
|
||||
let state: PlanState = {
|
||||
phase: null,
|
||||
supervisorModel: null,
|
||||
workerModel: null,
|
||||
workerRunId: null,
|
||||
workerPending: false,
|
||||
supervisorPaneId: null,
|
||||
approvalId: null,
|
||||
planVersion: null,
|
||||
autoIntervalMs: null,
|
||||
};
|
||||
let planningContextPending = false;
|
||||
let autoTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
let supervisorWakePending = false;
|
||||
let workerRegistration: { dispose(): void } | null = null;
|
||||
let workerRegistrationError: string | null = null;
|
||||
let unsubscribeWorkerCompletion: (() => void) | null = null;
|
||||
let workerLaunchPending = false;
|
||||
const workerCompletionsDuringLaunch = new Set<string>();
|
||||
// Set on session start and after compaction; the next supervisor call receives the whole plan.
|
||||
let resyncReason: string | null = "New session.";
|
||||
|
||||
const planRel = (ctx: ExtensionContext) => (state.planVersion === null ? PLAN_SHAPE : `${PLAN_DIR}/${ctx.sessionManager.getSessionId()}-v${state.planVersion}.md`);
|
||||
@@ -149,23 +130,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
pi.appendEntry<PlanState>(STATE, state);
|
||||
}
|
||||
|
||||
function setupWorker(ctx: ExtensionContext): void {
|
||||
workerRegistration?.dispose();
|
||||
workerRegistration = null;
|
||||
workerRegistrationError = null;
|
||||
try {
|
||||
workerRegistration = registerGoalSupervisor(pi.events, state.supervisorModel);
|
||||
} catch (error) {
|
||||
workerRegistrationError = error instanceof Error ? error.message : String(error);
|
||||
if (state.phase === "working") ctx.ui.notify(`Goal supervisor unavailable: ${workerRegistrationError}`, "warning");
|
||||
}
|
||||
}
|
||||
|
||||
function rememberWorkerRun(runId: string): void {
|
||||
state = { ...state, workerRunId: runId, workerPending: !workerCompletionsDuringLaunch.delete(runId) };
|
||||
persist();
|
||||
}
|
||||
|
||||
function beginReview(ctx: ExtensionContext): void {
|
||||
for (const goal of scanGoals(readPlan(ctx))) {
|
||||
rmSync(approvalPath(ctx.cwd, ctx.sessionManager.getSessionId(), goal.subject), { force: true });
|
||||
@@ -174,12 +138,35 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
persist();
|
||||
}
|
||||
|
||||
async function stopWorker(): Promise<boolean> {
|
||||
clearAutoTimer();
|
||||
if (!state.workerRunId) return !state.workerPending;
|
||||
function repositoryRoot(cwd: string): string {
|
||||
return execFileSync("git", ["rev-parse", "--show-toplevel"], { cwd, encoding: "utf8" }).trim();
|
||||
}
|
||||
|
||||
async function startSupervisor(ctx: ExtensionContext): Promise<void> {
|
||||
repositoryRoot(ctx.cwd);
|
||||
const sourceSessionFile = ctx.sessionManager.getSessionFile();
|
||||
if (!sourceSessionFile) throw new Error("The current session is not persisted, so it cannot be forked.");
|
||||
beginReview(ctx);
|
||||
state = {
|
||||
...state,
|
||||
supervisorPaneId: await openSupervisorPane({
|
||||
cwd: ctx.cwd,
|
||||
sourceSessionFile,
|
||||
workerSessionId: ctx.sessionManager.getSessionId(),
|
||||
planPath: planPath(ctx),
|
||||
approvalId: state.approvalId!,
|
||||
extensionPath: fileURLToPath(import.meta.url),
|
||||
model: state.supervisorModel,
|
||||
}),
|
||||
};
|
||||
persist();
|
||||
}
|
||||
|
||||
async function stopSupervisor(): Promise<boolean> {
|
||||
if (!state.supervisorPaneId) return true;
|
||||
try {
|
||||
await stopGoalSupervisor(pi.events, state.workerRunId);
|
||||
state = { ...state, workerPending: false };
|
||||
await closeSupervisorPane(state.supervisorPaneId);
|
||||
state = { ...state, supervisorPaneId: null };
|
||||
persist();
|
||||
return true;
|
||||
} catch {
|
||||
@@ -187,93 +174,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
}
|
||||
}
|
||||
|
||||
async function startOrResumeWorker(ctx: ExtensionContext, task: string, compactPlanning: boolean, signal?: AbortSignal): Promise<string> {
|
||||
if (workerLaunchPending) throw new Error("A goal-worker launch is already in progress.");
|
||||
if (!workerRegistration) setupWorker(ctx);
|
||||
if (!workerRegistration) throw new Error(`Goal worker unavailable: ${workerRegistrationError ?? "pi-subagents is not ready"}.`);
|
||||
workerLaunchPending = true;
|
||||
workerCompletionsDuringLaunch.clear();
|
||||
try {
|
||||
const runId = state.workerRunId
|
||||
? await resumeGoalSupervisor(pi.events, state.workerRunId, task, signal)
|
||||
: await startGoalSupervisor(pi.events, ctx.cwd, task, compactPlanning, state.workerModel, signal);
|
||||
rememberWorkerRun(runId);
|
||||
return runId;
|
||||
} finally {
|
||||
workerLaunchPending = false;
|
||||
workerCompletionsDuringLaunch.clear();
|
||||
}
|
||||
}
|
||||
|
||||
function watchWorkerCompletion(): void {
|
||||
unsubscribeWorkerCompletion?.();
|
||||
unsubscribeWorkerCompletion = pi.events.on("subagent:async-complete", (raw) => {
|
||||
if (!raw || typeof raw !== "object") return;
|
||||
const runId = (raw as { runId?: string }).runId;
|
||||
if (!runId) return;
|
||||
if (runId !== state.workerRunId) {
|
||||
if (workerLaunchPending) workerCompletionsDuringLaunch.add(runId);
|
||||
return;
|
||||
}
|
||||
state = { ...state, workerPending: false };
|
||||
persist();
|
||||
});
|
||||
}
|
||||
|
||||
function clearAutoTimer(): void {
|
||||
if (autoTimer !== null) clearTimeout(autoTimer);
|
||||
autoTimer = null;
|
||||
}
|
||||
|
||||
function activeGoals(ctx: ExtensionContext): boolean {
|
||||
return scanGoals(readPlan(ctx)).some((goal) => goal.status === "active" || goal.status === "open");
|
||||
}
|
||||
|
||||
function supervisorTask(ctx: ExtensionContext, instruction: string): string {
|
||||
const plan = readPlan(ctx);
|
||||
const checkpoints = scanGoals(plan)
|
||||
.filter((goal) => goal.status === "active" || goal.status === "open")
|
||||
.map((goal) => `- ${JSON.stringify(goal.subject)}: ${approvalPath(ctx.cwd, ctx.sessionManager.getSessionId(), goal.subject)}`)
|
||||
.join("\n");
|
||||
return `${instruction}\n\nYou are the retained goal-supervisor. Here is the complete current plan; inspect its exact goal blocks and cited evidence before directing or approving work.\nPlan path: ${planPath(ctx)}\nApproval ID: ${state.approvalId}\nNested worker model: ${state.workerModel ?? "pi-subagents default"}\nPass the exact approval ID to ApproveGoal. Keep checkpoint paths and the approval ID from the nested worker.\nPrivate approval checkpoints, one per current goal:\n${checkpoints || "(no open goals)"}\n\n${plan}`;
|
||||
}
|
||||
|
||||
async function directSupervisor(ctx: ExtensionContext, instruction: string, signal?: AbortSignal, compactPlanning = false): Promise<string> {
|
||||
beginReview(ctx);
|
||||
const task = supervisorTask(ctx, instruction);
|
||||
if (state.workerPending && state.workerRunId) {
|
||||
try {
|
||||
await steerGoalSupervisor(pi.events, state.workerRunId, task, signal);
|
||||
return state.workerRunId;
|
||||
} catch (error) {
|
||||
if (!terminalSteerError(error)) throw error;
|
||||
state = { ...state, workerPending: false };
|
||||
persist();
|
||||
}
|
||||
}
|
||||
return startOrResumeWorker(ctx, task, compactPlanning, signal);
|
||||
}
|
||||
|
||||
function wakeSupervisor(ctx: ExtensionContext, reason: string): void {
|
||||
if (supervisorWakePending || state.phase !== "working" || !activeGoals(ctx)) return;
|
||||
supervisorWakePending = true;
|
||||
void directSupervisor(ctx, `${reason}\nReview the current goal and either continue, redirect, or approve it through ApproveGoal.`)
|
||||
.catch((error) => ctx.ui.notify(`Goal supervisor check failed: ${error instanceof Error ? error.message : String(error)}`, "warning"))
|
||||
.finally(() => {
|
||||
supervisorWakePending = false;
|
||||
});
|
||||
}
|
||||
|
||||
function scheduleSupervisorCheck(ctx: ExtensionContext): void {
|
||||
if (autoTimer !== null || state.phase !== "working" || state.autoIntervalMs === null || !activeGoals(ctx)) return;
|
||||
autoTimer = setTimeout(() => {
|
||||
autoTimer = null;
|
||||
scheduleSupervisorCheck(ctx);
|
||||
wakeSupervisor(ctx, `The ${state.autoIntervalMs! / 60_000}-minute supervisor check is due.`);
|
||||
}, state.autoIntervalMs);
|
||||
autoTimer.unref();
|
||||
}
|
||||
|
||||
function updateWidget(ctx: ExtensionContext): void {
|
||||
if (state.phase === "planning") {
|
||||
ctx.ui.setStatus(STATUS_KEY, ctx.ui.theme.fg("warning", "planning"));
|
||||
@@ -288,9 +188,8 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
}
|
||||
const done = goals.filter((g) => g.status === "done").length;
|
||||
const liveGoals = goals.filter((g) => g.status === "active" || g.status === "open");
|
||||
const stateLabel = liveGoals.length > 0 ? " · supervising…" : " · complete";
|
||||
const auto = liveGoals.length > 0 && state.autoIntervalMs !== null ? ` · supervise ${state.autoIntervalMs / 60_000}m` : "";
|
||||
ctx.ui.setStatus(STATUS_KEY, ctx.ui.theme.fg("accent", `◷ ${done}/${goals.length} goals${stateLabel}${auto}`));
|
||||
const stateLabel = liveGoals.length > 0 ? " · supervised" : " · complete";
|
||||
ctx.ui.setStatus(STATUS_KEY, ctx.ui.theme.fg("accent", `◷ ${done}/${goals.length} goals${stateLabel}`));
|
||||
const mark: Record<GoalStatus, string> = { done: "✔", active: "▸", open: "◻", cancelled: "✗" };
|
||||
// Only live goals get lines so finished work never pushes current work off screen. The active
|
||||
// goal also shows its open subtasks: this file is the task list, so the widget is the task list.
|
||||
@@ -307,7 +206,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
// --- /goals: enter plan mode or configure supervision -- Pi/Codex -----------------------------
|
||||
|
||||
pi.registerCommand("goals", {
|
||||
description: `Plan goals, then supervise a retained worker. /goals <objective> | clear | auto [minutes|off] | model <supervisor> | worker-model <worker>`,
|
||||
description: `Plan goals, then open a visible supervisor session. /goals <objective> | clear | model <supervisor>`,
|
||||
handler: async (args, ctx) => {
|
||||
const arg = args.trim();
|
||||
if (arg === "clear") {
|
||||
@@ -316,62 +215,36 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
return;
|
||||
}
|
||||
const currentPlan = planRel(ctx);
|
||||
if (!(await stopWorker())) {
|
||||
ctx.ui.notify("Could not stop the retained supervisor; the plan remains connected.", "warning");
|
||||
if (!(await stopSupervisor())) {
|
||||
ctx.ui.notify("Could not close the visible supervisor; the plan remains connected.", "warning");
|
||||
return;
|
||||
}
|
||||
state = { ...state, phase: null, workerRunId: null, workerPending: false, approvalId: null, planVersion: null, autoIntervalMs: null };
|
||||
state = { ...state, phase: null, supervisorPaneId: null, approvalId: null, planVersion: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
ctx.ui.notify(`Disconnected from ${currentPlan}; the file remains on disk.`, "info");
|
||||
return;
|
||||
}
|
||||
if (arg === "auto" || arg.startsWith("auto ")) {
|
||||
const value = arg.slice("auto".length).trim();
|
||||
if (value === "off") {
|
||||
clearAutoTimer();
|
||||
state = { ...state, autoIntervalMs: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
ctx.ui.notify("Hourly goal supervision disabled.", "info");
|
||||
if (arg === "model" || arg.startsWith("model ")) {
|
||||
if (state.phase === "working") {
|
||||
ctx.ui.notify("Run /goals clear before changing the active supervisor model.", "warning");
|
||||
return;
|
||||
}
|
||||
if (state.phase !== "working") {
|
||||
ctx.ui.notify("Approve a plan with Ready before enabling supervision.", "warning");
|
||||
if (!(await stopSupervisor())) {
|
||||
ctx.ui.notify("Could not close the visible supervisor; its model was not changed.", "warning");
|
||||
return;
|
||||
}
|
||||
const minutes = value ? Number(value) : AUTO_DEFAULT_INTERVAL_MS / 60_000;
|
||||
if (!Number.isInteger(minutes) || minutes < 1) {
|
||||
ctx.ui.notify("Use /goals auto [whole minutes], or /goals auto off.", "warning");
|
||||
return;
|
||||
}
|
||||
clearAutoTimer();
|
||||
state = { ...state, autoIntervalMs: minutes * 60_000 };
|
||||
const ref = arg.slice("model".length).trim();
|
||||
state = { ...state, supervisorModel: ref || null, supervisorPaneId: null, approvalId: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
scheduleSupervisorCheck(ctx);
|
||||
ctx.ui.notify(`Goal supervision will check every ${minutes}m.`, "info");
|
||||
ctx.ui.notify(`Goal-supervisor model ${ref ? `set to ${ref}` : "reset to the current Pi default"}.`, "info");
|
||||
return;
|
||||
}
|
||||
if (arg === "model" || arg.startsWith("model ") || arg === "worker-model" || arg.startsWith("worker-model ")) {
|
||||
if (!(await stopWorker())) {
|
||||
ctx.ui.notify("Could not stop the retained supervisor; models were not changed.", "warning");
|
||||
return;
|
||||
}
|
||||
const worker = arg === "worker-model" || arg.startsWith("worker-model ");
|
||||
const command = worker ? "worker-model" : "model";
|
||||
const ref = arg.slice(command.length).trim();
|
||||
state = { ...state, [worker ? "workerModel" : "supervisorModel"]: ref || null, workerRunId: null, workerPending: false, approvalId: null };
|
||||
persist();
|
||||
setupWorker(ctx);
|
||||
ctx.ui.notify(`${worker ? "Implementation-worker" : "Goal-supervisor"} model ${ref ? `set to ${ref}` : "reset to pi-subagents default"}.`, "info");
|
||||
if (!(await stopSupervisor())) {
|
||||
ctx.ui.notify("Could not close the visible supervisor; no new plan was started.", "warning");
|
||||
return;
|
||||
}
|
||||
if (!(await stopWorker())) {
|
||||
ctx.ui.notify("Could not stop the retained supervisor; no new plan was started.", "warning");
|
||||
return;
|
||||
}
|
||||
state = { ...state, phase: "planning", workerRunId: null, workerPending: false, approvalId: null, planVersion: nextVersion(ctx), autoIntervalMs: null };
|
||||
state = { ...state, phase: "planning", supervisorPaneId: null, approvalId: null, planVersion: nextVersion(ctx) };
|
||||
planningContextPending = true;
|
||||
resyncReason = null;
|
||||
writePlan(ctx, "");
|
||||
@@ -399,10 +272,9 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
|
||||
// The phase snapshot enters context only when planning starts or context was lost.
|
||||
pi.on("before_agent_start", async (_event, ctx) => {
|
||||
supervisorWakePending = false;
|
||||
if (state.phase === "working") {
|
||||
return {
|
||||
systemPrompt: `${ctx.getSystemPrompt()}\n\nYou are the thin human-facing coordinator for ${planRel(ctx)}. The retained goal-supervisor owns nested-worker control and acceptance. Keep the human intent stable, inspect progress with read-only tools, and direct the supervisor through GuideGoalWorker. Built-in edit/write and write-like shell commands are blocked. CompleteGoal is a mechanical sign-off only: it fails closed unless the supervisor's private approval checkpoint still matches the exact goal, plan block, committed HEAD/tree, and clean worktree. This is not a filesystem sandbox: allowed verification scripts and other custom tools can still mutate. Do not approve implementation by prose alone. -- Pi/Codex`,
|
||||
systemPrompt: `${ctx.getSystemPrompt()}\n\nYou are the implementation worker for ${planRel(ctx)}. Keep the full conversation and do the work directly. A stronger read-only supervisor watches this session through pi-supervise and can steer you. Commit clean evidence before asking for sign-off. Stop when a goal appears complete so the supervisor can inspect a settled worker view. Call CompleteGoal only after the supervisor says it recorded approval. -- Pi/Codex`,
|
||||
};
|
||||
}
|
||||
if (!planningContextPending) return;
|
||||
@@ -434,12 +306,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
});
|
||||
|
||||
pi.on("tool_call", async (event, ctx) => {
|
||||
if (state.phase === "working" && event.toolName === "subagent" && !["list", "status"].includes(String((event.input as { action?: unknown }).action))) {
|
||||
return { block: true, reason: "The main coordinator may list or inspect subagents; pi-goals owns supervisor lifecycle and delegation." };
|
||||
}
|
||||
if (state.phase === "working" && event.toolName === "subagent_supervisor") {
|
||||
return { block: true, reason: "Direct the retained supervisor through GuideGoalWorker." };
|
||||
}
|
||||
if (state.phase === "planning") {
|
||||
if (PLAN_MODE_BLOCKED_TOOLS.includes(event.toolName)) {
|
||||
const target = (event.input as { path?: string }).path;
|
||||
@@ -451,14 +317,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (state.phase === "working") {
|
||||
if (PLAN_MODE_BLOCKED_TOOLS.includes(event.toolName)) {
|
||||
return { block: true, reason: "Working supervision is read-only: direct implementation and evidence writes to GuideGoalWorker. CompleteGoal is the explicit sign-off control." };
|
||||
}
|
||||
if (event.toolName === "bash" && !isSupervisorReadOnlyCommand(String((event.input as { command?: string }).command))) {
|
||||
return { block: true, reason: "Working supervision allows inspection and standard verification commands only. Direct file changes belong to GuideGoalWorker; this is not a full sandbox for custom tools or allowed scripts." };
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// A compaction loses context, so restore either the planning snapshot or the working plan once.
|
||||
@@ -469,10 +327,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
|
||||
// PI: Print after Pi settles. agent_end is still streaming, so its message queues behind the menu.
|
||||
pi.on("agent_settled", async (_event, ctx) => {
|
||||
if (state.phase === "working") {
|
||||
scheduleSupervisorCheck(ctx);
|
||||
return;
|
||||
}
|
||||
if (state.phase !== "planning" || !ctx.hasUI) return;
|
||||
let printed = "";
|
||||
while (true) {
|
||||
@@ -485,7 +339,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
printed = plan;
|
||||
pi.sendMessage({ customType: "plan", content: plan, display: true });
|
||||
}
|
||||
const choice = await ctx.ui.select(`Plan drafted in ${planRel(ctx)}.`, ["Ready", "Ready (compact)", "Refine", "Edit", "Cancel"]);
|
||||
const choice = await ctx.ui.select(`Plan drafted in ${planRel(ctx)}.`, ["Ready", "Refine", "Edit", "Cancel"]);
|
||||
if (choice === "Refine") {
|
||||
const notes = await ctx.ui.editor("What should change about the plan?", "");
|
||||
if (!notes?.trim()) continue;
|
||||
@@ -500,46 +354,28 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
continue;
|
||||
}
|
||||
if (choice === "Cancel") {
|
||||
if (!(await stopWorker())) {
|
||||
ctx.ui.notify("Could not stop the retained supervisor; the plan was not discarded.", "warning");
|
||||
return;
|
||||
}
|
||||
rmSync(planPath(ctx), { force: true });
|
||||
state = { ...state, phase: null, workerRunId: null, workerPending: false, approvalId: null, planVersion: null, autoIntervalMs: null };
|
||||
state = { ...state, phase: null, supervisorPaneId: null, approvalId: null, planVersion: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
ctx.ui.notify("Plan discarded.", "info");
|
||||
return;
|
||||
}
|
||||
if (choice !== "Ready" && choice !== "Ready (compact)") return;
|
||||
const startWorking = async (): Promise<boolean> => {
|
||||
state = { ...state, phase: "working", autoIntervalMs: AUTO_DEFAULT_INTERVAL_MS };
|
||||
if (choice !== "Ready") return;
|
||||
try {
|
||||
await startSupervisor(ctx);
|
||||
state = { ...state, phase: "working" };
|
||||
resyncReason = "The plan was approved.";
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
try {
|
||||
await directSupervisor(ctx, "Start by launching the foreground implementation worker. Then supervise the current plan.", undefined, choice === "Ready (compact)");
|
||||
scheduleSupervisorCheck(ctx);
|
||||
return true;
|
||||
} catch (error) {
|
||||
ctx.ui.notify(`Goal supervisor could not start: ${error instanceof Error ? error.message : String(error)}`, "warning");
|
||||
state = { ...state, phase: "planning", autoIntervalMs: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
return false;
|
||||
}
|
||||
};
|
||||
const started = await startWorking();
|
||||
if (!started || choice === "Ready") return;
|
||||
ctx.compact({
|
||||
onComplete: () => {
|
||||
resyncReason = "The main coordinator was compacted after the retained supervisor started.";
|
||||
ctx.ui.notify("Main-session compaction completed; the retained supervisor and worker kept their contexts.", "info");
|
||||
},
|
||||
onError: (error) => {
|
||||
ctx.ui.notify(`Main-session compaction failed; the retained supervisor continues: ${error.message}`, "warning");
|
||||
},
|
||||
});
|
||||
ctx.ui.notify(`Visible supervisor opened in Herdr pane ${state.supervisorPaneId}.`, "info");
|
||||
pi.sendUserMessage("The plan is approved. Begin implementation as the worker.");
|
||||
} catch (error) {
|
||||
ctx.ui.notify(`Goal supervisor could not start: ${error instanceof Error ? error.message : String(error)}`, "warning");
|
||||
state = { ...state, phase: "planning", supervisorPaneId: null, approvalId: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
}
|
||||
return;
|
||||
}
|
||||
});
|
||||
@@ -552,65 +388,13 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
state = {
|
||||
phase: last?.data?.phase ?? null,
|
||||
supervisorModel: last?.data?.supervisorModel ?? null,
|
||||
workerModel: last?.data?.workerModel ?? null,
|
||||
workerRunId: last?.data?.workerRunId ?? null,
|
||||
workerPending: last?.data?.workerPending ?? false,
|
||||
supervisorPaneId: last?.data?.supervisorPaneId ?? null,
|
||||
approvalId: last?.data?.approvalId ?? null,
|
||||
planVersion: last?.data?.planVersion ?? null,
|
||||
autoIntervalMs: last?.data?.autoIntervalMs ?? null,
|
||||
};
|
||||
watchWorkerCompletion();
|
||||
setupWorker(ctx);
|
||||
if (state.workerPending && await subagentWorkState(pi.events) === "idle") {
|
||||
state = { ...state, workerPending: false };
|
||||
persist();
|
||||
}
|
||||
planningContextPending = state.phase === "planning";
|
||||
resyncReason = state.phase === "working" ? "New session." : null;
|
||||
updateWidget(ctx);
|
||||
scheduleSupervisorCheck(ctx);
|
||||
});
|
||||
|
||||
pi.on("session_shutdown", async () => {
|
||||
await stopWorker();
|
||||
workerRegistration?.dispose();
|
||||
workerRegistration = null;
|
||||
unsubscribeWorkerCompletion?.();
|
||||
unsubscribeWorkerCompletion = null;
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "CheckGoalWork",
|
||||
label: "Check goal work",
|
||||
description: "Check whether pi-subagents or pi-processes still has active work before deciding that the goal worker stopped.",
|
||||
parameters: Type.Object({}),
|
||||
async execute(_id, _params, _signal, _onUpdate, _ctx) {
|
||||
try {
|
||||
const [subagents, processes] = await Promise.all([subagentWorkState(pi.events), Promise.resolve(processWorkState(pi.events))]);
|
||||
const unknown = subagents === "unknown" || processes === "unknown";
|
||||
return result(`subagents=${subagents}; processes=${processes}`, unknown);
|
||||
} catch (error) {
|
||||
return result(`Goal work status failed: ${error instanceof Error ? error.message : String(error)}`, true);
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "GuideGoalWorker",
|
||||
label: "Guide goal supervisor",
|
||||
description: "Send one concrete instruction to the retained goal-supervisor. A live supervisor is steered; a completed supervisor is resumed with its saved context and the full current plan.",
|
||||
parameters: Type.Object({
|
||||
instruction: Type.String({ description: "The next research or implementation action, with the evidence that should distinguish success from failure." }),
|
||||
}),
|
||||
async execute(_id, params, signal, _onUpdate, ctx) {
|
||||
if (state.phase !== "working") return result("Approve a plan with Ready before directing the goal supervisor.", true);
|
||||
try {
|
||||
const runId = await directSupervisor(ctx, params.instruction, signal);
|
||||
return result(`Instruction delivered to goal supervisor ${runId}.`);
|
||||
} catch (error) {
|
||||
return result(`Goal-supervisor guidance failed: ${error instanceof Error ? error.message : String(error)}`, true);
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
@@ -622,9 +406,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
}),
|
||||
async execute(_id, params, _signal, _onUpdate, ctx) {
|
||||
if (state.phase !== "working") return result("Planning is not approved. Choose Ready before signing off a goal.", true);
|
||||
if (state.workerPending) return result("Goal sign-off blocked while the retained supervisor is pending.", true);
|
||||
const [subagents, processes] = await Promise.all([subagentWorkState(pi.events), Promise.resolve(processWorkState(pi.events))]);
|
||||
if (subagents !== "idle" || processes !== "idle") return result(`Goal sign-off blocked: subagents=${subagents}; processes=${processes}.`, true);
|
||||
if (!state.approvalId) return result("Goal sign-off blocked: no current supervisor review.", true);
|
||||
const plan = readPlan(ctx);
|
||||
if (!plan.trim()) return result(`No plan file at ${planRel(ctx)}. Run /goals to draft one.`, true);
|
||||
@@ -676,18 +457,6 @@ function isPlanningReadOnlyCommand(command: string): boolean {
|
||||
});
|
||||
}
|
||||
|
||||
/** This blocks direct writes, not side effects hidden in allowed project scripts or custom tools. -- PI[gpt-5.6-sol] */
|
||||
export function isSupervisorReadOnlyCommand(command: string): boolean {
|
||||
if (/[|><`$\n\r]/.test(command)) return false;
|
||||
const safeArgs = "(?:\\s+[A-Za-z0-9_./:=,'\"@+%-]+)*";
|
||||
const inspection = new RegExp(`^(?:cd|pwd|ls|rg|grep|find|head|tail|wc|stat|test)${safeArgs}$`);
|
||||
const git = new RegExp(`^git\\s+(?:status|log|diff|show|branch|ls-files|grep|check-ignore)${safeArgs}$`);
|
||||
const verification = new RegExp(`^(?:npm\\s+test|npm\\s+run\\s+(?:test|typecheck|lint)|npx\\s+tsc\\s+--noEmit)${safeArgs}$`);
|
||||
return command.split(/&&|;/).every((raw) => {
|
||||
const part = raw.trim();
|
||||
return !mutatingReadCommand(part) && (inspection.test(part) || git.test(part) || verification.test(part));
|
||||
});
|
||||
}
|
||||
|
||||
/** Local time, not UTC: agents freehand-stamp their manual ## Log lines from the local clock they
|
||||
* see, so a UTC tool stamp made the trail read as two different afternoons (dogfood finding). */
|
||||
|
||||
+16
-18
@@ -2,9 +2,8 @@
|
||||
* pi-goals v2 — all model-facing text, in flow order.
|
||||
*
|
||||
* Design: the plan file is for LLMs and the human, not for TypeScript. No parser and no schema;
|
||||
* the skeleton below is a convention the drafting prompt teaches, the worker maintains with its
|
||||
* normal Edit tool, and the retained goal-supervisor reads natively. The harness provides format
|
||||
* guidance, one full-plan resync after context loss, and retained pi-subagents supervisor and worker sessions.
|
||||
* the skeleton below is a convention the drafting prompt teaches. The main session implements it,
|
||||
* while a visible forked Pi session supervises through pi-supervise.
|
||||
*
|
||||
* THE FOLD: everything above "## Log" is the short current-goal section. Everything below it
|
||||
* (Log, Learnings, Appendix) is durable memory: unlimited, read on demand, and sent in full at
|
||||
@@ -13,8 +12,8 @@
|
||||
* Flow:
|
||||
* SETUP (plan mode) 1. planDrafting — draft goals into the plan file (read-only), sent once
|
||||
* EXEC, after compact 2. resync — the WHOLE file back, once
|
||||
* SIGN-OFF, agent-side 3. completeGoal* — the one blessed tool's description
|
||||
* SUPERVISION worker.ts - retained supervisor and foreground worker
|
||||
* SIGN-OFF, worker-side 3. completeGoal* — the one blessed tool's description
|
||||
* SUPERVISION supervisor-session.ts — visible read-only supervisor
|
||||
*
|
||||
* The goal's test is the DISCRIMINATOR: the concrete observation that tells real success from the
|
||||
* named subtle failure mode. Evidence is empty at planning and filled at sign-off.
|
||||
@@ -69,7 +68,7 @@ Style: Make it easy for a busy and forgetfull user to review. Use ASD-STE100 Sim
|
||||
the same word for the same thing, and define a new terms at first use. Use redundant context for skim readers e.g. "our output - the cells, CV tag" is easy to read and reminds context. This covers the context
|
||||
paragraph and the appendix too, not just the checklist. No all-caps headers and no bold spam. Just write less, add your voice less, persuade less, and burden the reader less.
|
||||
|
||||
Write the plan file in roughly this shape -- the file is read directly by the human and the retained supervisor, so clarity beats conformance; small deviations are fine):
|
||||
Write the plan file in roughly this shape -- the file is read directly by the human and the visible supervisor, so clarity beats conformance; small deviations are fine):
|
||||
|
||||
# <short plan title>
|
||||
|
||||
@@ -89,7 +88,7 @@ Write the plan file in roughly this shape -- the file is read directly by the hu
|
||||
- subtle failure mode: <a way this could look done but isn't>
|
||||
- discriminator: <the concrete observation that tells real success from that failure>
|
||||
- verify: <optional shell command that exits 0 only when the discriminator passes; omit if not
|
||||
testable. The worker runs it and saves its output; the retained supervisor reads the evidence>
|
||||
testable. The worker runs it and saves its output; the visible supervisor reads the evidence>
|
||||
- tasks:
|
||||
1. [ ] <subtask>
|
||||
- evidence: (empty until sign-off)
|
||||
@@ -119,7 +118,7 @@ Conventions:
|
||||
none of the failure modes could fake. Ruling out failures is necessary, not sufficient.
|
||||
- Make the discriminator a concrete, checkable observation about a real artifact (a file, a test
|
||||
result, a committed diff, a metric), never about the plan file's own checkbox.
|
||||
- evidence stays empty at planning; the worker fills it and the retained supervisor checks it.
|
||||
- evidence stays empty at planning; the worker fills it and the visible supervisor checks it.
|
||||
Cite durable artifacts a future reader can open: committed files, test names, git diffs. .pi/ is
|
||||
usually gitignored, so files there prove things only at supervisor review time, not in history.
|
||||
- User-visible result: restate the original deliverable, not the proposed implementation. Every goal
|
||||
@@ -159,11 +158,10 @@ Ready.`;
|
||||
export function resync(plan: string, planRel: string, why: string): string {
|
||||
return `\
|
||||
<system-reminder>
|
||||
${why} This is the whole plan file (${planRel}), appendix included. You are the main coordinator.
|
||||
Keep the high-level goal and human intent stable; direct the retained goal-supervisor through
|
||||
GuideGoalWorker rather than doing implementation. The supervisor directs and approves the nested
|
||||
worker. The human's latest message outranks the plan: if it changes scope, direct the supervisor to
|
||||
have the worker amend the plan rather than preserving an obsolete decision.
|
||||
${why} This is the whole plan file (${planRel}), appendix included. You are the implementation worker.
|
||||
Keep the high-level goal and human intent stable and do the work directly. A visible read-only Pi
|
||||
session supervises you through pi-supervise. The human's latest message outranks the plan: if it
|
||||
changes scope, amend the plan rather than preserving an obsolete decision.
|
||||
|
||||
${plan}
|
||||
</system-reminder>`;
|
||||
@@ -180,11 +178,11 @@ export const completeGoalDescription =
|
||||
"output, rerun it or write that you couldn't -- an honest gap beats a plausible fabrication. If " +
|
||||
"the goal names a verify: command, direct the worker to run it and save its output to a file cited " +
|
||||
"in the evidence. The supervisor may run an allowed read-only verification command, but must not " +
|
||||
"create the evidence file itself. The retained goal-supervisor must reject a claimed pass with no " +
|
||||
"saved output. The read must show success POSITIVELY happened, not just that failures were avoided. " +
|
||||
"The supervisor records an approval checkpoint only after it inspected the current plan, repository, " +
|
||||
"evidence, and verify output with no active nested worker and a clean committed worktree. Then the main " +
|
||||
"coordinator calls this tool with the exact goal text. This tool independently checks that checkpoint " +
|
||||
"create the evidence file itself. The visible supervisor must reject a claimed pass with no saved " +
|
||||
"output. The read must show success POSITIVELY happened, not just that failures were avoided. The " +
|
||||
"supervisor records an approval checkpoint only after it inspected the current plan, repository, " +
|
||||
"evidence, verify output, and a stopped worker view with no active work. Then the worker calls this " +
|
||||
"tool with the exact goal text. This tool independently checks that checkpoint " +
|
||||
"against the exact current goal block, HEAD/tree, and clean worktree before it appends the sign-off to " +
|
||||
"## Log and ticks the goal [x]. If any check differs, it fails closed and requires a fresh supervisor review.";
|
||||
|
||||
|
||||
@@ -1,179 +0,0 @@
|
||||
import { readFileSync } from "node:fs";
|
||||
import { resolve } from "node:path";
|
||||
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
|
||||
import { Type } from "typebox";
|
||||
import { goalBlock, hashGoalBlock, repositoryState, writeApproval } from "./approval.js";
|
||||
import { isSupervisorReadOnlyCommand } from "./index.js";
|
||||
import { GOAL_WORKER_AGENT, processWorkState } from "./worker.js";
|
||||
|
||||
const COMPACTED_STATE = "pi-goals-supervisor-compacted";
|
||||
|
||||
function result(text: string, isError = false) {
|
||||
return { content: [{ type: "text" as const, text }], details: {}, isError };
|
||||
}
|
||||
|
||||
interface GoalBindings {
|
||||
compactPlanning?: boolean;
|
||||
workerModel?: string | null;
|
||||
}
|
||||
|
||||
function goalBindings(): GoalBindings {
|
||||
const raw = process.env.PI_SUBAGENT_EXTENSION_BINDINGS;
|
||||
if (!raw) return {};
|
||||
const binding = (JSON.parse(raw) as { "pi-goals/1"?: GoalBindings })["pi-goals/1"] ?? {};
|
||||
if (binding.workerModel !== undefined && binding.workerModel !== null && typeof binding.workerModel !== "string") throw new Error("pi-goals workerModel binding must be a string or null.");
|
||||
return binding;
|
||||
}
|
||||
|
||||
function messageLaunchesWorker(ctx: { sessionManager: { getBranch(): unknown[] } }): boolean {
|
||||
const entry = [...ctx.sessionManager.getBranch()].reverse().find((candidate) => {
|
||||
const value = candidate as { type?: unknown; message?: { role?: unknown } };
|
||||
return value.type === "message" && value.message?.role === "assistant";
|
||||
}) as { message?: { content?: unknown } } | undefined;
|
||||
if (!Array.isArray(entry?.message?.content)) return false;
|
||||
return entry.message.content.some((part) => {
|
||||
const value = part as { type?: unknown; name?: unknown; arguments?: Record<string, unknown> };
|
||||
return value.type === "toolCall" && value.name === "subagent" && value.arguments?.agent === GOAL_WORKER_AGENT;
|
||||
});
|
||||
}
|
||||
|
||||
export default function goalSupervisorRuntime(pi: ExtensionAPI): void {
|
||||
let compacting = false;
|
||||
let compactionDone = Promise.resolve();
|
||||
let currentTurn = -1;
|
||||
let completedWorkerTurn: number | null = null;
|
||||
let workerModel: string | null = null;
|
||||
const activeWorkerCalls = new Set<string>();
|
||||
|
||||
pi.on("session_before_compact", async (event) => {
|
||||
if (!compacting) return;
|
||||
const branchEntries = event.branchEntries as Array<{ id?: string; type?: string; message?: { role?: string } }>;
|
||||
const latestMessage = [...branchEntries].reverse().find((entry) => entry.type === "message" && ["user", "assistant"].includes(entry.message?.role ?? ""));
|
||||
return {
|
||||
compaction: {
|
||||
summary: "Planning is complete. The latest retained goal-supervisor task contains the current plan and approval paths; use it as the source of truth. -- PI[gpt-5.6-sol]",
|
||||
firstKeptEntryId: latestMessage?.id ?? event.preparation.firstKeptEntryId,
|
||||
tokensBefore: event.preparation.tokensBefore,
|
||||
details: { source: "pi-goals-plan-handoff" },
|
||||
},
|
||||
};
|
||||
});
|
||||
|
||||
pi.on("session_start", async (_event, ctx) => {
|
||||
const entries = ctx.sessionManager.getEntries();
|
||||
const bindings = goalBindings();
|
||||
workerModel = bindings.workerModel ?? null;
|
||||
if (bindings.compactPlanning !== true) return;
|
||||
if (entries.some((entry: { type?: string; customType?: string }) => entry.type === "custom" && entry.customType === COMPACTED_STATE)) return;
|
||||
compacting = true;
|
||||
compactionDone = new Promise<void>((resolvePromise, reject) => {
|
||||
ctx.compact({
|
||||
onComplete: () => {
|
||||
compacting = false;
|
||||
pi.appendEntry(COMPACTED_STATE, { version: 1 });
|
||||
resolvePromise();
|
||||
},
|
||||
onError: (error) => {
|
||||
compacting = false;
|
||||
reject(error);
|
||||
},
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
pi.on("before_agent_start", async () => {
|
||||
await compactionDone;
|
||||
});
|
||||
|
||||
pi.on("turn_start", async (event) => {
|
||||
activeWorkerCalls.clear();
|
||||
currentTurn = event.turnIndex;
|
||||
});
|
||||
|
||||
pi.on("tool_call", async (event) => {
|
||||
if (event.toolName === "edit" || event.toolName === "write") {
|
||||
return { block: true, reason: "Goal supervision is read-only. Direct project changes to the nested goal-worker." };
|
||||
}
|
||||
if (event.toolName === "bash" && !isSupervisorReadOnlyCommand(String((event.input as { command?: string }).command))) {
|
||||
return { block: true, reason: "Goal supervision allows inspection and standard verification commands only." };
|
||||
}
|
||||
if (event.toolName !== "subagent") return;
|
||||
const input = event.input as Record<string, unknown>;
|
||||
const allowedKeys = new Set(["agent", "task", "async", "context", ...(workerModel ? ["model"] : [])]);
|
||||
const unexpectedKeys = Object.keys(input).filter((key) => !allowedKeys.has(key));
|
||||
const validWorker = input.agent === GOAL_WORKER_AGENT
|
||||
&& typeof input.task === "string"
|
||||
&& input.task.trim().length > 0
|
||||
&& input.async === false
|
||||
&& input.context === "fork"
|
||||
&& (workerModel ? input.model === workerModel : input.model === undefined)
|
||||
&& unexpectedKeys.length === 0;
|
||||
if (!validWorker) {
|
||||
const model = workerModel ? `, model:${JSON.stringify(workerModel)}` : "";
|
||||
return { block: true, reason: `Launch only ${GOAL_WORKER_AGENT} with task, async:false, context:"fork"${model}, and no other fields.` };
|
||||
}
|
||||
if (activeWorkerCalls.size > 0) return { block: true, reason: "A foreground goal-worker is already running." };
|
||||
activeWorkerCalls.add(event.toolCallId);
|
||||
completedWorkerTurn = null;
|
||||
});
|
||||
|
||||
pi.on("tool_result", async (event) => {
|
||||
if (!activeWorkerCalls.delete(event.toolCallId)) return;
|
||||
if (!event.isError) completedWorkerTurn = currentTurn;
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "ApproveGoal",
|
||||
label: "Approve goal",
|
||||
executionMode: "sequential",
|
||||
description: "Record approval after inspecting the plan, repository, evidence, and saved verification output. Active or unknown work blocks approval.",
|
||||
parameters: Type.Object({
|
||||
approvalId: Type.String({ minLength: 1, description: "Exact approval ID from the latest main-coordinator direction." }),
|
||||
goal: Type.String({ description: "Exact current goal text from the approved plan." }),
|
||||
planPath: Type.String({ description: "Absolute path to the current plan file." }),
|
||||
checkpointPath: Type.String({ description: "Exact private approval-record path supplied by the main coordinator." }),
|
||||
inspectedPlan: Type.Literal(true),
|
||||
inspectedRepository: Type.Literal(true),
|
||||
inspectedEvidence: Type.Literal(true),
|
||||
inspectedVerifyOutput: Type.Literal(true),
|
||||
}),
|
||||
async execute(_id, params, _signal, _onUpdate, ctx) {
|
||||
if (messageLaunchesWorker(ctx) || activeWorkerCalls.size > 0 || completedWorkerTurn === null || completedWorkerTurn >= currentTurn) {
|
||||
return result("Cannot approve in a worker-launch message or before reviewing a finished worker on a later turn.", true);
|
||||
}
|
||||
const processes = processWorkState(pi.events);
|
||||
if (processes !== "idle") return result(`Cannot approve: processes=${processes}.`, true);
|
||||
const planPath = resolve(params.planPath);
|
||||
let plan: string;
|
||||
let repository: ReturnType<typeof repositoryState>;
|
||||
try {
|
||||
plan = readFileSync(planPath, "utf8");
|
||||
repository = repositoryState(ctx.cwd);
|
||||
} catch (error) {
|
||||
return result(`Cannot inspect approval inputs: ${error instanceof Error ? error.message : String(error)}`, true);
|
||||
}
|
||||
if (!repository.cleanWorktree) return result("Cannot approve with a dirty worktree. Commit the worker changes first.", true);
|
||||
const block = goalBlock(plan, params.goal);
|
||||
if (!block) return result(`Cannot approve: no unique open goal matches "${params.goal}".`, true);
|
||||
const path = resolve(params.checkpointPath);
|
||||
const approvalRoot = resolve(ctx.cwd, ".pi", "pi-goals", "approvals");
|
||||
if (!path.startsWith(`${approvalRoot}/`)) return result("Approval checkpoint must stay in private .pi/pi-goals/approvals state.", true);
|
||||
writeApproval(path, {
|
||||
version: 2,
|
||||
verdict: "accept",
|
||||
approvalId: params.approvalId,
|
||||
goal: params.goal,
|
||||
planPath,
|
||||
goalBlockHash: hashGoalBlock(block),
|
||||
repoRoot: repository.repoRoot,
|
||||
head: repository.head,
|
||||
tree: repository.tree,
|
||||
cleanWorktree: true,
|
||||
inspected: { plan: true, repository: true, evidence: true, verifyOutput: true },
|
||||
supervisor: { sessionId: ctx.sessionManager.getSessionId(), runId: process.env.PI_SUBAGENT_RUN_ID ?? null },
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
return result(`Approval recorded at ${path} for "${params.goal}".`);
|
||||
},
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,154 @@
|
||||
import { readFileSync } from "node:fs";
|
||||
import { resolve } from "node:path";
|
||||
import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent";
|
||||
import { Type } from "typebox";
|
||||
import { approvalPath, goalBlock, hashGoalBlock, repositoryState, writeApproval } from "./approval.js";
|
||||
|
||||
const BOOTSTRAPPED = "pi-goals-visible-supervisor-v1";
|
||||
const COMPACT_AT_TOKENS = 100_000;
|
||||
|
||||
interface SupervisorConfig {
|
||||
workerSessionId: string;
|
||||
ownerSessionId: string;
|
||||
planPath: string;
|
||||
approvalId: string;
|
||||
}
|
||||
|
||||
function result(text: string, isError = false) {
|
||||
return { content: [{ type: "text" as const, text }], details: {}, isError };
|
||||
}
|
||||
|
||||
function requiredEnv(name: string): string {
|
||||
const value = process.env[name]?.trim();
|
||||
if (!value) throw new Error(`${name} is required in a pi-goals supervisor session.`);
|
||||
return value;
|
||||
}
|
||||
|
||||
function config(): SupervisorConfig {
|
||||
return {
|
||||
workerSessionId: requiredEnv("PI_GOALS_WORKER_ID"),
|
||||
ownerSessionId: requiredEnv("PI_GOALS_OWNER_SESSION_ID"),
|
||||
planPath: resolve(requiredEnv("PI_GOALS_PLAN_PATH")),
|
||||
approvalId: requiredEnv("PI_GOALS_APPROVAL_ID"),
|
||||
};
|
||||
}
|
||||
|
||||
function latestWorkerView(ctx: ExtensionContext): string | null {
|
||||
for (const entry of [...ctx.sessionManager.getBranch()].reverse()) {
|
||||
const message = (entry as { type?: string; message?: { role?: string; content?: unknown[] } }).message;
|
||||
if ((entry as { type?: string }).type !== "message" || message?.role !== "user" || !Array.isArray(message.content)) continue;
|
||||
for (const part of message.content) {
|
||||
const text = (part as { type?: string; text?: string }).type === "text" ? (part as { text?: string }).text : undefined;
|
||||
if (text?.startsWith("The worker ")) return text;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function supervisorPrompt(settings: SupervisorConfig): string {
|
||||
return `You are the visible pi-goals supervisor for ${settings.planPath}. You are a stronger, read-only reviewer. The other Pi session is the implementation worker and keeps the full conversation. You keep the high-level intent from the compacted planning conversation, the complete plan, and pi-supervise worker views.
|
||||
|
||||
Use pi-supervise to inspect and steer the worker. Give one concrete instruction when work is incomplete. Do not edit files. For each open goal, inspect its exact plan block, repository state, cited evidence, and saved verify output. When its discriminator is positively satisfied and the worker view says no work is active, call ApproveGoal. Then call steer and tell the worker to call CompleteGoal with the exact goal text. Do not call done until every plan goal is [x]. -- PI[gpt-5.6-sol]`;
|
||||
}
|
||||
|
||||
export function isVisibleSupervisor(): boolean {
|
||||
return process.env.PI_GOALS_ROLE === "supervisor";
|
||||
}
|
||||
|
||||
export function registerVisibleSupervisor(pi: ExtensionAPI): void {
|
||||
const settings = config();
|
||||
let compacting = false;
|
||||
|
||||
pi.on("before_agent_start", async (_event, ctx) => ({
|
||||
systemPrompt: `${ctx.getSystemPrompt()}\n\n${supervisorPrompt(settings)}`,
|
||||
}));
|
||||
|
||||
pi.on("session_start", async (_event, ctx) => {
|
||||
const entries = ctx.sessionManager.getEntries();
|
||||
if (entries.some((entry: { type?: string; customType?: string }) => entry.type === "custom" && entry.customType === BOOTSTRAPPED)) return;
|
||||
if (!pi.getCommands().some((command) => command.name === "supervise" && command.source === "extension")) {
|
||||
ctx.ui.notify("pi-goals supervisor needs the @wassname2/pi-supervise extension.", "error");
|
||||
return;
|
||||
}
|
||||
compacting = true;
|
||||
ctx.compact({
|
||||
customInstructions: `Preserve the user's decisions, preferences, and high-level objective from planning. Preserve unresolved risks and the plan path ${settings.planPath}. Remove implementation chatter. This summary is for a read-only supervisor that will judge and steer another Pi session.`,
|
||||
onComplete: () => {
|
||||
compacting = false;
|
||||
pi.appendEntry(BOOTSTRAPPED, { version: 1, workerSessionId: settings.workerSessionId, planPath: settings.planPath });
|
||||
const sendCommand = pi.sendUserMessage as (content: string, options: { expandPromptTemplates: boolean }) => void;
|
||||
sendCommand(`/supervise @${settings.workerSessionId} ${settings.planPath}`, { expandPromptTemplates: true });
|
||||
},
|
||||
onError: (error) => {
|
||||
compacting = false;
|
||||
ctx.ui.notify(`Supervisor compaction failed: ${error.message}`, "error");
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
pi.on("agent_settled", async (_event, ctx) => {
|
||||
if (compacting || (ctx.getContextUsage()?.tokens ?? 0) < COMPACT_AT_TOKENS) return;
|
||||
compacting = true;
|
||||
ctx.compact({
|
||||
customInstructions: `Keep the user's high-level intent, current plan state, unresolved risks, approval decisions, and the supervisor's own concise findings. Remove old worker views and implementation detail.`,
|
||||
onComplete: () => {
|
||||
compacting = false;
|
||||
ctx.ui.notify("Supervisor context compacted at 100k tokens.", "info");
|
||||
},
|
||||
onError: (error) => {
|
||||
compacting = false;
|
||||
ctx.ui.notify(`Supervisor compaction failed: ${error.message}`, "error");
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "ApproveGoal",
|
||||
label: "Approve goal",
|
||||
executionMode: "sequential",
|
||||
description: "Record approval after inspecting the current goal, repository, evidence, saved verify output, and a stopped worker view with no active work.",
|
||||
parameters: Type.Object({
|
||||
goal: Type.String({ description: "Exact text after goal: in the plan." }),
|
||||
inspectedPlan: Type.Literal(true),
|
||||
inspectedRepository: Type.Literal(true),
|
||||
inspectedEvidence: Type.Literal(true),
|
||||
inspectedVerifyOutput: Type.Literal(true),
|
||||
}),
|
||||
async execute(_id, params, _signal, _onUpdate, ctx) {
|
||||
const view = latestWorkerView(ctx);
|
||||
if (!view?.startsWith("The worker stopped.")) return result("Cannot approve without a current stopped-worker view.", true);
|
||||
const pendingTool = view.match(/^tool calls with no result: (?!none$)(.+)$/m);
|
||||
const pendingChild = view.match(/^child pi processes still running: (?!none$)(.+)$/m);
|
||||
if (pendingTool || pendingChild) return result(`Cannot approve while work is active: ${(pendingTool ?? pendingChild)![1]}`, true);
|
||||
let plan: string;
|
||||
let repository: ReturnType<typeof repositoryState>;
|
||||
try {
|
||||
plan = readFileSync(settings.planPath, "utf8");
|
||||
repository = repositoryState(ctx.cwd);
|
||||
} catch (error) {
|
||||
return result(`Cannot inspect approval inputs: ${error instanceof Error ? error.message : String(error)}`, true);
|
||||
}
|
||||
if (!repository.cleanWorktree) return result("Cannot approve with a dirty worktree. Commit the worker changes first.", true);
|
||||
const block = goalBlock(plan, params.goal);
|
||||
if (!block) return result(`Cannot approve: no unique open goal matches "${params.goal}".`, true);
|
||||
if (/evidence:\s*\(empty until sign-off\)/i.test(block)) return result("Cannot approve while the goal evidence is empty.", true);
|
||||
const path = approvalPath(ctx.cwd, settings.ownerSessionId, params.goal);
|
||||
writeApproval(path, {
|
||||
version: 2,
|
||||
verdict: "accept",
|
||||
approvalId: settings.approvalId,
|
||||
goal: params.goal,
|
||||
planPath: settings.planPath,
|
||||
goalBlockHash: hashGoalBlock(block),
|
||||
repoRoot: repository.repoRoot,
|
||||
head: repository.head,
|
||||
tree: repository.tree,
|
||||
cleanWorktree: true,
|
||||
inspected: { plan: true, repository: true, evidence: true, verifyOutput: true },
|
||||
supervisor: { sessionId: ctx.sessionManager.getSessionId(), runId: null },
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
return result(`Approval recorded for "${params.goal}". Now steer the worker to call CompleteGoal.`);
|
||||
},
|
||||
});
|
||||
}
|
||||
-186
@@ -1,186 +0,0 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { fileURLToPath } from "node:url";
|
||||
|
||||
const REGISTER_EVENT = "pi-subagents:runtime-agent-register:v1";
|
||||
const RPC_REQUEST_EVENT = "subagents:rpc:v1:request";
|
||||
const RPC_REPLY_PREFIX = "subagents:rpc:v1:reply:";
|
||||
const RPC_VERSION = 1;
|
||||
const RPC_TIMEOUT_MS = 15_000;
|
||||
export const SUPERVISOR_AGENT = "goal-supervisor";
|
||||
export const GOAL_WORKER_AGENT = "pi-goals-worker-v1";
|
||||
|
||||
interface EventBus {
|
||||
on(event: string, handler: (data: unknown) => void): () => void;
|
||||
emit(event: string, data: unknown): void;
|
||||
}
|
||||
|
||||
interface Registration {
|
||||
dispose(): void;
|
||||
}
|
||||
|
||||
interface RpcData {
|
||||
text: string;
|
||||
details?: Record<string, unknown>;
|
||||
asyncSnapshot?: AsyncSnapshot;
|
||||
}
|
||||
|
||||
interface AsyncNode {
|
||||
id: string;
|
||||
state: string;
|
||||
children?: AsyncNode[];
|
||||
}
|
||||
|
||||
interface AsyncSnapshot {
|
||||
kind: string;
|
||||
version: number;
|
||||
omitted: { runs: number; children: number; byteLimitExceeded: boolean };
|
||||
runs: AsyncNode[];
|
||||
}
|
||||
|
||||
export type WorkState = "active" | "idle" | "unknown";
|
||||
|
||||
export const supervisorSystemPrompt = `You are the retained goal supervisor. The main Pi session only coordinates with the human.
|
||||
Your forked planning history may be compacted before your first turn. Launch ${GOAL_WORKER_AGENT} in the foreground with exactly
|
||||
agent, task, async:false, context:"fork", and, when named in the current direction, that worker model. Wait for its result; do not use
|
||||
bg_wait or worker run IDs. Read the current plan, repository, cited evidence, and saved verification output yourself after the
|
||||
worker finishes. Do not edit project files. Use read/search and standard verification commands only. The worker must commit its
|
||||
changes before approval. If the evidence needs a correction, launch a new foreground ${GOAL_WORKER_AGENT} with one concrete task
|
||||
and wait for it. On a later turn, when HEAD is committed, the worktree is clean, and the evidence proves the discriminator, call
|
||||
ApproveGoal with the current approval ID. Only ApproveGoal creates acceptance. -- Pi/Codex`;
|
||||
|
||||
function registerRuntimeAgent(events: EventBus, name: string, definition: Record<string, unknown>): Registration {
|
||||
const request: Record<string, unknown> = { version: 1, name, definition };
|
||||
events.emit(REGISTER_EVENT, request);
|
||||
const result = request.result as { ok?: boolean; registration?: Registration; error?: Error } | undefined;
|
||||
if (!result) throw new Error("pi-subagents is not installed or not ready.");
|
||||
if (!result.ok || !result.registration) throw result.error ?? new Error(`pi-subagents rejected the ${name} agent.`);
|
||||
return result.registration;
|
||||
}
|
||||
|
||||
export function registerGoalSupervisor(events: EventBus, model: string | null): Registration {
|
||||
const supervisorRuntime = fileURLToPath(new URL("./supervisor-runtime.ts", import.meta.url));
|
||||
return registerRuntimeAgent(events, SUPERVISOR_AGENT, {
|
||||
description: "Read-only supervisor that owns a foreground implementation worker.",
|
||||
systemPrompt: supervisorSystemPrompt,
|
||||
tools: ["read", "grep", "find", "ls", "bash", "subagent", "ApproveGoal"],
|
||||
allowNestedSubagents: true,
|
||||
subagentOnlyExtensions: [supervisorRuntime],
|
||||
...(model ? { model } : {}),
|
||||
systemPromptMode: "replace",
|
||||
thinking: "low",
|
||||
inheritProjectContext: false,
|
||||
inheritGlobalContext: false,
|
||||
inheritSkills: false,
|
||||
defaultContext: "fork",
|
||||
defaultAsync: true,
|
||||
defaultProgress: true,
|
||||
});
|
||||
}
|
||||
|
||||
async function rpc(events: EventBus, method: "spawn" | "resume" | "steer" | "status" | "stop", params: Record<string, unknown>, signal?: AbortSignal): Promise<RpcData> {
|
||||
if (signal?.aborted) throw new Error("Goal-worker request aborted.");
|
||||
const requestId = randomUUID();
|
||||
return new Promise((resolve, reject) => {
|
||||
let timer: ReturnType<typeof setTimeout>;
|
||||
const replyEvent = `${RPC_REPLY_PREFIX}${requestId}`;
|
||||
const cleanup = () => {
|
||||
clearTimeout(timer);
|
||||
unsubscribe();
|
||||
signal?.removeEventListener("abort", onAbort);
|
||||
};
|
||||
const onAbort = () => {
|
||||
cleanup();
|
||||
reject(new Error("Goal-worker request aborted."));
|
||||
};
|
||||
const unsubscribe = events.on(replyEvent, (raw) => {
|
||||
const reply = raw as { success?: boolean; data?: RpcData; error?: { message?: string } };
|
||||
cleanup();
|
||||
if (!reply.success || !reply.data) reject(new Error(reply.error?.message ?? `pi-subagents ${method} failed.`));
|
||||
else resolve(reply.data);
|
||||
});
|
||||
timer = setTimeout(() => {
|
||||
cleanup();
|
||||
reject(new Error(`pi-subagents ${method} did not reply within ${RPC_TIMEOUT_MS / 1000}s.`));
|
||||
}, RPC_TIMEOUT_MS);
|
||||
timer.unref();
|
||||
signal?.addEventListener("abort", onAbort, { once: true });
|
||||
events.emit(RPC_REQUEST_EVENT, { version: RPC_VERSION, requestId, method, params, source: { extension: "pi-goals" } });
|
||||
});
|
||||
}
|
||||
|
||||
function asyncRunId(data: RpcData): string {
|
||||
const runId = data.details?.asyncId ?? data.details?.runId;
|
||||
if (typeof runId !== "string" || !runId) throw new Error("pi-subagents returned no async run ID.");
|
||||
return runId;
|
||||
}
|
||||
|
||||
export async function startGoalSupervisor(events: EventBus, cwd: string, task: string, compactPlanning: boolean, workerModel: string | null, signal?: AbortSignal): Promise<string> {
|
||||
const data = await rpc(events, "spawn", {
|
||||
agent: SUPERVISOR_AGENT,
|
||||
task,
|
||||
cwd,
|
||||
context: "fork",
|
||||
async: true,
|
||||
mission: false,
|
||||
extensionBindings: { "pi-goals/1": { compactPlanning, workerModel } },
|
||||
}, signal);
|
||||
return asyncRunId(data);
|
||||
}
|
||||
|
||||
export async function resumeGoalSupervisor(events: EventBus, runId: string, task: string, signal?: AbortSignal): Promise<string> {
|
||||
return asyncRunId(await rpc(events, "resume", { id: runId, message: task }, signal));
|
||||
}
|
||||
|
||||
export async function steerGoalSupervisor(events: EventBus, runId: string, task: string, signal?: AbortSignal): Promise<void> {
|
||||
await rpc(events, "steer", { id: runId, message: task, mode: "steer" }, signal);
|
||||
}
|
||||
|
||||
export async function stopGoalSupervisor(events: EventBus, runId: string): Promise<void> {
|
||||
try {
|
||||
await rpc(events, "stop", { id: runId });
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
if (!/not found|already completed|\bis (?:complete|completed|failed|partial|paused|stopped|rejected)\b/i.test(message)) throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export function terminalSteerError(error: unknown): boolean {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
return /not found|already completed|not running|\bis (?:complete|completed|failed|partial|paused|stopped|rejected)\b/i.test(message);
|
||||
}
|
||||
|
||||
function activeNode(node: AsyncNode): boolean {
|
||||
return node.state === "queued" || node.state === "running" || node.state === "stopping" || Boolean(node.children?.some(activeNode));
|
||||
}
|
||||
|
||||
function validSnapshot(snapshot: AsyncSnapshot | undefined): snapshot is AsyncSnapshot {
|
||||
return snapshot?.kind === "pi-subagents.async-status-snapshot" && snapshot.version === 1 && snapshot.omitted.runs === 0 && snapshot.omitted.children === 0 && !snapshot.omitted.byteLimitExceeded;
|
||||
}
|
||||
|
||||
async function asyncSnapshot(events: EventBus): Promise<AsyncSnapshot | undefined> {
|
||||
return (await rpc(events, "status", {})).asyncSnapshot;
|
||||
}
|
||||
|
||||
export async function subagentWorkState(events: EventBus): Promise<WorkState> {
|
||||
const snapshot = await asyncSnapshot(events);
|
||||
if (!validSnapshot(snapshot)) return "unknown";
|
||||
return snapshot.runs.some(activeNode) ? "active" : "idle";
|
||||
}
|
||||
|
||||
export interface ProcessInfo {
|
||||
status: string;
|
||||
}
|
||||
|
||||
export function processWorkState(events: EventBus): WorkState {
|
||||
let replied = false;
|
||||
let processes: ProcessInfo[] = [];
|
||||
events.emit("processes:request:list", {
|
||||
reply(value: ProcessInfo[]) {
|
||||
replied = true;
|
||||
processes = value;
|
||||
},
|
||||
});
|
||||
if (!replied || !Array.isArray(processes)) return "unknown";
|
||||
const terminal = new Set(["finished", "failed", "exited", "killed"]);
|
||||
return processes.every((process) => terminal.has(process.status)) ? "idle" : "active";
|
||||
}
|
||||
Reference in new issue
Block a user