Add persistent plan steward

This commit is contained in:
wassname2
2026-09-03 14:17:39 +08:00
parent f09d443d88
commit 3f0aadfffa
7 changed files with 909 additions and 33 deletions
+223 -4
View File
@@ -16,12 +16,24 @@ function setup(
const tools = new Map<string, any>();
const entries: Array<{ type: string; customType: string; data: unknown }> = [];
const events: string[] = [];
const messages: Array<{ content: string; display?: boolean }> = [];
const messages: Array<{ content: string; display?: boolean; customType?: string }> = [];
const busHandlers = new Map<string, Set<(value: unknown) => unknown>>();
const bus = {
on(name: string, handler: (value: unknown) => unknown) {
const handlers = busHandlers.get(name) ?? new Set();
handlers.add(handler);
busHandlers.set(name, handlers);
return () => handlers.delete(handler);
},
emit(name: string, value: unknown) {
for (const handler of busHandlers.get(name) ?? []) void handler(value);
},
};
const ctx = {
cwd,
hasUI: true,
isIdle: () => true,
sessionManager: { getSessionId: () => "session-a", getEntries: () => entries },
sessionManager: { getSessionId: () => "session-a", getSessionFile: () => join(cwd, "session-a.jsonl"), getEntries: () => entries },
ui: {
theme: { fg: (_kind: string, text: string) => text },
setStatus: () => {},
@@ -42,14 +54,15 @@ function setup(
on: (name: string, handler: any) => hooks.set(name, handler),
appendEntry: (customType: string, data: unknown) => entries.push({ type: "custom", customType, data }),
registerTool: (tool: any) => tools.set(tool.name, tool),
sendMessage: (message: { content: string; display?: boolean }) => {
events: bus,
sendMessage: (message: { content: string; display?: boolean; customType?: string }) => {
events.push("display");
messages.push(message);
},
sendUserMessage: (message: string) => messages.push({ content: message }),
};
piGoalsExtension(pi as unknown as ExtensionAPI);
return { commands, ctx, cwd, entries, events, hooks, messages, tools };
return { bus, commands, ctx, cwd, entries, events, hooks, messages, tools };
}
describe("/goals draft flow", () => {
@@ -154,6 +167,212 @@ describe("/goals draft flow", () => {
}
});
it("forks a persistent steward at Ready and resumes it before sign-off", async () => {
const flow = setup(["Ready"]);
flow.bus.on("subagents:rpc:v1:request", (value: unknown) => {
const request = value as { requestId: string; method: string; params: Record<string, unknown> };
const runId = request.method === "spawn" ? "plan-review-run" : "signoff-review-run";
flow.bus.emit(`subagents:rpc:v1:reply:${request.requestId}`, {
version: 1,
requestId: request.requestId,
success: true,
data: { details: { runId } },
});
});
try {
await flow.hooks.get("session_start")({}, flow.ctx);
await flow.commands.get("goals").handler("steward on", flow.ctx);
await flow.commands.get("goals").handler("objective", flow.ctx);
const planPath = join(flow.cwd, ".pi/plan/session-a-v1.md");
writeFileSync(planPath, "# Plan\n\n## User-visible result\n\nA file exists.\n\n## Goals\n\n1. [ ] goal: make the file\n - discriminator: the file can be read\n");
await flow.hooks.get("agent_settled")({}, flow.ctx);
expect(flow.entries.at(-1)?.data).toMatchObject({
phase: "reviewing",
stewardReview: { kind: "plan", runId: "plan-review-run" },
});
expect(flow.messages.some((message) => message.content.includes("Work the goals"))).toBe(false);
const blockedDuringReview = await flow.hooks.get("tool_call")({ toolName: "write", input: { path: "README.md" } }, flow.ctx);
expect(blockedDuringReview?.block).toBe(true);
const reviewContext = await flow.hooks.get("before_agent_start")({}, flow.ctx);
expect(reviewContext.message.content).toContain("[PLAN STEWARD REVIEW]");
flow.bus.emit("subagent:async-complete", {
runId: "plan-review-run",
success: true,
results: [{ structuredOutput: {
decision: "approve",
reason: "The goals preserve the requested result.",
nextAction: "Start work.",
contractDrift: [],
unresolvedDecisions: [],
} }],
});
await new Promise((resolve) => setImmediate(resolve));
expect(flow.entries.at(-1)?.data).toMatchObject({ phase: "working", stewardRunId: "plan-review-run" });
expect(flow.messages.some((message) => message.content.includes("Work the goals"))).toBe(true);
const firstSignoff = await flow.tools.get("CompleteGoal").execute("", { goal: "make the file" }, undefined, undefined, flow.ctx);
expect(firstSignoff.isError).toBe(false);
expect(firstSignoff.content[0].text).toContain("Sign-off paused");
expect(flow.entries.at(-1)?.data).toMatchObject({
stewardReview: { kind: "signoff", runId: "signoff-review-run", goal: "make the file" },
});
flow.bus.emit("subagent:async-complete", {
runId: "signoff-review-run",
success: true,
results: [{ structuredOutput: {
decision: "approve",
reason: "The sign-off remains in scope.",
nextAction: "Run the fresh evidence review.",
contractDrift: [],
unresolvedDecisions: [],
} }],
});
await new Promise((resolve) => setImmediate(resolve));
expect(flow.entries.at(-1)?.data).toMatchObject({
stewardRunId: "signoff-review-run",
stewardApproval: { goal: "make the file" },
});
expect(flow.messages.at(-1)?.content).toContain("Call CompleteGoal again");
} finally {
rmSync(flow.cwd, { recursive: true, force: true });
}
});
it("invalidates a plan approval when the human corrects it during review", async () => {
const flow = setup(["Ready"]);
flow.bus.on("subagents:rpc:v1:request", (value: unknown) => {
const request = value as { requestId: string };
flow.bus.emit(`subagents:rpc:v1:reply:${request.requestId}`, {
version: 1, requestId: request.requestId, success: true, data: { details: { runId: "review-run" } },
});
});
try {
await flow.hooks.get("session_start")({}, flow.ctx);
await flow.commands.get("goals").handler("steward on", flow.ctx);
await flow.commands.get("goals").handler("objective", flow.ctx);
const planPath = join(flow.cwd, ".pi/plan/session-a-v1.md");
writeFileSync(planPath, "# Plan\n\n## Goals\n\n1. [ ] goal: make the file\n\n## Log\n\n## Interview\n");
await flow.hooks.get("agent_settled")({}, flow.ctx);
await flow.hooks.get("input")({ text: "The file must be CSV.", source: "interactive" }, flow.ctx);
flow.bus.emit("subagent:async-complete", {
runId: "review-run",
success: true,
results: [{ structuredOutput: {
decision: "approve", reason: "The old plan was sound.", nextAction: "Start.", contractDrift: [], unresolvedDecisions: [],
} }],
});
await new Promise((resolve) => setImmediate(resolve));
expect(flow.entries.at(-1)?.data).toMatchObject({ phase: "planning" });
expect(flow.messages.some((message) => message.content.includes("Work the goals"))).toBe(false);
expect(flow.messages.at(-1)?.content).toContain("plan changed while the steward reviewed it");
} finally {
rmSync(flow.cwd, { recursive: true, force: true });
}
});
it("reconciles a completed pending steward review after session restart", async () => {
const flow = setup([]);
const methods: string[] = [];
flow.bus.on("subagents:rpc:v1:request", (value: unknown) => {
const request = value as { requestId: string; method: string };
methods.push(request.method);
if (request.method === "resume") {
flow.bus.emit("subagent:async-complete", {
runId: "lost-review",
success: true,
results: [{ structuredOutput: {
decision: "approve", reason: "Stale completion.", nextAction: "Start.", contractDrift: [], unresolvedDecisions: [],
} }],
});
}
flow.bus.emit(`subagents:rpc:v1:reply:${request.requestId}`, {
version: 1,
requestId: request.requestId,
success: true,
data: request.method === "status" ? { text: "State: complete" } : { details: { runId: "recovered-review" } },
});
});
try {
mkdirSync(join(flow.cwd, ".pi/plan"), { recursive: true });
writeFileSync(join(flow.cwd, ".pi/plan/session-a-v1.md"), "# Plan\n\n## Goals\n\n1. [ ] goal: make the file\n");
flow.entries.push({
type: "custom",
customType: "pi-goals-state",
data: {
phase: "reviewing", judgeModel: null, planVersion: 1, autoIntervalMs: null, autoPaused: false,
stewardEnabled: true, stewardRunId: "lost-review", approvedPlan: null, stewardApproval: null,
stewardReview: { kind: "plan", runId: "lost-review", snapshotHash: "old" },
},
});
await flow.hooks.get("session_start")({}, flow.ctx);
expect(methods).toEqual(["status", "resume"]);
expect(flow.entries.at(-1)?.data).toMatchObject({ phase: "reviewing", stewardReview: { runId: "recovered-review" } });
expect(flow.messages.some((message) => message.content.includes("Work the goals"))).toBe(false);
} finally {
rmSync(flow.cwd, { recursive: true, force: true });
}
});
it("resumes the same steward after it requests a plan revision", async () => {
const flow = setup(["Ready", "Ready"]);
const methods: string[] = [];
flow.bus.on("subagents:rpc:v1:request", (value: unknown) => {
const request = value as { requestId: string; method: string };
methods.push(request.method);
flow.bus.emit(`subagents:rpc:v1:reply:${request.requestId}`, {
version: 1,
requestId: request.requestId,
success: true,
data: { details: { runId: methods.length === 1 ? "first-review" : "revised-review" } },
});
});
try {
await flow.hooks.get("session_start")({}, flow.ctx);
await flow.commands.get("goals").handler("steward on", flow.ctx);
await flow.commands.get("goals").handler("objective", flow.ctx);
const planPath = join(flow.cwd, ".pi/plan/session-a-v1.md");
writeFileSync(planPath, "# Plan\n\n## Goals\n\n1. [ ] goal: make it better\n");
await flow.hooks.get("agent_settled")({}, flow.ctx);
flow.bus.emit("subagent:async-complete", {
runId: "first-review",
success: true,
results: [{ structuredOutput: {
decision: "revise_plan",
reason: "The result is not observable.",
nextAction: "Name the artifact.",
contractDrift: [],
unresolvedDecisions: [],
} }],
});
await new Promise((resolve) => setImmediate(resolve));
expect(flow.entries.at(-1)?.data).toMatchObject({ phase: "planning", stewardRunId: "first-review" });
writeFileSync(planPath, "# Plan\n\n## Goals\n\n1. [ ] goal: create report.html\n");
await flow.hooks.get("agent_settled")({}, flow.ctx);
expect(methods).toEqual(["spawn", "resume"]);
expect(flow.entries.at(-1)?.data).toMatchObject({ stewardReview: { runId: "revised-review" } });
} finally {
rmSync(flow.cwd, { recursive: true, force: true });
}
});
it("refuses to enable a steward after an unreviewed plan is already working", async () => {
const flow = setup(["Ready"]);
try {
await flow.hooks.get("session_start")({}, flow.ctx);
await flow.commands.get("goals").handler("objective", flow.ctx);
const planPath = join(flow.cwd, ".pi/plan/session-a-v1.md");
writeFileSync(planPath, "# Plan\n\n## Goals\n\n1. [ ] goal: make the file\n");
await flow.hooks.get("agent_settled")({}, flow.ctx);
await flow.commands.get("goals").handler("steward on", flow.ctx);
expect(flow.entries.at(-1)?.data).toMatchObject({ phase: "working", stewardEnabled: false, stewardRunId: null });
} finally {
rmSync(flow.cwd, { recursive: true, force: true });
}
});
it("edits a plan in Pi and cancels without starting work", async () => {
const original = "# Plan\n\n## Goals\n\n1. [ ] goal: original\n";
const edited = "# Plan\n\n## Goals\n\n1. [ ] goal: edited\n";
+64
View File
@@ -0,0 +1,64 @@
import { describe, expect, it } from "vitest";
import { parseStewardDecision, rpcRunId, stewardCompletion, stewardContract } from "../src/steward.js";
const decision = {
decision: "approve",
reason: "The goal remains faithful.",
nextAction: "Send it to the evidence judge.",
contractDrift: [],
unresolvedDecisions: [],
};
describe("persistent steward protocol", () => {
it("extracts run ids from pi-subagents RPC replies", () => {
expect(rpcRunId({ details: { runId: "run-1" } })).toBe("run-1");
expect(rpcRunId({ runId: "run-2" })).toBe("run-2");
expect(rpcRunId({ details: {} })).toBeNull();
});
it("accepts only the bounded structured decision", () => {
expect(parseStewardDecision(decision)).toEqual(decision);
expect(parseStewardDecision({ ...decision, decision: "done" })).toBeNull();
expect(parseStewardDecision({ ...decision, unresolvedDecisions: "none" })).toBeNull();
expect(parseStewardDecision({ ...decision, unresolvedDecisions: ["Choose the published scope."] })?.decision).toBe("needs_user");
expect(parseStewardDecision({ ...decision, contractDrift: ["The output changed."] })?.decision).toBe("revise_plan");
});
it("keeps the approved contract compact across progress and evidence growth", () => {
const plan = `1. [/] goal: make the file
- discriminator: the file can be read
- tasks:
1. [x] write it
- evidence:
- huge quoted log
- another artifact
2. [ ] goal: publish it`;
expect(stewardContract(plan)).toBe(`1. [ ] goal: make the file
- discriminator: the file can be read
- tasks:
1. [ ] write it
- evidence: (checked separately by the fresh evidence judge)
2. [ ] goal: publish it`);
expect(stewardContract(plan, { preserveGoalStatus: true })).toContain("1. [/] goal: make the file");
});
it("correlates a structured completion by run id", () => {
expect(stewardCompletion({
runId: "run-3",
success: true,
results: [{ structuredOutput: decision }],
})).toEqual({ runId: "run-3", decision, error: null });
expect(stewardCompletion({
runId: "run-4",
success: false,
summary: "child failed",
results: [{}],
})).toEqual({ runId: "run-4", decision: null, error: "child failed" });
expect(stewardCompletion({
runId: "run-5",
success: true,
results: [{ structuredOutput: decision, effects: { fileMutation: { status: "observed", attempted: true } } }],
})?.error).toBe("steward attempted or produced a file mutation");
});
});