mirror of
https://github.com/wassname/pi-plan.git
synced 2026-10-03 12:42:00 +08:00
Own visible supervision in pi-goals
Co-Authored-By: PI[Kimi K3] <288921227+claudypoo@users.noreply.github.com>
This commit is contained in:
1 parent
6c86405841
commit
19fa8d7a7b
14 files changed
+398
-292
No files matched your search
+15
-8
@@ -1,17 +1,18 @@
|
||||
import { execFile } from "node:child_process";
|
||||
import { promisify } from "node:util";
|
||||
import { supervisorReady } from "./mailbox.js";
|
||||
|
||||
const execFileAsync = promisify(execFile);
|
||||
const STARTUP_TIMEOUT_MS = 5 * 60_000;
|
||||
|
||||
interface LaunchSupervisorInput {
|
||||
cwd: string;
|
||||
sourceSessionFile: string;
|
||||
workerSessionId: string;
|
||||
workerIntercomId: string;
|
||||
planPath: string;
|
||||
approvalId: string;
|
||||
mailboxPath: string;
|
||||
extensionPath: string;
|
||||
superviseExtensionPath: string | null;
|
||||
model: string | null;
|
||||
}
|
||||
|
||||
@@ -49,16 +50,14 @@ export function supervisorCommand(input: LaunchSupervisorInput): string {
|
||||
const env = [
|
||||
"PI_GOALS_ROLE=supervisor",
|
||||
`PI_GOALS_WORKER_ID=${input.workerSessionId}`,
|
||||
`PI_GOALS_WORKER_INTERCOM_ID=${input.workerIntercomId}`,
|
||||
`PI_GOALS_PLAN_PATH=${input.planPath}`,
|
||||
`PI_GOALS_APPROVAL_ID=${input.approvalId}`,
|
||||
`PI_GOALS_OWNER_SESSION_ID=${input.workerSessionId}`,
|
||||
`PI_GOALS_MAILBOX_PATH=${input.mailboxPath}`,
|
||||
];
|
||||
const args = [
|
||||
"pi",
|
||||
"--no-extensions",
|
||||
"-e", "npm:pi-intercom",
|
||||
"-e", process.env.PI_GOALS_SUPERVISE_EXTENSION ?? input.superviseExtensionPath ?? "npm:@wassname2/pi-supervise@0.0.4",
|
||||
"-e", input.extensionPath,
|
||||
"--fork", input.sourceSessionFile,
|
||||
"--name", `goals-supervisor-${input.workerSessionId.slice(0, 8)}`,
|
||||
@@ -67,18 +66,26 @@ export function supervisorCommand(input: LaunchSupervisorInput): string {
|
||||
return `env ${[...env, ...args].map(shellQuote).join(" ")}`;
|
||||
}
|
||||
|
||||
export async function waitForSupervisorReady(mailboxPath: string, timeoutMs = STARTUP_TIMEOUT_MS): Promise<void> {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (!supervisorReady(mailboxPath)) {
|
||||
if (Date.now() >= deadline) throw new Error("The visible supervisor did not become ready.");
|
||||
await new Promise((resolve) => setTimeout(resolve, 100));
|
||||
}
|
||||
}
|
||||
|
||||
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"], false);
|
||||
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.");
|
||||
await herdr(["pane", "run", paneId, supervisorCommand(input)]);
|
||||
try {
|
||||
await herdr(["pane", "run", paneId, supervisorCommand(input)]);
|
||||
await waitForSupervisorReady(input.mailboxPath);
|
||||
return paneId;
|
||||
} catch (error) {
|
||||
await closeSupervisorPane(paneId);
|
||||
throw error;
|
||||
throw new Error(`Supervisor startup incomplete in Herdr pane ${paneId}; inspect that pane. ${error instanceof Error ? error.message : String(error)}`);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+72
-18
@@ -1,6 +1,6 @@
|
||||
/**
|
||||
* 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.
|
||||
* plan while a compacted, visible fork supervises it through a durable mailbox.
|
||||
*
|
||||
* 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,
|
||||
@@ -22,9 +22,10 @@ import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-a
|
||||
import { Type } from "typebox";
|
||||
import { approvalMatches, approvalPath, goalBlock, hashGoalBlock, readApproval, repositoryState } from "./approval.js";
|
||||
import { closeSupervisorPane, openSupervisorPane } from "./herdr.js";
|
||||
import { createMailbox, workerSteersAfter, writeWorkerView } from "./mailbox.js";
|
||||
import { completeGoalDescription, completeGoalParamDescription, planDrafting, planningState, resync } from "./prompts.js";
|
||||
import { SUPERVISOR_STARTUP_TIMEOUT_MS, workerPiSupervise } from "./supervise.js";
|
||||
import { isVisibleSupervisor, registerVisibleSupervisor } from "./supervisor-session.js";
|
||||
import { workerView } from "./worker-view.js";
|
||||
|
||||
const STATE = "pi-goals-state";
|
||||
const STATUS_KEY = "pi-goals";
|
||||
@@ -95,6 +96,8 @@ interface PlanState {
|
||||
supervisorModel: string | null;
|
||||
supervisorPaneId: string | null;
|
||||
approvalId: string | null;
|
||||
mailboxPath: string | null;
|
||||
lastSteer: number;
|
||||
planVersion: number | null;
|
||||
}
|
||||
|
||||
@@ -109,6 +112,8 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
supervisorModel: null,
|
||||
supervisorPaneId: null,
|
||||
approvalId: null,
|
||||
mailboxPath: null,
|
||||
lastSteer: 0,
|
||||
planVersion: null,
|
||||
};
|
||||
let planningContextPending = false;
|
||||
@@ -135,15 +140,12 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
for (const goal of scanGoals(readPlan(ctx))) {
|
||||
rmSync(approvalPath(ctx.cwd, ctx.sessionManager.getSessionId(), goal.subject), { force: true });
|
||||
}
|
||||
state = { ...state, approvalId: randomUUID() };
|
||||
const approvalId = randomUUID();
|
||||
const mailbox = createMailbox(ctx.cwd, ctx.sessionManager.getSessionId(), approvalId, planPath(ctx));
|
||||
state = { ...state, approvalId, mailboxPath: mailbox.path, lastSteer: 0 };
|
||||
persist();
|
||||
}
|
||||
|
||||
function loadedPiSuperviseExtensionPath(): string | null {
|
||||
const tool = pi.getAllTools().find((candidate) => candidate.name === "worker_view") as { sourceInfo?: { path?: unknown } } | undefined;
|
||||
return typeof tool?.sourceInfo?.path === "string" ? tool.sourceInfo.path : null;
|
||||
}
|
||||
|
||||
function repositoryRoot(cwd: string): string {
|
||||
return execFileSync("git", ["rev-parse", "--show-toplevel"], { cwd, encoding: "utf8" }).trim();
|
||||
}
|
||||
@@ -152,7 +154,6 @@ export default function piGoalsExtension(pi: ExtensionAPI): 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.");
|
||||
const worker = await workerPiSupervise(pi);
|
||||
beginReview(ctx);
|
||||
let paneId: string | null = null;
|
||||
try {
|
||||
@@ -160,14 +161,12 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
cwd: ctx.cwd,
|
||||
sourceSessionFile,
|
||||
workerSessionId: ctx.sessionManager.getSessionId(),
|
||||
workerIntercomId: worker.intercomId,
|
||||
planPath: planPath(ctx),
|
||||
approvalId: state.approvalId!,
|
||||
mailboxPath: state.mailboxPath!,
|
||||
extensionPath: fileURLToPath(import.meta.url),
|
||||
superviseExtensionPath: loadedPiSuperviseExtensionPath(),
|
||||
model: state.supervisorModel,
|
||||
});
|
||||
await worker.waitForPair(SUPERVISOR_STARTUP_TIMEOUT_MS);
|
||||
} catch (error) {
|
||||
if (paneId) throw new Error(`Supervisor startup failed in Herdr pane ${paneId}; it remains open for inspection. ${error instanceof Error ? error.message : String(error)}`);
|
||||
throw error;
|
||||
@@ -176,7 +175,43 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
persist();
|
||||
}
|
||||
|
||||
let workerTurns = 0;
|
||||
let viewTimer: ReturnType<typeof setInterval> | undefined;
|
||||
let steerTimer: ReturnType<typeof setInterval> | undefined;
|
||||
|
||||
function mailbox(ctx: ExtensionContext) {
|
||||
if (!state.approvalId || !state.mailboxPath) throw new Error("No active supervisor mailbox.");
|
||||
return createMailbox(ctx.cwd, ctx.sessionManager.getSessionId(), state.approvalId, planPath(ctx));
|
||||
}
|
||||
|
||||
function publishWorkerView(ctx: ExtensionContext, reason: "ready" | "settled" | "turns" | "interval"): void {
|
||||
if (state.phase !== "working") return;
|
||||
writeWorkerView(mailbox(ctx), reason, workerView(ctx.sessionManager.getBranch(), reason));
|
||||
}
|
||||
|
||||
function deliverWorkerSteers(ctx: ExtensionContext): void {
|
||||
if (state.phase !== "working") return;
|
||||
for (const steer of workerSteersAfter(mailbox(ctx).path, state.lastSteer)) {
|
||||
state = { ...state, lastSteer: steer.sequence };
|
||||
persist();
|
||||
pi.sendUserMessage(`[supervisor] ${steer.instruction}`, { deliverAs: "followUp" });
|
||||
}
|
||||
}
|
||||
|
||||
function startWorkerTimers(ctx: ExtensionContext): void {
|
||||
if (!viewTimer) viewTimer = setInterval(() => publishWorkerView(ctx, "interval"), 60 * 60_000);
|
||||
if (!steerTimer) steerTimer = setInterval(() => deliverWorkerSteers(ctx), 1_000);
|
||||
}
|
||||
|
||||
function stopWorkerTimers(): void {
|
||||
if (viewTimer) clearInterval(viewTimer);
|
||||
if (steerTimer) clearInterval(steerTimer);
|
||||
viewTimer = undefined;
|
||||
steerTimer = undefined;
|
||||
}
|
||||
|
||||
async function stopSupervisor(): Promise<boolean> {
|
||||
stopWorkerTimers();
|
||||
if (!state.supervisorPaneId) return true;
|
||||
try {
|
||||
await closeSupervisorPane(state.supervisorPaneId);
|
||||
@@ -233,7 +268,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
ctx.ui.notify("Could not close the visible supervisor; the plan remains connected.", "warning");
|
||||
return;
|
||||
}
|
||||
state = { ...state, phase: null, supervisorPaneId: null, approvalId: null, planVersion: null };
|
||||
state = { ...state, phase: null, supervisorPaneId: null, approvalId: null, mailboxPath: null, lastSteer: 0, planVersion: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
ctx.ui.notify(`Disconnected from ${currentPlan}; the file remains on disk.`, "info");
|
||||
@@ -249,7 +284,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
return;
|
||||
}
|
||||
const ref = arg.slice("model".length).trim();
|
||||
state = { ...state, supervisorModel: ref || null, supervisorPaneId: null, approvalId: null };
|
||||
state = { ...state, supervisorModel: ref || null, supervisorPaneId: null, approvalId: null, mailboxPath: null, lastSteer: 0 };
|
||||
persist();
|
||||
ctx.ui.notify(`Goal-supervisor model ${ref ? `set to ${ref}` : "reset to the current Pi default"}.`, "info");
|
||||
return;
|
||||
@@ -258,7 +293,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
ctx.ui.notify("Could not close the visible supervisor; no new plan was started.", "warning");
|
||||
return;
|
||||
}
|
||||
state = { ...state, phase: "planning", supervisorPaneId: null, approvalId: null, planVersion: nextVersion(ctx) };
|
||||
state = { ...state, phase: "planning", supervisorPaneId: null, approvalId: null, mailboxPath: null, lastSteer: 0, planVersion: nextVersion(ctx) };
|
||||
planningContextPending = true;
|
||||
resyncReason = null;
|
||||
writePlan(ctx, "");
|
||||
@@ -288,7 +323,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
pi.on("before_agent_start", async (_event, ctx) => {
|
||||
if (state.phase === "working") {
|
||||
return {
|
||||
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`,
|
||||
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 its durable mailbox 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[Kimi K3]`,
|
||||
};
|
||||
}
|
||||
if (!planningContextPending) return;
|
||||
@@ -317,6 +352,11 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
|
||||
pi.on("turn_end", async (_event, ctx) => {
|
||||
updateWidget(ctx);
|
||||
if (state.phase !== "working") return;
|
||||
workerTurns++;
|
||||
if (workerTurns < 50) return;
|
||||
workerTurns = 0;
|
||||
publishWorkerView(ctx, "turns");
|
||||
});
|
||||
|
||||
pi.on("tool_call", async (event, ctx) => {
|
||||
@@ -341,6 +381,11 @@ 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") {
|
||||
deliverWorkerSteers(ctx);
|
||||
publishWorkerView(ctx, "settled");
|
||||
return;
|
||||
}
|
||||
if (state.phase !== "planning" || !ctx.hasUI) return;
|
||||
let printed = "";
|
||||
while (true) {
|
||||
@@ -369,7 +414,7 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
}
|
||||
if (choice === "Cancel") {
|
||||
rmSync(planPath(ctx), { force: true });
|
||||
state = { ...state, phase: null, supervisorPaneId: null, approvalId: null, planVersion: null };
|
||||
state = { ...state, phase: null, supervisorPaneId: null, approvalId: null, mailboxPath: null, lastSteer: 0, planVersion: null };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
ctx.ui.notify("Plan discarded.", "info");
|
||||
@@ -381,12 +426,14 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
state = { ...state, phase: "working" };
|
||||
resyncReason = "The plan was approved.";
|
||||
persist();
|
||||
startWorkerTimers(ctx);
|
||||
publishWorkerView(ctx, "ready");
|
||||
updateWidget(ctx);
|
||||
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 };
|
||||
state = { ...state, phase: "planning", supervisorPaneId: null, approvalId: null, mailboxPath: null, lastSteer: 0 };
|
||||
persist();
|
||||
updateWidget(ctx);
|
||||
}
|
||||
@@ -404,13 +451,20 @@ export default function piGoalsExtension(pi: ExtensionAPI): void {
|
||||
supervisorModel: last?.data?.supervisorModel ?? null,
|
||||
supervisorPaneId: last?.data?.supervisorPaneId ?? null,
|
||||
approvalId: last?.data?.approvalId ?? null,
|
||||
mailboxPath: last?.data?.mailboxPath ?? null,
|
||||
lastSteer: last?.data?.lastSteer ?? 0,
|
||||
planVersion: last?.data?.planVersion ?? null,
|
||||
};
|
||||
planningContextPending = state.phase === "planning";
|
||||
resyncReason = state.phase === "working" ? "New session." : null;
|
||||
if (state.phase === "working") startWorkerTimers(ctx);
|
||||
updateWidget(ctx);
|
||||
});
|
||||
|
||||
pi.on("session_shutdown", async () => {
|
||||
stopWorkerTimers();
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "CompleteGoal",
|
||||
label: "Goal signoff",
|
||||
|
||||
+106
@@ -0,0 +1,106 @@
|
||||
import { mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from "node:fs";
|
||||
import { join, resolve } from "node:path";
|
||||
|
||||
const VERSION = 1;
|
||||
const VIEWS = "views";
|
||||
const STEERS = "steers";
|
||||
|
||||
export interface SupervisorMailbox {
|
||||
path: string;
|
||||
workerSessionId: string;
|
||||
planPath: string;
|
||||
approvalId: string;
|
||||
}
|
||||
|
||||
export interface WorkerView {
|
||||
version: 1;
|
||||
sequence: number;
|
||||
reason: "ready" | "settled" | "turns" | "interval";
|
||||
text: string;
|
||||
timestamp: string;
|
||||
}
|
||||
|
||||
export interface WorkerSteer {
|
||||
version: 1;
|
||||
sequence: number;
|
||||
instruction: string;
|
||||
timestamp: string;
|
||||
}
|
||||
|
||||
function writeJson(path: string, value: object): void {
|
||||
const temporary = `${path}.${process.pid}.tmp`;
|
||||
writeFileSync(temporary, `${JSON.stringify(value)}\n`);
|
||||
renameSync(temporary, path);
|
||||
}
|
||||
|
||||
function readJson<T>(path: string): T | null {
|
||||
try {
|
||||
return JSON.parse(readFileSync(path, "utf8")) as T;
|
||||
} catch (error: unknown) {
|
||||
if ((error as NodeJS.ErrnoException).code === "ENOENT") return null;
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
function sequence(name: string, prefix: string): number | null {
|
||||
const match = new RegExp(`^${prefix}-(\\d+)\\.json$`).exec(name);
|
||||
return match ? Number(match[1]) : null;
|
||||
}
|
||||
|
||||
function nextSequence(directory: string, prefix: string): number {
|
||||
return Math.max(0, ...readdirSync(directory).flatMap((name) => {
|
||||
const value = sequence(name, prefix);
|
||||
return value === null ? [] : [value];
|
||||
})) + 1;
|
||||
}
|
||||
|
||||
export function createMailbox(cwd: string, workerSessionId: string, approvalId: string, planPath: string): SupervisorMailbox {
|
||||
const path = resolve(cwd, ".pi", "goals-supervision", workerSessionId, approvalId);
|
||||
mkdirSync(join(path, VIEWS), { recursive: true });
|
||||
mkdirSync(join(path, STEERS), { recursive: true });
|
||||
return { path, workerSessionId, approvalId, planPath };
|
||||
}
|
||||
|
||||
export function readyMailbox(path: string): void {
|
||||
writeJson(join(path, "ready.json"), { version: VERSION, timestamp: new Date().toISOString() });
|
||||
}
|
||||
|
||||
export function supervisorReady(path: string): boolean {
|
||||
return readJson<{ version?: number }>(join(path, "ready.json"))?.version === VERSION;
|
||||
}
|
||||
|
||||
export function writeWorkerView(mailbox: SupervisorMailbox, reason: WorkerView["reason"], text: string): WorkerView {
|
||||
const directory = join(mailbox.path, VIEWS);
|
||||
const view: WorkerView = { version: 1, sequence: nextSequence(directory, "view"), reason, text, timestamp: new Date().toISOString() };
|
||||
writeJson(join(directory, `view-${view.sequence}.json`), view);
|
||||
return view;
|
||||
}
|
||||
|
||||
export function workerViewsAfter(path: string, after: number): WorkerView[] {
|
||||
const directory = join(path, VIEWS);
|
||||
return readdirSync(directory)
|
||||
.flatMap((name) => {
|
||||
const value = sequence(name, "view");
|
||||
return value !== null && value > after ? [readJson<WorkerView>(join(directory, name))] : [];
|
||||
})
|
||||
.filter((view): view is WorkerView => view !== null)
|
||||
.sort((a, b) => a.sequence - b.sequence);
|
||||
}
|
||||
|
||||
export function writeWorkerSteer(path: string, instruction: string): WorkerSteer {
|
||||
const directory = join(path, STEERS);
|
||||
const steer: WorkerSteer = { version: 1, sequence: nextSequence(directory, "steer"), instruction, timestamp: new Date().toISOString() };
|
||||
writeJson(join(directory, `steer-${steer.sequence}.json`), steer);
|
||||
return steer;
|
||||
}
|
||||
|
||||
export function workerSteersAfter(path: string, after: number): WorkerSteer[] {
|
||||
const directory = join(path, STEERS);
|
||||
return readdirSync(directory)
|
||||
.flatMap((name) => {
|
||||
const value = sequence(name, "steer");
|
||||
return value !== null && value > after ? [readJson<WorkerSteer>(join(directory, name))] : [];
|
||||
})
|
||||
.filter((steer): steer is WorkerSteer => steer !== null)
|
||||
.sort((a, b) => a.sequence - b.sequence);
|
||||
}
|
||||
+2
-2
@@ -3,7 +3,7 @@
|
||||
*
|
||||
* 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 main session implements it,
|
||||
* while a visible forked Pi session supervises through pi-supervise.
|
||||
* while a visible forked Pi session supervises through a durable mailbox.
|
||||
*
|
||||
* 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
|
||||
@@ -161,7 +161,7 @@ export function resync(plan: string, planRel: string, why: string): string {
|
||||
<system-reminder>
|
||||
${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
|
||||
session supervises you through a durable mailbox. The human's latest message outranks the plan: if it
|
||||
changes scope, amend the plan rather than preserving an obsolete decision.
|
||||
|
||||
${plan}
|
||||
|
||||
@@ -1,52 +0,0 @@
|
||||
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
|
||||
|
||||
const PAIR_EVENT = "pi-supervise:pair:v1";
|
||||
const WORKER_STATE_EVENT = "pi-supervise:worker-state:v1";
|
||||
const WORKER_PAIRED_EVENT = "pi-supervise:worker-paired:v1";
|
||||
const API_READY_EVENT = "pi-supervise:api-ready:v1";
|
||||
const TIMEOUT_MS = 15_000;
|
||||
export const SUPERVISOR_STARTUP_TIMEOUT_MS = 5 * 60_000;
|
||||
|
||||
type Events = { emit(name: string, value: unknown): boolean; on(name: string, handler: (value: any) => void): void };
|
||||
|
||||
function wait<T>(start: (resolve: (value: T) => void, reject: (error: Error) => void) => void, message: string, timeoutMs = TIMEOUT_MS): Promise<T> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const timer = setTimeout(() => reject(new Error(message)), timeoutMs);
|
||||
start((value) => { clearTimeout(timer); resolve(value); }, (error) => { clearTimeout(timer); reject(error); });
|
||||
});
|
||||
}
|
||||
|
||||
export function pairWithPiSupervise(pi: ExtensionAPI, workerIntercomId: string, goal: string): Promise<void> {
|
||||
const events = (pi as unknown as { events: Events }).events;
|
||||
return wait((resolve, reject) => events.emit(PAIR_EVENT, { version: 1, workerIntercomId, goal, resolve, reject }), "pi-supervise did not accept the visible-supervisor pairing request.");
|
||||
}
|
||||
|
||||
export interface WorkerPiSupervise {
|
||||
intercomId: string;
|
||||
waitForPair(timeoutMs?: number): Promise<void>;
|
||||
}
|
||||
|
||||
export function workerPiSupervise(pi: ExtensionAPI, timeoutMs = TIMEOUT_MS): Promise<WorkerPiSupervise> {
|
||||
const events = (pi as unknown as { events: Events }).events;
|
||||
let paired = false;
|
||||
let resolvePair: (() => void) | undefined;
|
||||
events.on(WORKER_PAIRED_EVENT, () => {
|
||||
paired = true;
|
||||
resolvePair?.();
|
||||
});
|
||||
return wait((resolve, reject) => {
|
||||
let resolved = false;
|
||||
const request = () => events.emit(WORKER_STATE_EVENT, (state: { intercomId?: string; paired?: boolean }) => {
|
||||
if (resolved) return;
|
||||
if (!state.intercomId) return reject(new Error("pi-supervise returned no worker intercom ID."));
|
||||
if (state.paired) return reject(new Error("This worker is already paired with a supervisor. Stop that supervision before selecting Ready."));
|
||||
resolved = true;
|
||||
resolve({
|
||||
intercomId: state.intercomId,
|
||||
waitForPair: (pairTimeoutMs = timeoutMs) => paired ? Promise.resolve() : wait((pairResolve) => { resolvePair = pairResolve; }, "The visible supervisor did not pair with this worker.", pairTimeoutMs),
|
||||
});
|
||||
});
|
||||
events.on(API_READY_EVENT, request);
|
||||
request();
|
||||
}, "pi-supervise did not publish this worker's intercom state.", timeoutMs);
|
||||
}
|
||||
+46
-28
@@ -3,18 +3,19 @@ import { resolve } from "node:path";
|
||||
import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent";
|
||||
import { Type } from "typebox";
|
||||
import { approvalPath, goalBlock, hashGoalBlock, repositoryState, verifyOutputPath, writeApproval } from "./approval.js";
|
||||
import { pairWithPiSupervise } from "./supervise.js";
|
||||
import { readyMailbox, workerViewsAfter, writeWorkerSteer } from "./mailbox.js";
|
||||
|
||||
const BOOTSTRAPPED = "pi-goals-visible-supervisor-v1";
|
||||
const BOOTSTRAPPED = "pi-goals-visible-supervisor-v2";
|
||||
const INITIAL_COMPACT_AT_TOKENS = 20_000;
|
||||
const COMPACT_AT_TOKENS = 100_000;
|
||||
const WRITER_TOOLS = new Set(["bash", "edit", "write", "multi_edit", "multiedit", "apply_patch", "notebook_edit", "edit_file", "write_file", "quick_edit", "target_edit"]);
|
||||
|
||||
interface SupervisorConfig {
|
||||
workerSessionId: string;
|
||||
workerIntercomId: string;
|
||||
ownerSessionId: string;
|
||||
planPath: string;
|
||||
approvalId: string;
|
||||
mailboxPath: string;
|
||||
}
|
||||
|
||||
function result(text: string, isError = false) {
|
||||
@@ -30,10 +31,10 @@ function requiredEnv(name: string): string {
|
||||
function config(): SupervisorConfig {
|
||||
return {
|
||||
workerSessionId: requiredEnv("PI_GOALS_WORKER_ID"),
|
||||
workerIntercomId: requiredEnv("PI_GOALS_WORKER_INTERCOM_ID"),
|
||||
ownerSessionId: requiredEnv("PI_GOALS_OWNER_SESSION_ID"),
|
||||
planPath: resolve(requiredEnv("PI_GOALS_PLAN_PATH")),
|
||||
approvalId: requiredEnv("PI_GOALS_APPROVAL_ID"),
|
||||
mailboxPath: resolve(requiredEnv("PI_GOALS_MAILBOX_PATH")),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -67,9 +68,9 @@ function latestWorkerView(ctx: ExtensionContext): string | 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 and pi-supervise worker views. The complete plan at ${settings.planPath} is the source of truth; read it directly after every compaction.
|
||||
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 and worker views. The complete plan at ${settings.planPath} is the source of truth; read it directly after every compaction.
|
||||
|
||||
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 a saved nonempty verification-output file. When its discriminator is positively satisfied and the worker view says no work is active, call ApproveGoal with that repository-relative path. 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]`;
|
||||
Use SteerWorker to 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 a saved nonempty verification-output file. When its discriminator is positively satisfied and the worker view says no work is active, call ApproveGoal with that repository-relative path. Then call SteerWorker and tell the worker to call CompleteGoal with the exact goal text. Do not call done until every plan goal is [x]. -- PI[Kimi K3]`;
|
||||
}
|
||||
|
||||
export function isVisibleSupervisor(): boolean {
|
||||
@@ -80,6 +81,15 @@ export function registerVisibleSupervisor(pi: ExtensionAPI): void {
|
||||
const settings = config();
|
||||
let compacting = false;
|
||||
let bootstrapping = false;
|
||||
let deliveredView = 0;
|
||||
let viewTimer: ReturnType<typeof setInterval> | undefined;
|
||||
|
||||
const deliverWorkerViews = (): void => {
|
||||
for (const view of workerViewsAfter(settings.mailboxPath, deliveredView)) {
|
||||
deliveredView = view.sequence;
|
||||
pi.sendUserMessage(view.text, { deliverAs: "followUp" });
|
||||
}
|
||||
};
|
||||
|
||||
const bootstrap = async (ctx: ExtensionContext): Promise<void> => {
|
||||
if (bootstrapping) return;
|
||||
@@ -87,9 +97,14 @@ export function registerVisibleSupervisor(pi: ExtensionAPI): void {
|
||||
if (entries.some((entry: { type?: string; customType?: string }) => entry.type === "custom" && entry.customType === BOOTSTRAPPED)) return;
|
||||
bootstrapping = true;
|
||||
try {
|
||||
await pairWithPiSupervise(pi, settings.workerIntercomId, settings.planPath);
|
||||
pi.appendEntry(BOOTSTRAPPED, { version: 1, workerSessionId: settings.workerSessionId, planPath: settings.planPath });
|
||||
pi.sendUserMessage("Supervision is paired. Inspect the worker and give its next concrete instruction.");
|
||||
const active = pi.getActiveTools();
|
||||
pi.setActiveTools(active.filter((tool) => !WRITER_TOOLS.has(tool.toLowerCase())));
|
||||
const writers = pi.getActiveTools().filter((tool) => WRITER_TOOLS.has(tool.toLowerCase()));
|
||||
if (writers.length) throw new Error(`Could not remove supervisor writing tools: ${writers.join(", ")}`);
|
||||
pi.appendEntry(BOOTSTRAPPED, { version: 2, workerSessionId: settings.workerSessionId, planPath: settings.planPath });
|
||||
readyMailbox(settings.mailboxPath);
|
||||
viewTimer = setInterval(deliverWorkerViews, 1_000);
|
||||
deliverWorkerViews();
|
||||
} catch (error) {
|
||||
ctx.ui.notify(`Supervisor startup failed: ${error instanceof Error ? error.message : String(error)}`, "error");
|
||||
}
|
||||
@@ -119,16 +134,15 @@ export function registerVisibleSupervisor(pi: ExtensionAPI): void {
|
||||
pi.on("session_start", async (_event, ctx) => {
|
||||
setImmediate(() => { bootstrapAfterInitialCompaction(ctx); });
|
||||
});
|
||||
|
||||
pi.on("before_agent_start", async (_event, ctx) => {
|
||||
return { systemPrompt: `${ctx.getSystemPrompt()}\n\n${supervisorPrompt(settings)}` };
|
||||
pi.on("session_shutdown", async () => {
|
||||
if (viewTimer) clearInterval(viewTimer);
|
||||
});
|
||||
|
||||
pi.on("before_agent_start", async (_event, ctx) => ({ systemPrompt: `${ctx.getSystemPrompt()}\n\n${supervisorPrompt(settings)}` }));
|
||||
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.`,
|
||||
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. The canonical plan remains ${settings.planPath}.`,
|
||||
onComplete: () => {
|
||||
compacting = false;
|
||||
ctx.ui.notify("Supervisor context compacted at 100k tokens.", "info");
|
||||
@@ -140,6 +154,20 @@ export function registerVisibleSupervisor(pi: ExtensionAPI): void {
|
||||
});
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "SteerWorker",
|
||||
label: "Steer worker",
|
||||
executionMode: "sequential",
|
||||
description: "Write one concrete instruction for the implementation worker.",
|
||||
parameters: Type.Object({ instruction: Type.String({ description: "Concrete next instruction for the worker." }) }),
|
||||
async execute(_id, params) {
|
||||
const instruction = params.instruction.trim();
|
||||
if (!instruction) return result("A worker instruction cannot be empty.", true);
|
||||
const steer = writeWorkerSteer(settings.mailboxPath, instruction);
|
||||
return result(`Worker instruction ${steer.sequence} recorded.`);
|
||||
},
|
||||
});
|
||||
|
||||
pi.registerTool({
|
||||
name: "ApproveGoal",
|
||||
label: "Approve goal",
|
||||
@@ -171,20 +199,10 @@ export function registerVisibleSupervisor(pi: ExtensionAPI): void {
|
||||
if (!verifiedOutput) return result("Cannot approve without a nonempty repository-relative verification-output file.", true);
|
||||
const path = approvalPath(ctx.cwd, settings.ownerSessionId, params.goal);
|
||||
writeApproval(path, {
|
||||
version: 3,
|
||||
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 },
|
||||
verifyOutputPath: verifiedOutput,
|
||||
supervisor: { sessionId: ctx.sessionManager.getSessionId(), runId: null },
|
||||
timestamp: new Date().toISOString(),
|
||||
version: 3, 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 }, verifyOutputPath: verifiedOutput,
|
||||
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.`);
|
||||
},
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
export interface SessionBlock {
|
||||
type?: string;
|
||||
id?: string;
|
||||
name?: string;
|
||||
text?: string;
|
||||
}
|
||||
|
||||
export interface SessionMessage {
|
||||
role?: string;
|
||||
content?: string | SessionBlock[];
|
||||
toolCallId?: string;
|
||||
}
|
||||
|
||||
export interface SessionEntry {
|
||||
type?: string;
|
||||
summary?: string;
|
||||
message?: SessionMessage;
|
||||
}
|
||||
|
||||
function text(message: SessionMessage): string {
|
||||
if (typeof message.content === "string") return message.content;
|
||||
return (message.content ?? []).flatMap((block) => {
|
||||
if (block.type === "text" && block.text) return [block.text];
|
||||
if (block.type === "toolCall") return [`tool: ${block.name ?? "unknown"}`];
|
||||
return [];
|
||||
}).join("\n");
|
||||
}
|
||||
|
||||
function outstandingTools(entries: SessionEntry[]): string[] {
|
||||
const calls = new Map<string, string>();
|
||||
const results = new Set<string>();
|
||||
for (const entry of entries) {
|
||||
for (const block of Array.isArray(entry.message?.content) ? entry.message.content : []) {
|
||||
if (block.type === "toolCall" && block.id) calls.set(block.id, block.name ?? "unknown");
|
||||
}
|
||||
if (entry.message?.role === "toolResult" && entry.message.toolCallId) results.add(entry.message.toolCallId);
|
||||
}
|
||||
return [...calls].filter(([id]) => !results.has(id)).map(([, name]) => name);
|
||||
}
|
||||
|
||||
export function workerView(entries: SessionEntry[], reason: "ready" | "settled" | "turns" | "interval"): string {
|
||||
const summary = [...entries].reverse().find((entry) => entry.type === "compaction" && entry.summary)?.summary;
|
||||
const recent = entries.flatMap((entry) => entry.type === "message" && entry.message ? [text(entry.message)] : []).filter(Boolean).slice(-12).join("\n\n").slice(-12_000);
|
||||
const outstanding = outstandingTools(entries);
|
||||
const state = reason === "settled" ? "stopped" : reason === "ready" ? "is ready to begin" : "is still working";
|
||||
return `The worker ${state}.\n\nreview trigger: ${reason}\ntool calls with no result: ${outstanding.join(", ") || "none"}\n\n${summary ? `last compaction summary:\n${summary}\n\n` : ""}recent worker transcript:\n${recent || "none"}`;
|
||||
}
|
||||
Reference in new issue
Block a user