diff --git a/src/index.ts b/src/index.ts index e7b45d5..a3167a2 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,4 +1,5 @@ // Pi/OpenAI: Plan and supervise in the main chat; delegate implementation to a visible worker. +import { execFileSync } from "node:child_process"; import { createHash, randomUUID } from "node:crypto"; import { type FSWatcher, mkdirSync, readdirSync, readFileSync, watch, writeFileSync } from "node:fs"; import { basename, dirname, isAbsolute, join, relative, resolve, sep } from "node:path"; @@ -59,8 +60,9 @@ const WORKER = "goals-worker"; const REPORT = "pi-goals-report", REVIEW = "pi-goals-report-review", REVIEW_DRAFT = "pi-goals-review-draft", REVIEW_REMINDER = "pi-goals-review-reminder"; const RUN = "pi-goals-worker-run", STOP = "pi-goals-worker-stop"; interface WorkerStop { type: "stopped"; entryId: string; to: string; requestId: string; plan: string; text: string; identity: Peer; } -interface Report { id: string; plan: string; session: string; sessionFile: string; requestId: string; text: string; supersedes?: string; } -interface ReportReview { id: string; report: string; verdict: string; content: string; continuation: string; } +interface Report { id: string; plan: string; session: string; sessionFile: string; requestId: string; task?: string; text: string; } +interface ReportReview { id: string; reportId?: string; report?: string; verdict: string; content: string; continuation: string; } +const reviewedReportId = (review: ReportReview) => review.reportId ?? review.report; const WIDGET_GOAL_LIMIT = 3; type Mode = "chat" | "planning" | "supervising" | "paused" | "solo"; type GoalStatus = "open" | "active" | "reported" | "done" | "cancelled"; @@ -70,7 +72,7 @@ interface Peer { interface State { mode: Mode; plan?: string; - worker?: { sessionFile?: string; intercomId?: string; paneId?: string; requestId?: string; parentId?: string; identity?: Peer }; + worker?: { sessionFile?: string; intercomId?: string; paneId?: string; requestId?: string; parentId?: string; task?: string; identity?: Peer }; parent?: { intercomId: string; requestId: string }; workerStopped?: boolean; pausedFrom?: "solo" | "supervising"; @@ -356,7 +358,13 @@ export default function mainSupervisor(pi: ExtensionAPI) { paneId: process.env.HERDR_PANE_ID ?? "", model: ctx.model ? ctx.model.provider + "/" + ctx.model.id : undefined }; } const records = (ctx: ExtensionContext, type: string): T[] => ctx.sessionManager.getBranch().flatMap(entry => entry.type === "custom" && entry.customType === type ? [entry.data as T] : []); - const pendingReports = (ctx: ExtensionContext) => records(ctx, REPORT).filter(report => !records(ctx, REVIEW).some(review => review.report === report.id) && !records(ctx, REPORT).some(newer => newer.session === report.session && newer.supersedes === report.id)); + const pendingReports = (ctx: ExtensionContext) => records(ctx, REPORT).filter(report => !records(ctx, REVIEW).some(review => reviewedReportId(review) === report.id)); + const reportLabel = (report: Report) => { + const revision = report.id.slice(report.id.lastIndexOf(":") + 1); + const task = report.task ? `${report.task.trim().replace(/\s+/g, " ").slice(0, 100)} — ` : ""; + const summary = report.text.trim().replace(/\s+/g, " ").slice(0, 160) || "No worker summary"; + return `revision ${revision}: ${task}${summary} (reportId ${report.id})`; + }; function recordReport(ctx: ExtensionContext, report: Report, wake = true) { if (records(ctx, REPORT).some(saved => saved.id === report.id)) return; pi.appendEntry(REPORT, report); @@ -365,11 +373,11 @@ export default function mainSupervisor(pi: ExtensionAPI) { } function remindReports(ctx: ExtensionContext) { if (state.child || state.mode !== "supervising") return; - const ids = pendingReports(ctx).map(report => report.id); - const fingerprint = digest(JSON.stringify(ids)); - if (!ids.length || records(ctx, REVIEW_REMINDER).at(-1) === fingerprint) return; + const reports = pendingReports(ctx); + const fingerprint = digest(JSON.stringify(reports.map(report => report.id))); + if (!reports.length || records(ctx, REVIEW_REMINDER).at(-1) === fingerprint) return; pi.appendEntry(REVIEW_REMINDER, fingerprint); - send(pendingReportReviews(ids)); + send(pendingReportReviews(reports.map(reportLabel))); } pi.registerEntryRenderer(REVIEW, entry => new Text((entry.data as ReportReview).content, 0, 0)); function registerChannel(ctx: ExtensionContext) { @@ -387,12 +395,12 @@ export default function mainSupervisor(pi: ExtensionAPI) { const stopped = stop?.type === "custom" ? stop.data as WorkerStop : undefined; entryId = run && run.id !== stopped?.entryId ? `${run.id}:disconnected` : stopped?.entryId || branch.filter(entry => entry.type === "message" && entry.message.role === "assistant").at(-1)?.id || entryId; } catch { /* Unknown history remains a review obligation, not proof of exit. */ } - recordReport(ctx, { id: `${event.sessionId}:${entryId}`, plan: state.plan, session: event.sessionId, sessionFile: state.worker.sessionFile || "", requestId: state.worker.requestId!, text: nativeMessages.disconnected }); + recordReport(ctx, { id: `${event.sessionId}:${entryId}`, plan: state.plan, session: event.sessionId, sessionFile: state.worker.sessionFile || "", requestId: state.worker.requestId!, task: state.worker.task, text: nativeMessages.disconnected }); } if (event.type !== "message" || !event.payload || typeof event.payload !== "object") return; const data = event.payload as { type?: string; to?: string; requestId?: string; plan?: string; sessionFile?: string; text?: string; identity?: Peer; entryId?: string; review?: ReportReview }; if (data.type === "review" && state.child && state.parent && event.fromSessionId === state.parent.intercomId && records(ctx, STATE).some(saved => saved.child && saved.parent?.intercomId === event.fromSessionId && saved.parent.requestId === data.requestId && saved.plan === data.plan) && data.sessionFile === ctx.sessionManager.getSessionFile() && data.review) { - const prior = records(ctx, REVIEW).find(saved => saved.report === data.review!.report); + const prior = records(ctx, REVIEW).find(saved => reviewedReportId(saved) === reviewedReportId(data.review!)); const review = prior || data.review; if (!prior && review.verdict === "changes_requested" && (state.mode === "paused" || data.plan !== state.plan || data.requestId !== state.parent.requestId || !review.continuation.trim())) return; if (!prior) { @@ -404,7 +412,7 @@ export default function mainSupervisor(pi: ExtensionAPI) { } if (data.type === "review_saved" && !state.child && data.review) { const draft = records(ctx, REVIEW_DRAFT).find(review => review.id === data.review!.id); - const report = records(ctx, REPORT).find(report => report.id === draft?.report); + const report = records(ctx, REPORT).find(report => report.id === (draft && reviewedReportId(draft))); if (!draft || !report || event.fromSessionId !== report.session || data.requestId !== report.requestId || !records(ctx, STATE).some(saved => saved.worker && saved.worker.parentId === data.to && saved.worker.requestId === report.requestId && saved.worker.intercomId === report.session) || data.plan !== report.plan) return; try { const saved = savedSession(report.sessionFile).getBranch().some(entry => entry.type === "custom" && entry.customType === REVIEW && JSON.stringify(entry.data) === JSON.stringify(draft)); @@ -420,9 +428,9 @@ export default function mainSupervisor(pi: ExtensionAPI) { workerRevision++; save(); sendAttachment(state.plan, event.fromSessionId, nativeMessages.attached(data.sessionFile)); } - if (data.type === "stopped" && event.fromSessionId === worker.intercomId && typeof data.text === "string") { + if (data.type === "stopped" && event.fromSessionId === worker.intercomId && typeof data.entryId === "string" && data.entryId && typeof data.text === "string") { if (data.identity) { worker.identity = data.identity; worker.sessionFile = data.identity.sessionFile; save(); } - const report: Report = { id: `${event.fromSessionId}:${data.entryId || digest(data.text)}`, plan: state.plan, session: event.fromSessionId, sessionFile: worker.sessionFile!, requestId: data.requestId, text: data.text }; + const report: Report = { id: `${event.fromSessionId}:${data.entryId}`, plan: state.plan, session: event.fromSessionId, sessionFile: worker.sessionFile!, requestId: data.requestId, task: worker.task, text: data.text }; recordReport(ctx, report); } }, @@ -439,25 +447,10 @@ export default function mainSupervisor(pi: ExtensionAPI) { if (entry.type !== "custom" || entry.customType !== STOP) continue; const stopped = entry.data as WorkerStop; if (stopped.to !== worker.parentId || stopped.requestId !== worker.requestId || stopped.plan !== owner.plan) continue; - recordReport(ctx, { id: `${worker.intercomId}:${stopped.entryId}`, plan: stopped.plan, session: worker.intercomId!, sessionFile: worker.sessionFile!, requestId: stopped.requestId, text: stopped.text }, false); + recordReport(ctx, { id: `${worker.intercomId}:${stopped.entryId}`, plan: stopped.plan, session: worker.intercomId!, sessionFile: worker.sessionFile!, requestId: stopped.requestId, task: worker.task, text: stopped.text }, false); } } catch { /* Unavailable saved history is not evidence of a stopped worker. */ } } - const aliases = new Map(); - for (const entry of ctx.sessionManager.getBranch()) { - if (entry.type !== "custom_message" || entry.customType !== "intercom_message") continue; - const received = entry.details as { from?: { id?: string }; message?: { id?: string; timestamp?: number; content?: { text?: string }; retryOf?: string; supersedes?: string } } | undefined; - const sender = received?.from?.id, message = received?.message; - if (!sender || !message?.id || typeof message.content?.text !== "string" || !message.content.text.trim() || /^(?:ok(?:ay)?|thanks|received|acknowledged)[.!]?$/i.test(message.content.text.trim())) continue; - const ownerEntry = ctx.sessionManager.getBranch().filter(saved => saved.type === "custom" && saved.customType === STATE && (!message.timestamp || !saved.timestamp || Date.parse(saved.timestamp) <= message.timestamp) && (saved.data as State).worker?.intercomId === sender).at(-1); - const owner = ownerEntry?.type === "custom" ? ownerEntry.data as State : undefined; - if (!owner?.worker?.requestId || !owner.worker.sessionFile || !owner.plan) continue; - const id = `${sender}:${message.id}`, retry = message.retryOf ? aliases.get(`${sender}:${message.retryOf}`) || `${sender}:${message.retryOf}` : undefined; - if (retry && records(ctx, REPORT).some(report => report.id === retry && report.text === message.content!.text)) { aliases.set(id, retry); continue; } - aliases.set(id, id); - const supersedes = message.supersedes ? aliases.get(`${sender}:${message.supersedes}`) || `${sender}:${message.supersedes}` : undefined; - recordReport(ctx, { id, session: sender, sessionFile: owner.worker.sessionFile, plan: owner.plan, requestId: owner.worker.requestId, text: message.content.text, supersedes }, false); - } if (ctx.isIdle()) remindReports(ctx); } const reportStop = (text: string) => { @@ -568,7 +561,7 @@ export default function mainSupervisor(pi: ExtensionAPI) { lastWorkingSet = foldPlan(snapshot.text); pendingUpkeep = undefined; const pending = !state.child && state.mode === "supervising" ? pendingReports(ctx) : []; - return { systemPrompt: `${event.systemPrompt}\n\n${role}${pending.length ? `\n${pendingReportReviews(pending.map(report => `${report.id} (${report.sessionFile})`))}` : ""}`, ...(message ? { message } : {}) }; + return { systemPrompt: `${event.systemPrompt}\n\n${role}${pending.length ? `\n${pendingReportReviews(pending.map(reportLabel))}` : ""}`, ...(message ? { message } : {}) }; }); pi.on("input", (event, ctx) => { if (event.source !== "extension") { pauseCheckIn = false; return; } @@ -647,7 +640,7 @@ export default function mainSupervisor(pi: ExtensionAPI) { `Plan: ${state.plan ?? "none"}`, `Preferred worker model (plan, not configuration): ${notedPlanValue("preferred worker model") ?? "inherit"}`, `Last observed worker model: ${state.worker?.identity?.model ?? "unconfirmed"}; verify current choice before claiming configuration.`, - `Pending report reviews: ${pendingReports(ctx).map(report => report.id).join(", ") || "none"}`, + `Pending worker revision reviews: ${pendingReports(ctx).map(reportLabel).join("; ") || "none"}`, `Recorded worker session: ${state.worker?.sessionFile ?? "not recorded"}`, `Worker Intercom: ${state.worker?.intercomId ?? "unconfirmed"}; native pane: ${state.worker?.paneId ?? "unconfirmed"}`, notedPlanValue("worker session") ? `Worker session noted in plan: ${notedPlanValue("worker session")}` : "", @@ -770,7 +763,7 @@ export default function mainSupervisor(pi: ExtensionAPI) { if (self.length !== 1) return result(nativeMessages.noIdentity); if (stamp !== generation || opening || state.worker || signal?.aborted) return result(messages.cancelled); const requestId = randomUUID(); - state.worker = { requestId, parentId: self[0].id }; state.workerStopped = false; workerRevision++; opening = true; save(); + state.worker = { requestId, parentId: self[0].id, task: params.task }; state.workerStopped = false; workerRevision++; opening = true; save(); try { // Stock open sends startup only to a newly created context; existing panes receive nothing. const pane = await openProjectPane({ cwd: ctx.cwd, message: workerAssignment(plan, self[0].id, requestId, params.task, model), focus: false, signal }); @@ -814,34 +807,43 @@ export default function mainSupervisor(pi: ExtensionAPI) { pi.registerTool({ name: "review_subagent", label: "Review worker report", description: reportReviewDescription, parameters: Type.Object({ - report: Type.String({ minLength: 1 }), goal: sourceQuote, + reportId: Type.String({ minLength: 1 }), goal: sourceQuote, evidence: Type.Array(Type.Object({ path: Type.String({ minLength: 1 }), entryId: Type.Optional(Type.String()), quote: Type.Optional(Type.String()), observation: Type.String({ minLength: 1 }) }), { minItems: 1 }), observation: Type.String({ minLength: 1 }), unmet: Type.String({ minLength: 1 }), verdict: Type.String({ enum: ["accepted", "changes_requested", "blocked"] }), continuation: Type.Optional(Type.String()), }), async execute(_id, params, signal, _update, ctx) { if (state.child || state.mode !== "supervising") throw new Error("Only the active parent supervisor can review worker reports."); - const report = records(ctx, REPORT).find(report => report.id === params.report); - if (!report) throw new Error("Unknown owned report; inspect pending report reviews in /goals status."); + const report = records(ctx, REPORT).find(report => report.id === params.reportId); + if (!report) throw new Error("Unknown owned report ID; use the reportId shown in /goals status."); if (!["accepted", "changes_requested", "blocked"].includes(params.verdict) || !params.goal.quote.trim() || !params.observation.trim() || !params.unmet.trim() || params.verdict === "changes_requested" && !params.continuation?.trim()) throw new Error("Supply the verdict, inspection, unmet requirements (or none), and a concrete continuation for changes_requested."); const sources = [params.goal, ...params.evidence].map(source => { - const path = resolve(ctx.cwd, source.path), bytes = readFileSync(path); - if (!bytes.length) throw new Error(`Empty evidence: ${path}`); - let text = bytes.toString("utf8"); - if (source.entryId) { - const entry = savedSession(path).getBranch().find(entry => entry.id === source.entryId); - if (!entry) throw new Error(`Missing session entry: ${path}#${source.entryId}`); - const content = entry.type === "message" && "content" in entry.message ? entry.message.content : entry.type === "custom_message" ? entry.content : undefined; - text = typeof content === "string" ? content : Array.isArray(content) ? content.filter(part => part.type === "text").map(part => part.text).join("\n") : JSON.stringify(entry); + const gitEvidence = /^git:([0-9a-f]{7,64}):(.+)$/i.exec(source.path); + let path = source.path, bytes: Buffer, text: string; + if (gitEvidence) { + const [, commit, gitPath] = gitEvidence; + if (source.entryId || !gitPath.split("/").every(part => part && part !== "." && part !== "..")) throw new Error(`Invalid Git evidence: ${source.path}`); + try { bytes = Buffer.from(execFileSync("git", ["show", `${commit}:${gitPath}`], { cwd: ctx.cwd, encoding: "buffer" })); } + catch { throw new Error(`Unavailable Git evidence: ${source.path}`); } + text = bytes.toString("utf8"); + } else { + path = resolve(ctx.cwd, source.path); bytes = readFileSync(path); text = bytes.toString("utf8"); + if (source.entryId) { + const entry = savedSession(path).getBranch().find(entry => entry.id === source.entryId); + if (!entry) throw new Error(`Missing session entry: ${path}#${source.entryId}`); + const content = entry.type === "message" && "content" in entry.message ? entry.message.content : entry.type === "custom_message" ? entry.content : undefined; + text = typeof content === "string" ? content : Array.isArray(content) ? content.filter(part => part.type === "text").map(part => part.text).join("\n") : JSON.stringify(entry); + } } + if (!bytes.length) throw new Error(`Empty evidence: ${path}`); const binary = bytes.includes(0) || /\.(png|jpe?g|gif|webp|pdf|mp4)$/i.test(path); if (!source.quote?.trim() && !binary || source.quote && !text.includes(source.quote)) throw new Error(`Quote does not match source: ${path}`); if ("observation" in source && !source.observation.trim()) throw new Error(`Describe the inspected evidence: ${path}`); return `${path}${source.entryId ? `#${source.entryId}` : ""}\n${source.quote ? `> ${source.quote}` : "[non-text capture]"}${"observation" in source ? `\nObserved: ${source.observation}` : ""}`; }); const content = reportReviewContent(report.id, report.sessionFile, sources, params.observation, params.unmet, params.verdict, params.continuation || ""); - const review: ReportReview = { id: digest(content), report: report.id, verdict: params.verdict, content, continuation: params.continuation || "" }; - if (records(ctx, REVIEW).some(saved => saved.report === report.id)) return result("This report already has a delivered review; a new revision needs its own report."); + const review: ReportReview = { id: digest(content), reportId: report.id, verdict: params.verdict, content, continuation: params.continuation || "" }; + if (records(ctx, REVIEW).some(saved => reviewedReportId(saved) === report.id)) return result("This worker revision already has a delivered review; a later stop report is a new revision."); if (!channel?.snapshot().connected || !channel.snapshot().supported || signal?.aborted) throw new Error("Review delivery unavailable; report remains pending."); const payload = { type: "review", to: report.session, sessionFile: report.sessionFile, requestId: report.requestId, plan: report.plan, review }; if (Buffer.byteLength(JSON.stringify(payload)) > 16000) throw new Error("Review exceeds Intercom's 16 KiB limit; shorten the quotes and retain source references."); diff --git a/src/prompts.ts b/src/prompts.ts index 68c3c8a..c031971 100644 --- a/src/prompts.ts +++ b/src/prompts.ts @@ -191,9 +191,9 @@ export function workerAttachment(plan: string, session: string, text: string): s return `Worker attachment for ${plan}, exact Intercom session ${session}:\n${text}\nMetadata only; no acknowledgement or review turn requested.`; } // Pi/OpenAI: supervisor-authored report reviews, separate from goal completion. -export const reportReviewDescription = "Review an owned worker report after inspecting its actual artifacts. Quote the assigned goal/task and evidence from files (optional saved-session entryId selects decoded message text). State observations and unmet requirements; use accepted, changes_requested or blocked. Changes requested need a concrete continuation. Text quotes are checked, not their relevance or quality. Non-text evidence needs a nonempty capture and specific observation. Delivery stays pending until the worker saves the visible review. Acceptance never completes a goal or wakes/closes the worker."; +export const reportReviewDescription = "Review one owned worker revision after inspecting its actual artifacts. reportId is the exact session:revision token shown in /goals status, not report prose. Quote the assigned goal/task and evidence from files; git:: reads an immutable tracked revision. An optional saved-session entryId selects decoded message text. State observations and unmet requirements; use accepted, changes_requested or blocked. Changes requested need a concrete continuation. Text quotes are checked, not their relevance or quality. Non-text evidence needs a nonempty capture and specific observation. Delivery stays pending until the worker saves the visible review. Acceptance never completes a goal or wakes/closes the worker."; export const reportReviewContent = (report: string, sessionFile: string, sources: string[], observation: string, unmet: string, verdict: string, continuation: string) => `Worker review: ${verdict}\nReport: ${report}\nSaved session: ${sessionFile}\n\nAssigned goal/task:\n${sources[0]}\n\nEvidence:\n${sources.slice(1).join("\n\n")}\n\nInspected: ${observation}\nUnmet: ${unmet}\nContinuation: ${continuation || "none"}\nThis is a report review, not CompleteGoal.\n— Pi supervisor`; -export const pendingReportReviews = (reports: string[]) => `Pending worker reviews: ${reports.join(", ")}. Inspect their saved reports and actual artifacts, then use review_subagent. Independent authorized work may continue; receipts and generic replies do not resolve reviews.`; +export const pendingReportReviews = (reports: string[]) => `Pending worker revision reviews: ${reports.join("; ")}. Inspect each saved stop report and actual artifacts, then use review_subagent with its reportId. Independent authorized work may continue; attachment receipts and ordinary Intercom messages are not review obligations.`; export function workerReview(plan: string, session: string, text: string): string { return `Worker event for ${plan}, exact Intercom session ${session}:\n${text}\nThis is a report, not completion approval. Inspect actual artifacts and saved messages; if correction is needed, send it to the same session. Preserve its visible review conversation. Respect pauses; do not reply merely to acknowledge.`; } diff --git a/test/goals.test.ts b/test/goals.test.ts index b5b52ea..6d620da 100644 --- a/test/goals.test.ts +++ b/test/goals.test.ts @@ -1,3 +1,4 @@ +import { execFileSync } from "node:child_process"; import { existsSync, mkdirSync, mkdtempSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync } from "node:fs"; import { access, readFile, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; @@ -1094,23 +1095,25 @@ it("opens no-focus, records explicit attachment only, and wakes review only for expect(worker).toMatchObject({ paneId: "native-pane", intercomId: "worker-id", sessionFile: "/tmp/native-worker.jsonl" }); expect(f.messages.at(-1)).toMatchObject({ message: { customType: "pi-goals-supervision", display: true, content: expect.stringContaining("Metadata only; no acknowledgement or review turn requested") }, options: { triggerTurn: false } }); expect(f.messages.at(-1).savedPrompt).toBeUndefined(); - const notice = { type: "stopped", to: worker.parentId, requestId: worker.requestId, plan: f.path, text: "Blocked: input missing" }; + const notice = { type: "stopped", to: worker.parentId, requestId: worker.requestId, plan: f.path, entryId: "revision-1", text: "Blocked: input missing" }; const count = f.messages.length; for (const fromSessionId of [worker.parentId, "foreign-id"]) f.event({ type: "message", fromSessionId, payload: notice }); f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, plan: "/foreign.md" } }); f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, requestId: "stale" } }); + f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, entryId: undefined } }); expect(f.messages).toHaveLength(count); - for (const text of ["Blocked: input missing", "Done: output.txt", "Error: execution failed"]) { - f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, text } }); - expect(f.messages.at(-2)?.message.content).toContain(text); - expect(f.messages.at(-1)?.message.content).toContain("Pending worker reviews:"); - } + f.event({ type: "message", fromSessionId: "worker-id", payload: notice }); + expect(f.messages.at(-2)?.message.content).toContain("Blocked: input missing"); + expect(f.messages.at(-1)?.message.content).toContain("Pending worker revision reviews:"); + const afterFirstRevision = f.messages.length; + for (const text of ["Done: output.txt", "Error: execution failed"]) f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, text } }); + expect(f.messages).toHaveLength(afterFirstRevision); expect(readFileSync(f.path, "utf8")).not.toContain("[āœ“]"); await f.command("stop"); - f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, text: "New report during pause" } }); + f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, entryId: "revision-2", text: "New report during pause" } }); expect(f.messages.at(-1).options).toEqual({ deliverAs: "nextTurn" }); await f.command("clear"); - const cleared = f.messages.length; f.event({ type: "message", fromSessionId: "worker-id", payload: notice }); + const cleared = f.messages.length; f.event({ type: "message", fromSessionId: "worker-id", payload: { ...notice, entryId: "revision-3" } }); expect(f.messages).toHaveLength(cleared); }); @@ -1184,9 +1187,9 @@ it.each(["inherit", "plan", "explicit"])("hands off %s model policy without clai await f.command("status"); expect(f.ctx.ui.notify).toHaveBeenLastCalledWith(expect.stringContaining("Last observed worker model: offline/inherited"), "info"); if (model) { - f.event({ type: "message", fromSessionId: "worker", payload: { type: "stopped", to: worker.parentId, requestId: worker.requestId, plan: f.path, text: "Requested missing/unavailable is unavailable; unrelated work can continue." } }); + f.event({ type: "message", fromSessionId: "worker", payload: { type: "stopped", to: worker.parentId, requestId: worker.requestId, plan: f.path, entryId: "model-unavailable", text: "Requested missing/unavailable is unavailable; unrelated work can continue." } }); expect(f.messages.at(-2)?.message.content).toContain("missing/unavailable is unavailable"); - expect(f.messages.at(-1)?.message.content).toContain("Pending worker reviews:"); + expect(f.messages.at(-1)?.message.content).toContain("Pending worker revision reviews:"); } }); @@ -1220,7 +1223,12 @@ it("reviews a saved worker revision through inspection, silent delivery, retry a parent.hooks.get("session_start")({}, parent.ctx); expect(parent.hooks.get("before_agent_start")({ systemPrompt: "base" }, parent.ctx).systemPrompt).toContain(id); const artifact = join(parent.ctx.cwd, "output.txt"); writeFileSync(artifact, "first output\n"); - const form = { report: id, goal: { path: parent.path, quote: "goal: first output" }, evidence: [{ path: artifact, quote: "invented", observation: "Read output.txt" }], observation: "Inspected actual output and assigned goal", unmet: "none", verdict: "accepted" }; + execFileSync("git", ["init"], { cwd: parent.ctx.cwd }); + execFileSync("git", ["add", "output.txt"], { cwd: parent.ctx.cwd }); + execFileSync("git", ["-c", "user.name=Test", "-c", "user.email=test@example.com", "commit", "-m", "worker evidence"], { cwd: parent.ctx.cwd }); + const commit = execFileSync("git", ["rev-parse", "HEAD"], { cwd: parent.ctx.cwd, encoding: "utf8" }).trim(); + writeFileSync(artifact, "changed after reported revision\n"); + const form = { reportId: id, goal: { path: parent.path, quote: "goal: first output" }, evidence: [{ path: `git:${commit}:output.txt`, quote: "invented", observation: "Read the reported revision" }], observation: "Inspected actual output and assigned goal", unmet: "none", verdict: "accepted" }; const review = () => parent.tools.get("review_subagent").execute("review", form, undefined, undefined, parent.ctx); await expect(review()).rejects.toThrow("Quote does not match"); form.evidence[0].quote = "first output"; @@ -1228,17 +1236,18 @@ it("reviews a saved worker revision through inspection, silent delivery, retry a parent.channel.publish.mockImplementationOnce(() => {}); // Publish success is not saved delivery. await review(); await parent.command("status"); expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain(id); + expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Implement first output"); await review(); await parent.command("status"); - expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Pending report reviews: none"); + expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Pending worker revision reviews: none"); expect(worker.messages).toHaveLength(workerTurns); expect(readFileSync(sm.getSessionFile()!, "utf8")).toContain("Worker review: accepted"); await review(); worker.hooks.get("session_shutdown")(); parent.event({ type: "session_left", sessionId: workerId }); await parent.command("status"); - expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Pending report reviews: none"); + expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Pending worker revision reviews: none"); worker.hooks.get("session_start")({}, worker.ctx); await parent.command("stop"); - form.report = report("Revision failed", "error"); + form.reportId = report("Revision failed", "error"); const paused = parent.messages.length; await parent.hooks.get("agent_settled")({}, parent.ctx); expect(parent.messages).toHaveLength(paused); @@ -1247,9 +1256,9 @@ it("reviews a saved worker revision through inspection, silent delivery, retry a await expect(review()).rejects.toThrow("concrete continuation"); await parent.tools.get("review_subagent").execute("revision", { ...form, unmet: "Output still needs correction", continuation: "Correct output.txt and rerun verification." }, undefined, undefined, parent.ctx); expect(worker.messages.at(-1)).toMatchObject({ savedPrompt: true, message: { content: expect.stringContaining("Correct output.txt") } }); - form.report = report("Cancelled while correcting", "aborted"); form.verdict = "blocked"; + form.reportId = report("Cancelled while correcting", "aborted"); form.verdict = "blocked"; await review(); await parent.command("status"); - expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Pending report reviews: none"); + expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain("Pending worker revision reviews: none"); expect(readFileSync(parent.path, "utf8")).not.toContain("[āœ“]"); // Lost notification and cancellation before any new assistant message: durable run identity. @@ -1257,7 +1266,7 @@ it("reviews a saved worker revision through inspection, silent delivery, retry a worker.channel.publish.mockImplementationOnce(() => {}); worker.hooks.get("agent_end")({ messages: [] }, worker.ctx); const missed = `${workerId}:${worker.channel.publish.mock.lastCall?.[0].entryId}`; - expect(missed).not.toBe(form.report); + expect(missed).not.toBe(form.reportId); // Same preserved worker, newly approved plan and request: old reviews stay in history. const history = sm.getBranch(), oldPlan = readFileSync(parent.path, "utf8"); await parent.command("clear"); await parent.command("new Next output"); @@ -1283,7 +1292,7 @@ it("reviews a saved worker revision through inspection, silent delivery, retry a parent.hooks.get("session_start")({}, parent.ctx); // Reconcile the missed stop from real saved worker history. await parent.command("status"); expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain(missed); - form.report = missed; await review(); // Old-plan blocked review remains deliverable after retargeting. + form.reportId = missed; await review(); // Old-plan blocked review remains deliverable after retargeting. await parent.command("status"); expect(parent.ctx.ui.notify.mock.lastCall?.[0]).not.toContain(missed); expect(parent.ctx.ui.notify.mock.lastCall?.[0]).toContain(nextReport); @@ -1299,11 +1308,8 @@ it("reviews a saved worker revision through inspection, silent delivery, retry a ordinary("ack-only", "OK"); ordinary("foreign", "Unowned report", {}, "foreign-peer"); parent.hooks.get("session_start")({}, parent.ctx); await parent.command("status"); const pending = parent.ctx.ui.notify.mock.lastCall?.[0]; - expect(pending).toContain(`${workerId}:ordinary-b`); - for (const excluded of ["ordinary-a", "retry-b", "retry-again", "ack-only", "foreign-peer"]) expect(pending).not.toContain(excluded); - form.report = `${workerId}:ordinary-b`; await review(); - await parent.command("status"); - expect(parent.ctx.ui.notify.mock.lastCall?.[0]).not.toContain("ordinary-b"); + for (const nonReviewable of ["ordinary-a", "ordinary-b", "retry-b", "retry-again", "ack-only", "foreign-peer"]) expect(pending).not.toContain(nonReviewable); + expect(pending).toContain(nextReport); expect(readFileSync(parent.path, "utf8")).toBe(oldPlan); expect((await worker.tools.get("CompleteGoal").execute("deny", { goal: "first output", evidence: [], observation: "claim" }, undefined, undefined, worker.ctx)).content[0].text).toContain("only to the active parent"); });