diff --git a/apps/server/src/agent/planner.ts b/apps/server/src/agent/planner.ts index bb2aba7d..f955cdbf 100644 --- a/apps/server/src/agent/planner.ts +++ b/apps/server/src/agent/planner.ts @@ -249,11 +249,18 @@ and cannot change GitHub. Ground the plan in what those reading tools return.`; ? `Your working directory, ${workspace.cwd}, is a local checkout of ${repository} verified against its origin. Its branch and working tree may differ from what the repository tools read.` - : `Your working directory, ${workspace.cwd}, is a scratch directory Chopin created -empty for this document. It is not a checkout and holds no repository files, so read + : `Your working directory, ${workspace.cwd}, is a scratch directory Chopin keeps for +this document. It is not a checkout and holds no repository files, so read ${repository} through the repository tools.`; let questions = `\`ask_user_question\` and \`workflow\` questions appear to the document's members as Decisions. If one expires unanswered, proceed on your best judgement and say what you assumed.`; - return [PROMPT, reading, place, questions, bootstrap].filter(Boolean).join("\n\n"); + let surface = `The document's members use Chopin in a browser. They cannot run slash commands, +terminal commands, or Atomic CLI commands, so never tell them to use \`/workflow connect\`, +\`/workflow status\`, \`/tasks\`, \`atomic\`, or anything else typed into a terminal, even when a +tool result suggests it. When you start a workflow, say in plain words what it will do and +where to follow it: Chat shows a card for each run with its stages and status, its questions +appear under Decisions, **Stop Planner** pauses it, and **Resume Planner** resumes it. Offer to +check on or steer a run yourself with your \`workflow\` and \`intercom\` tools when someone asks.`; + return [PROMPT, reading, place, questions, surface, bootstrap].filter(Boolean).join("\n\n"); } diff --git a/apps/server/src/chat/invoke.test.ts b/apps/server/src/chat/invoke.test.ts index a20c02a5..71549cbd 100644 --- a/apps/server/src/chat/invoke.test.ts +++ b/apps/server/src/chat/invoke.test.ts @@ -1,4 +1,4 @@ -import { afterEach, expect, test } from "bun:test"; +import { afterEach, beforeEach, expect, test } from "bun:test"; import { mkdtemp, readdir, realpath, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -7,7 +7,7 @@ import { ActiveOwnerBindings } from "../agent/active-owner"; import { Admission } from "../auth/admission"; import { Sessions } from "../auth/session"; import { fullPlanner } from "../harness/atomic/full"; -import { removeWorkspaces } from "../harness/atomic/workspace"; +import { forgetWorkspaces } from "../harness/atomic/workspace"; import { openPlannerSession } from "../harness/session"; import * as Plan from "../plan/service"; import { MemoryStorage } from "../storage/memory/adapter"; @@ -15,6 +15,8 @@ import { configured, LOCAL } from "../testing/config"; import * as Chat from "./service"; import type { Server } from "bun"; +import type { ActiveOwnerBinding } from "../agent/active-owner"; +import type { Chat as Wire } from "@chopin/protocol"; import type { HostedAuth } from "../auth/routes"; import type { Config } from "../config"; import type { GitHub } from "../github/client"; @@ -23,9 +25,17 @@ import type { Socket, SocketData } from "../wire"; const ATOMIC = { HARNESS: "atomic", HARNESS_AUTH: "ai-gateway", MODEL: "stub/model" }; let cleanups: Array<() => Promise> = []; +let previousState = process.env.XDG_STATE_HOME; +beforeEach(async () => { + let state = await mkdtemp(join(tmpdir(), "chopin-invoke-state-")); + process.env.XDG_STATE_HOME = state; + cleanups.push(() => rm(state, { recursive: true, force: true })); +}); afterEach(async () => { for (let cleanup of cleanups.splice(0)) await cleanup(); - await removeWorkspaces(); + forgetWorkspaces(); + if (previousState === undefined) delete process.env.XDG_STATE_HOME; + else process.env.XDG_STATE_HOME = previousState; }); function grant(accessToken: string) { @@ -409,3 +419,208 @@ test("HARNESS=atomic runs every Planner session full in hosted and local configu } } }); + +type Runs = { active: string[]; paused: string[]; cards: Wire.Run[] }; + +function card(status: Wire.Run["status"]): Wire.Run { + let ended = status === "finished" || status === "failed" || status === "stopped"; + return { + id: "run-1", + name: "plan-review", + status, + started: 1_000, + updated: 1_600, + ...(ended ? { ended: 1_600 } : {}), + stages: [{ + id: "run-1:draft", + name: "draft-1", + status: status === "waiting" ? "awaiting_input" : ended ? "completed" : "running", + started: 1_000, + }], + waiting: status === "waiting" ? 1 : 0, + }; +} + +/** A Planner session that owns workflow runs, driven by the test. */ +function runningSession() { + let runs: Runs = { active: ["run-1"], paused: [], cards: [card("running")] }; + let listeners = new Set<(runs: Runs) => void>(); + let calls: string[] = []; + let session = { + async stream(prompt: string) { + calls.push(`stream ${prompt}`); + return { + fullStream: (async function*() { + yield { type: "finish" }; + })(), + } as never; + }, + async destroy() { + calls.push("destroy"); + }, + runs: () => runs, + watchRuns(listener: (runs: Runs) => void) { + listeners.add(listener); + return () => listeners.delete(listener); + }, + async pauseRuns() { + calls.push("pause"); + }, + async resumeRuns() { + calls.push("resume"); + }, + async pauseRun(runId: string) { + calls.push(`pause ${runId}`); + }, + async resumeRun(runId: string) { + calls.push(`resume ${runId}`); + }, + }; + return { + session, + calls, + set(next: Runs) { + runs = next; + for (let listener of listeners) listener(next); + }, + }; +} + +test("a Planner that still owns workflow runs outlives its turn, pauses, resumes, and is let go when they finish", async () => { + let { context, events, user } = await setup(configured(ATOMIC)); + let planner = runningSession(); + let opened = 0; + context.openPlannerSession = async () => { + opened++; + return { ok: true, value: planner.session }; + }; + let holds = 0; + context.hold = () => { + holds++; + return () => holds--; + }; + let ws = { data: { handle: "ana" } } as unknown as Socket; + let said = () => context.chat.entries.at(-1)?.text; + let runs = () => (events.filter(event => event.kind === "chat:state").at(-1)?.runs as + | Wire.Run[] + | undefined); + + expect(await Chat.invoke(context, user, "Run the workflow")).toBeUndefined(); + await context.chat.running; + expect(planner.calls).toEqual(["stream @ana: Run the workflow"]); + expect(context.chat.busy).toBe(false); + expect(context.chat.runs).toEqual([card("running")]); + expect(runs()).toEqual([card("running")]); + expect(holds).toBe(1); + + await Chat.abort(context, ws); + expect(planner.calls.at(-1)).toBe("pause"); + expect(said()).toBe("@ana stopped the Planner and paused its workflows."); + planner.set({ active: [], paused: ["run-1"], cards: [card("paused")] }); + expect(runs()?.map(run => run.status)).toEqual(["paused"]); + expect(planner.calls).not.toContain("destroy"); + + await Chat.resume(context, ws); + expect(planner.calls.at(-1)).toBe("resume"); + expect(said()).toBe("@ana resumed the Planner's workflows."); + planner.set({ active: ["run-1"], paused: [], cards: [card("waiting")] }); + expect(runs()?.[0]).toMatchObject({ + status: "waiting", + waiting: 1, + stages: [{ status: "awaiting_input" }], + }); + + expect(await Chat.invoke(context, user, "How is it going?")).toBeUndefined(); + await context.chat.running; + expect(opened).toBe(1); + expect(planner.calls.filter(call => call.startsWith("stream"))).toHaveLength(2); + expect(planner.calls).not.toContain("destroy"); + + planner.set({ active: [], paused: [], cards: [card("finished")] }); + expect(context.chat.entries.some(entry => entry.text === "plan-review finished after 10 min.")) + .toBe(true); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(planner.calls.at(-1)).toBe("destroy"); + expect(context.chat.retained).toBeUndefined(); + expect(context.chat.runs).toEqual([card("finished")]); + expect(runs()?.map(run => run.status)).toEqual(["finished"]); + expect(holds).toBe(0); + + await Chat.resume(context, ws); + await Chat.abort(context, ws); + expect(planner.calls.at(-1)).toBe("destroy"); +}); + +test("a retained Planner is let go when its owner's binding ends, and sessions without runs still end with their turn", async () => { + let { context, user } = await setup(configured(ATOMIC)); + let planner = runningSession(); + let bindings: ActiveOwnerBinding[] = []; + let resolve = context.activeOwner!; + context.activeOwner = async () => { + let binding = await resolve(); + if (binding) bindings.push(binding); + return binding; + }; + context.openPlannerSession = async () => ({ ok: true, value: planner.session }); + expect(await Chat.invoke(context, user, "Run the workflow")).toBeUndefined(); + await context.chat.running; + expect(context.chat.retained).toBeDefined(); + bindings[0]!.release(); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(context.chat.retained).toBeUndefined(); + expect(planner.calls.at(-1)).toBe("destroy"); + + let quiet = { + ...runningSession().session, + runs: (): Runs => ({ active: [], paused: [], cards: [] }), + }; + let destroyed = 0; + quiet.destroy = async () => { + destroyed++; + }; + context.openPlannerSession = async () => ({ ok: true, value: quiet }); + expect(await Chat.invoke(context, user, "Just answer")).toBeUndefined(); + await context.chat.running; + expect(destroyed).toBe(1); + expect(context.chat.retained).toBeUndefined(); +}); + +test("run cards keep ended runs until a new run starts, and a released session's live runs become stopped", () => { + let live = card("running"); + let done = { ...card("finished"), id: "run-0" }; + expect(Chat.mergeRuns([done], [live], 2_000)).toEqual([live]); + expect(Chat.mergeRuns([done, live], [card("waiting")], 2_000)).toEqual([done, card("waiting")]); + let released = Chat.mergeRuns([done, card("waiting")], undefined, 2_000); + expect(released.map(run => [run.id, run.status, run.waiting])).toEqual([ + ["run-0", "finished", 0], + ["run-1", "stopped", 0], + ]); + expect(released[1]?.ended).toBe(2_000); + expect(Chat.mergeRuns([{ ...card("paused"), ended: undefined }], undefined, 2_000)[0]) + .toMatchObject({ status: "stopped", ended: 2_000 }); + expect(Chat.restoreRuns([])).toBeUndefined(); +}); + +test("one run can be paused and resumed by name while others keep going", async () => { + let { context, user } = await setup(configured(ATOMIC)); + let planner = runningSession(); + context.openPlannerSession = async () => ({ ok: true, value: planner.session }); + let ws = { data: { handle: "ana" } } as unknown as Socket; + expect(await Chat.invoke(context, user, "Run the workflow")).toBeUndefined(); + await context.chat.running; + + await Chat.controlRun(context, ws, { kind: "chat:resume-run", runId: "run-1" }); + await Chat.controlRun(context, ws, { kind: "chat:pause-run", runId: "unknown" }); + await Chat.controlRun(context, ws, { kind: "chat:pause-run", runId: 7 }); + expect(planner.calls.filter(call => call.includes("run-1") || call.includes("unknown"))).toEqual( + [], + ); + + await Chat.controlRun(context, ws, { kind: "chat:pause-run", runId: "run-1" }); + expect(planner.calls.at(-1)).toBe("pause run-1"); + expect(context.chat.entries.at(-1)?.text).toBe("@ana paused plan-review."); + planner.set({ active: [], paused: ["run-1"], cards: [card("paused")] }); + await Chat.controlRun(context, ws, { kind: "chat:resume-run", runId: "run-1" }); + expect(planner.calls.at(-1)).toBe("resume run-1"); + expect(context.chat.entries.at(-1)?.text).toBe("@ana resumed plan-review."); +}); diff --git a/apps/server/src/chat/service.ts b/apps/server/src/chat/service.ts index c88749af..38487d5b 100644 --- a/apps/server/src/chat/service.ts +++ b/apps/server/src/chat/service.ts @@ -116,8 +116,12 @@ export type Chat = { pendingSends: number; /** Per-turn owner for revocation and credential rotation fencing. */ openingOwner?: { sessionId: string; generation: number; revision: number }; - /** Current disposable session; never reused by the next turn. */ + /** Current session. Disposable per turn, except a full Planner session kept while it owns workflow runs. */ agent?: PlannerSession; + /** The Planner session kept past its turn, with its owner binding, until its workflow runs finish. */ + retained?: Retained; + /** Live and paused workflow runs of the retained session. */ + runs?: Wire.Runs; turnController?: AbortController; /** Fences a turn invalidated while opening or streaming. */ lifecycle: number; @@ -179,14 +183,17 @@ export function create(): Chat { } /** - * Come back with what was said before. + * Come back with what was said before, and the workflow runs shown then. * * An entry that was still streaming when the process went away never finished, * so the flag is cleared: it is as complete as it is ever going to be, and - * leaving it set would show a spinner nothing will ever stop. + * leaving it set would show a spinner nothing will ever stop. Runs still live + * then have no session now, so they come back stopped. */ -export function restore(entries: Wire.Entry[]): Chat { +export function restore(entries: Wire.Entry[], runs?: Wire.Run[]): Chat { + let restoredRuns = restoreRuns(runs); return { + ...(restoredRuns ? { runs: restoredRuns } : {}), entries: entries.map(entry => { let { streaming: _streaming, ...rest } = entry; return rest; @@ -343,6 +350,7 @@ function state(chat: Chat, server: Server, room: string): void { ts: 0, busy: chat.busy, ...(chat.turn ? { turn: chat.turn } : {}), + ...(chat.runs ? { runs: chat.runs } : {}), }); } @@ -395,6 +403,7 @@ export function greet(chat: Chat, ws: Socket): void { entries: chat.entries.map(publicEntry), busy: chat.busy, ...(chat.turn ? { turn: chat.turn } : {}), + ...(chat.runs ? { runs: chat.runs } : {}), queued: visible(chat), }); } @@ -427,6 +436,14 @@ export type Room = { text: string; question: string; }) => Promise; + /** Keeps the room loaded while a retained Planner still owns workflow runs; returns the release. */ + hold?: () => () => void; +}; + +type Retained = { + session: PlannerSession; + binding: ActiveOwnerBinding; + release: () => Promise; }; /** @@ -774,19 +791,189 @@ export function unqueue(context: Room, ws: Socket, msg: Request): queued(chat, server, room); } -/** Stop the running turn. Anyone may, and the transcript says who did. */ +/** + * Stop the running turn and pause the Planner's workflow runs, so both hold + * until someone resumes them. Anyone may, and the transcript says who did. + */ export async function abort(context: Room, ws: Socket): Promise { let { chat, room, server } = context; - if (!chat.busy || !chat.turnController) return; - chat.turnController.abort(); + let turn = chat.busy ? chat.turnController : undefined; + let session = chat.retained?.session ?? chat.agent; + let live = (session?.runs?.()?.active.length ?? 0) > 0; + if (!turn && !live) return; + turn?.abort(); + let paused = false; + if (live && session?.pauseRuns) { + try { + await session.pauseRuns(); + paused = true; + } catch (err) { + console.error("[chat] pausing workflow runs failed:", err); + } + } say(chat, server, room, { id: ulid(), author: { kind: "system" }, - text: `@${ws.data.handle} stopped the turn.`, + text: paused + ? `@${ws.data.handle} stopped the Planner and paused its workflows.` + : `@${ws.data.handle} stopped the turn.`, ts: now(), }); } +/** Resume the workflow runs the Planner paused. Anyone may, and the transcript says who did. */ +export async function resume(context: Room, ws: Socket): Promise { + let { chat, room, server } = context; + let session = chat.retained?.session; + if (!session?.resumeRuns || !session.runs?.()?.paused.length) return; + let text = `@${ws.data.handle} resumed the Planner's workflows.`; + try { + await session.resumeRuns(); + } catch (err) { + console.error("[chat] resuming workflow runs failed:", err); + text = "The Planner's workflows could not be resumed."; + } + say(chat, server, room, { id: ulid(), author: { kind: "system" }, text, ts: now() }); +} + +/** Pause or resume one workflow run of the retained Planner. Anyone may, and the transcript says who did. */ +export async function controlRun( + context: Room, + ws: Socket, + frame: { kind: "chat:pause-run" | "chat:resume-run"; runId?: unknown }, +): Promise { + let { chat, room, server } = context; + let session = chat.retained?.session; + let run = typeof frame.runId === "string" + ? chat.runs?.find(candidate => candidate.id === frame.runId) + : undefined; + let runs = session?.runs?.(); + if (!session || !run) return; + let pausing = frame.kind === "chat:pause-run"; + if (pausing ? !runs?.active.includes(run.id) : !runs?.paused.includes(run.id)) return; + let control = pausing ? session.pauseRun : session.resumeRun; + if (!control) return; + let text = `@${ws.data.handle} ${pausing ? "paused" : "resumed"} ${run.name}.`; + try { + await control(run.id); + } catch (err) { + console.error(`[chat] ${pausing ? "pausing" : "resuming"} workflow run failed:`, err); + text = `${run.name} could not be ${pausing ? "paused" : "resumed"}.`; + } + say(chat, server, room, { id: ulid(), author: { kind: "system" }, text, ts: now() }); +} + +const ENDED_RUN: Partial> = { + finished: "finished", + blocked: "ended blocked", + failed: "failed", + stopped: "was stopped", +}; + +/** + * The run cards Chat shows after the retained session reports `incoming`. Ended + * cards stay until a new run starts in the document. With no report, the session + * was let go, so a card still marked live becomes stopped. + */ +export function mergeRuns( + previous: Wire.Run[], + incoming: Wire.Run[] | undefined, + at: number, +): Wire.Run[] { + if (!incoming) return previous.map(run => ENDED_RUN[run.status] ? run : stopped(run, at)); + let known = new Set(previous.map(run => run.id)); + if (incoming.some(run => !known.has(run.id))) return incoming; + let reported = new Set(incoming.map(run => run.id)); + return [...previous.filter(run => ENDED_RUN[run.status] && !reported.has(run.id)), ...incoming]; +} + +function stopped(run: Wire.Run, at: number): Wire.Run { + return { ...run, status: "stopped", ended: run.ended ?? at, waiting: 0 }; +} + +/** Run cards restored with the document; no session survives a reload of the room, so live ones were stopped. */ +export function restoreRuns(runs: Wire.Run[] | undefined): Wire.Run[] | undefined { + return runs?.length ? mergeRuns(runs, undefined, now()) : undefined; +} + +/** + * Show the session's runs, keeping ended ones as summary cards until the next + * run starts, store them with the document so a reload keeps them, and say once + * in the transcript when one ends. + */ +function publishRuns( + context: Room, + runs: { active: string[]; paused: string[]; cards?: Wire.Run[] } | undefined, +): void { + let { chat, room, server } = context; + let before = new Map((chat.runs ?? []).map(run => [run.id, run.status])); + let cards = mergeRuns(chat.runs ?? [], runs?.cards, now()); + for (let run of cards) { + let ended = ENDED_RUN[run.status]; + let previous = before.get(run.id); + if (!ended || !previous || ENDED_RUN[previous]) continue; + let minutes = Math.max(1, Math.round(((run.ended ?? run.updated) - run.started) / 60)); + say(chat, server, room, { + id: ulid(), + author: { kind: "system" }, + text: `${run.name} ${ended} after ${minutes} min.`, + ts: now(), + }); + } + chat.runs = cards.length ? cards : undefined; + state(chat, server, room); + context.persist().catch(err => console.error("[chat] storing workflow runs failed:", err)); +} + +/** + * Keep a Planner session that still owns workflow runs instead of destroying + * it with its turn: the runs belong to the session, and the next turn reuses + * it so the Planner can still see and steer them. It is let go once every run + * has finished, when its owner binding ends, or when the chat closes. + */ +function retain( + context: Room, + opened: { session: PlannerSession; binding: ActiveOwnerBinding }, +): boolean { + let { chat } = context; + let runs = opened.session.runs?.(); + if (chat.retained && chat.retained.session !== opened.session) return false; + if (!runs || !(runs.active.length || runs.paused.length) || opened.binding.signal.aborted) { + return false; + } + if (!chat.retained) { + let unhold = context.hold?.(); + let stopWatching = opened.session.watchRuns?.(next => { + publishRuns(context, next); + if (!chat.busy && !next.active.length && !next.paused.length) void chat.retained?.release(); + }); + let released: Promise | undefined; + let retained: Retained = { + ...opened, + release: () => + released ??= (async () => { + if (chat.retained === retained) chat.retained = undefined; + stopWatching?.(); + opened.binding.signal.removeEventListener("abort", ended); + try { + await opened.session.destroy(); + } finally { + opened.binding.release(); + unhold?.(); + if (chat.agent === opened.session) chat.agent = undefined; + if (!chat.busy) chat.owner = undefined; + publishRuns(context, undefined); + } + })(), + }; + let ended = () => void retained.release(); + opened.binding.signal.addEventListener("abort", ended, { once: true }); + chat.retained = retained; + } + publishRuns(context, runs); + return true; +} + function currentMemberRequest(chat: Chat): ActiveMemberRequest | undefined { let active = chat.activeRequest; if ( @@ -938,6 +1125,14 @@ async function repositorySession( currentEntryId?: string, currentReferences: Wire.Reference[] = [], ): Promise<{ session: PlannerSession; binding: ActiveOwnerBinding }> { + let retained = context.chat.retained; + if (retained) { + if (!retained.binding.signal.aborted && await retained.binding.revalidate()) { + context.chat.agent = retained.session; + return { session: retained.session, binding: retained.binding }; + } + await retained.release(); + } let { ownership, owner, repository } = await resolveOwner( context.auth, context.repository, @@ -1153,12 +1348,16 @@ async function run( } finally { chat.activeRequest = undefined; chat.turnController = undefined; - try { - await opened?.session.destroy(); - } finally { - opened?.binding.release(); - if (chat.agent === opened?.session) chat.agent = undefined; - chat.owner = undefined; + let kept = !!opened && !chat.closed && retain(context, opened); + if (!kept) { + try { + if (chat.retained?.session === opened?.session) await chat.retained?.release(); + else await opened?.session.destroy(); + } finally { + opened?.binding.release(); + if (chat.agent === opened?.session) chat.agent = undefined; + chat.owner = undefined; + } } chat.messageIds = undefined; chat.writing = undefined; @@ -1444,4 +1643,5 @@ export async function close(chat: Chat): Promise { clearTimeout(chat.lingering); chat.lingering = undefined; await Promise.all([chat.sending, chat.running]); + await chat.retained?.release(); } diff --git a/apps/server/src/harness/atomic/adapter.ts b/apps/server/src/harness/atomic/adapter.ts index 81f1f31f..c19d5840 100644 --- a/apps/server/src/harness/atomic/adapter.ts +++ b/apps/server/src/harness/atomic/adapter.ts @@ -23,7 +23,7 @@ import { import { mkdtemp, readFile, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { fullPlanner } from "./full"; +import { fullPlanner, workflowRuns } from "./full"; import type { HarnessV1, @@ -484,7 +484,10 @@ export function createAtomicAdapter( systemPrompt: ATOMIC_DEFAULT_SYSTEM_PROMPT, appendSystemPrompt: [], }), - extensionFactories: [atomicTurnExtension(policy)], + extensionFactories: [ + atomicTurnExtension(policy), + ...(full ? [workflowRuns(full)] : []), + ], }); await loader.reload(); let created = await createAgentSession({ @@ -502,6 +505,9 @@ export function createAtomicAdapter( customTools: turn.tools.map(hostTool), }); let session = created.session; + if (full) { + full.workflows = session.workflows; + } let leak = full ? undefined : hostLeak(session, created.extensionsResult.extensions.length, hostNames); diff --git a/apps/server/src/harness/atomic/full.test.ts b/apps/server/src/harness/atomic/full.test.ts index dd31ba08..8a755590 100644 --- a/apps/server/src/harness/atomic/full.test.ts +++ b/apps/server/src/harness/atomic/full.test.ts @@ -10,10 +10,18 @@ import { ATOMIC_RESULT_TOOL_NAME, createAtomicAdapter, } from "./adapter"; -import { registerFullPlanner } from "./full"; +import { + classifyRuns, + type FullPlanner, + pauseOwnedRuns, + registerFullPlanner, + resumeOwnedRuns, + runCards, + untilUnpaused, +} from "./full"; import { startStubModelServer } from "../pi/model-stub"; import { hostInputRoom } from "../../testing/decisions"; -import type { HostInput, QuestionParams } from "@bastani/atomic"; +import type { HostInput, QuestionParams, SessionWorkflows } from "@bastani/atomic"; let params: QuestionParams = { questions: [{ @@ -306,3 +314,281 @@ test("worker sessions stay isolated, even beside a full Planner session on the s expect(beside.workerRequests[0]!.toolNames).toEqual(["host_tool"]); expect(beside.workerRequests[0]!.system).toBe("CHOPIN-INSTRUCTIONS-MARKER"); }); + +test("paused roots are paused, roots whose run ended are neither, and every other root is live", () => { + let root = (rootRunId: string, state: "working" | "idle" | "blocked", reason: string) => + ({ + rootRunId, + ownerSessionId: "session", + state, + reason, + activeExecutionCount: 0, + actionableBlockCount: 0, + needsAttention: false, + }) as Parameters[0] extends Iterable ? T : never; + expect(classifyRuns([ + root("drafting", "working", "executing"), + root("asking", "blocked", "awaiting_input"), + root("between-steps", "idle", "quiescent"), + root("held", "idle", "paused"), + root("done", "idle", "quiescent"), + ], new Set(["done"]))).toEqual({ + active: ["drafting", "asking", "between-steps"], + paused: ["held"], + }); + expect(classifyRuns([])).toEqual({ active: [], paused: [] }); +}); + +test("run cards fold lifecycle events into ordered stages, Decisions waits, and paused or finished status", async () => { + let { foldLifecycle, runCards } = await import("./full"); + let cards = new Map(); + let event = (target: object, at: number) => + ({ + type: "workflow_lifecycle", + eventId: `e${at}`, + cursor: { epoch: "e", revision: at }, + runId: "run-1", + rootRunId: "run-1", + ownerSessionId: "session", + occurredAt: at * 1000, + observedAt: at * 1000, + delivery: "live", + target, + }) as never; + foldLifecycle( + cards, + event({ kind: "run", runId: "run-1", status: "running" }, 100), + "plan-review", + ); + foldLifecycle( + cards, + event( + { kind: "stage", runId: "run-1", stageId: "a", stageName: "draft-1", status: "running" }, + 101, + ), + ); + foldLifecycle( + cards, + event({ kind: "prompt", runId: "run-1", stageId: "a", promptId: "p1", status: "opened" }, 160), + ); + foldLifecycle( + cards, + event({ + kind: "stage", + runId: "run-1", + stageId: "a", + stageName: "draft-1", + status: "awaiting_input", + }, 160), + ); + let [waiting] = runCards(cards, { active: ["run-1"], paused: [] }); + expect(waiting).toMatchObject({ + id: "run-1", + name: "plan-review", + status: "waiting", + waiting: 1, + started: 100, + stages: [{ id: "run-1:a", name: "draft-1", status: "awaiting_input", started: 101 }], + }); + expect(runCards(cards, { active: [], paused: ["run-1"] })[0]!.status).toBe("paused"); + foldLifecycle( + cards, + event( + { kind: "prompt", runId: "run-1", stageId: "a", promptId: "p1", status: "answered" }, + 200, + ), + ); + foldLifecycle( + cards, + event({ + kind: "stage", + runId: "run-1", + stageId: "a", + stageName: "draft-1", + status: "completed", + }, 300), + ); + foldLifecycle( + cards, + event({ + kind: "stage", + runId: "run-1", + stageId: "b", + stageName: "reviewer-a-1", + status: "running", + }, 301), + ); + let [running] = runCards(cards, { active: ["run-1"], paused: [] }); + expect(running!.status).toBe("running"); + expect(running!.waiting).toBe(0); + expect(running!.stages.map(stage => [stage.name, stage.status, stage.ended])).toEqual([ + ["draft-1", "completed", 300], + ["reviewer-a-1", "running", undefined], + ]); + foldLifecycle(cards, event({ kind: "run", runId: "run-1", status: "completed" }, 900)); + expect(runCards(cards, { active: [], paused: [] })[0]).toMatchObject({ + status: "finished", + ended: 900, + }); +}); + +test("Stop and Resume use the session's run control, and a partial pause is reported", async () => { + let calls: unknown[] = []; + let outcome = { action: "pause", runId: "--all", status: "paused", message: "Paused 1 run(s)." }; + let workflows = { + pause: async (target: unknown) => { + calls.push(["pause", target]); + return outcome; + }, + listRuns: async (filter: unknown) => { + calls.push(["listRuns", filter]); + return [{ runId: "run-1" }, { runId: "run-2" }]; + }, + resume: async (runId: string) => { + calls.push(["resume", runId]); + return { action: "resume", runId, status: "running", message: "" }; + }, + } as unknown as SessionWorkflows; + await pauseOwnedRuns(workflows); + await resumeOwnedRuns(workflows); + expect(calls).toEqual([ + ["pause", { all: true }], + ["listRuns", { status: "paused" }], + ["resume", "run-1"], + ["resume", "run-2"], + ]); + outcome = { + ...outcome, + status: "partial", + failedRuns: [{ runId: "run-2", reason: "pause_failed", message: "busy" }], + } as typeof outcome; + await expect(pauseOwnedRuns(workflows)).rejects.toThrow("run-2: busy"); +}); + +test("an answer to a paused run's question is held until the run resumes, including child runs", async () => { + let planner: FullPlanner = { + cwd: "/", + humanInput: {} as HostInput, + runs: { active: [], paused: ["root"], cards: [] }, + rootOf: runId => runId === "child" ? "root" : runId, + }; + let released: string[] = []; + let controller = new AbortController(); + let held = untilUnpaused(planner, "child", controller.signal).then(() => released.push("child")); + await untilUnpaused(planner, "other", controller.signal).then(() => released.push("other")); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(released).toEqual(["other"]); + planner.runs = { active: ["root"], paused: [], cards: [] }; + for (let wake of planner.waiters ?? []) wake(); + await held; + expect(released).toEqual(["other", "child"]); + + planner.runs = { active: [], paused: ["root"], cards: [] }; + let aborted = new AbortController(); + let pending = untilUnpaused(planner, "root", aborted.signal); + aborted.abort(); + await pending; + expect(planner.waiters?.size ?? 0).toBe(0); +}); + +test("a run whose stage waits on a question counts as waiting even without prompt events", () => { + let cards = new Map([["root", { + id: "root", + name: "plan-review", + status: "running" as const, + started: 0, + updated: 0, + stages: [{ id: "root:s", name: "draft-1", status: "awaiting_input" as const }], + waiting: 0, + stageIndex: new Map(), + prompts: new Set(), + }]]); + expect(runCards(cards, { active: ["root"], paused: [] })[0]).toMatchObject({ + status: "waiting", + waiting: 1, + }); +}); + +test("a new run clears ended cards, keeps live ones, and lists only the twelve most recent stages", async () => { + let { foldLifecycle } = await import("./full"); + let cards = new Map(); + let event = (rootRunId: string, target: object, at: number) => + ({ + type: "workflow_lifecycle", + eventId: `${rootRunId}-${at}`, + cursor: { epoch: "e", revision: at }, + runId: rootRunId, + rootRunId, + ownerSessionId: "session", + occurredAt: at * 1000, + observedAt: at * 1000, + delivery: "live", + target, + }) as never; + foldLifecycle( + cards, + event("done", { kind: "run", runId: "done", status: "running" }, 1), + "first", + ); + foldLifecycle(cards, event("done", { kind: "run", runId: "done", status: "completed" }, 2)); + foldLifecycle( + cards, + event("live", { kind: "run", runId: "live", status: "running" }, 3), + "second", + ); + expect([...cards.keys()]).toEqual(["live"]); + foldLifecycle( + cards, + event("other", { kind: "run", runId: "other", status: "running" }, 4), + "third", + ); + expect([...cards.keys()]).toEqual(["live", "other"]); + + for (let index = 0; index < 15; index++) { + foldLifecycle( + cards, + event("live", { + kind: "stage", + runId: "live", + stageId: `s${index}`, + stageName: `stage-${index}`, + status: "completed", + }, 10 + index), + ); + } + let [card] = runCards(cards, { active: ["live", "other"], paused: [] }); + expect(card?.stages).toHaveLength(12); + expect(card?.stages[0]?.name).toBe("stage-3"); + expect(card?.earlierStages).toBe(3); +}); + +test("ctx.tool steps appear among the stages, marked as tool steps", async () => { + let { foldLifecycle } = await import("./full"); + let cards = new Map(); + let event = (target: object, at: number) => + ({ + type: "workflow_lifecycle", + eventId: `t${at}`, + cursor: { epoch: "e", revision: at }, + runId: "run-1", + rootRunId: "run-1", + ownerSessionId: "session", + occurredAt: at * 1000, + observedAt: at * 1000, + delivery: "live", + target, + }) as never; + foldLifecycle(cards, event({ kind: "run", runId: "run-1", status: "running" }, 1), "demo"); + let tool = (toolNodeId: string, toolName: string, status: string, at: number) => + foldLifecycle(cards, event({ kind: "tool", runId: "run-1", toolNodeId, toolName, status }, at)); + tool("t1", "prepare", "running", 2); + tool("t1", "prepare", "cached", 5); + tool("t2", "work", "running", 6); + tool("t3", "finish", "cancelled", 7); + let [card] = runCards(cards, { active: ["run-1"], paused: [] }); + expect(card?.stages).toEqual([ + { id: "run-1:t1", name: "prepare", kind: "tool", status: "completed", started: 2, ended: 5 }, + { id: "run-1:t2", name: "work", kind: "tool", status: "running", started: 6 }, + { id: "run-1:t3", name: "finish", kind: "tool", status: "skipped", ended: 7 }, + ]); +}); diff --git a/apps/server/src/harness/atomic/full.ts b/apps/server/src/harness/atomic/full.ts index fc79c5a5..241ea259 100644 --- a/apps/server/src/harness/atomic/full.ts +++ b/apps/server/src/harness/atomic/full.ts @@ -1,8 +1,76 @@ -import type { HostInput } from "@bastani/atomic"; +import type { + ExtensionAPI, + ExtensionFactory, + HostInput, + SessionWorkflows, + WorkflowActivitySubscription, + WorkflowLifecycleEvent, + WorkflowRootActivity, + WorkflowToolNodeStatus, +} from "@bastani/atomic"; +import type { Chat as Wire } from "@chopin/protocol"; + +/** + * Root workflow runs a Planner session owns: still live, or paused and + * resumable, plus what each run's card shows. + */ +export type PlannerRuns = { active: string[]; paused: string[]; cards: Wire.Run[] }; + +export type FullPlanner = { + cwd: string; + humanInput: HostInput; + /** Latest known runs; absent until Atomic reports its workflow activity as ready. */ + runs?: PlannerRuns; + /** Called with every change to `runs`. */ + onRuns?: (runs: PlannerRuns) => void; + /** The live session's run control; set once the session exists. */ + workflows?: SessionWorkflows; + /** The root run that owns a run, as learned from lifecycle events. */ + rootOf?: (runId: string) => string; + /** Woken on every change to `runs`. */ + waiters?: Set<() => void>; +}; -export type FullPlanner = { cwd: string; humanInput: HostInput }; let planners = new Map(); +/** Pauses every run the session owns. A run whose stage waits on a question counts as paused. */ +export async function pauseOwnedRuns(workflows: SessionWorkflows): Promise { + let outcome = await workflows.pause({ all: true }); + if (outcome.status === "partial") { + let still = (outcome.failedRuns ?? []).map(run => `${run.runId}: ${run.message}`).join("; "); + throw new Error(`Some workflow runs are still active: ${still}`); + } +} + +/** Resumes every paused run the session owns. */ +export async function resumeOwnedRuns(workflows: SessionWorkflows): Promise { + for (let run of await workflows.listRuns({ status: "paused" })) await workflows.resume(run.runId); +} + +/** + * Resolves once the root run owning `runId` is not paused. Atomic leaves a + * question open when its run pauses, so an answer given meanwhile is held here + * until the run resumes, keeping the paused run fully stopped. + */ +export async function untilUnpaused( + planner: FullPlanner, + runId: string, + signal: AbortSignal, +): Promise { + let paused = () => planner.runs?.paused.includes(planner.rootOf?.(runId) ?? runId) ?? false; + while (paused() && !signal.aborted) { + await new Promise(resolve => { + let wake = () => { + planner.waiters?.delete(wake); + signal.removeEventListener("abort", wake); + resolve(); + }; + (planner.waiters ??= new Set()).add(wake); + signal.addEventListener("abort", wake, { once: true }); + }); + } +} + /** Every Atomic Planner session is registered; background workers never are. */ export function registerFullPlanner(sessionId: string, planner: FullPlanner): () => void { if (planners.has(sessionId)) throw new Error("Atomic Planner session is already registered"); @@ -15,3 +83,215 @@ export function registerFullPlanner(sessionId: string, planner: FullPlanner): () export function fullPlanner(sessionId: string): FullPlanner | undefined { return planners.get(sessionId); } + +/** + * A root Atomic reports as paused is paused; one whose run has ended, as told + * by its lifecycle events, is neither; any other root is live. Activity alone + * cannot say a run ended: a root is idle and quiescent whenever nothing is + * counted as executing, which also happens between steps of a live run. + */ +export function classifyRuns( + roots: Iterable, + ended: ReadonlySet = new Set(), +): Omit { + let runs = { active: [] as string[], paused: [] as string[] }; + for (let root of roots) { + if (root.reason === "paused") runs.paused.push(root.rootRunId); + else if (!ended.has(root.rootRunId)) runs.active.push(root.rootRunId); + } + return runs; +} + +/** Stages kept per run; earlier ones are counted, not listed. */ +const STORED_STAGES = 12; + +type Card = Wire.Run & { stageIndex: Map; prompts: Set }; + +const ENDED: Record = { + completed: "finished", + skipped: "finished", + failed: "failed", + blocked: "blocked", + killed: "stopped", + cancelled: "stopped", +}; + +/** A `ctx.tool` step shown with the stage statuses it corresponds to. */ +const TOOL_STEP: Record = { + pending: "pending", + running: "running", + completed: "completed", + cached: "completed", + failed: "failed", + cancelled: "skipped", +}; + +/** Records one stage or tool step, in the order the run first reached it. */ +function foldStep( + card: Card, + key: string, + name: string, + status: Wire.RunStage["status"], + seconds: number, + kind?: "tool", +): void { + let index = card.stageIndex.get(key); + if (index === undefined) { + index = card.stages.length; + card.stageIndex.set(key, index); + card.stages.push({ id: key, name, status, ...(kind ? { kind } : {}) }); + } + let stage = card.stages[index]!; + stage.status = status; + if (status === "running" && stage.started === undefined) stage.started = seconds; + if (status === "completed" || status === "failed" || status === "skipped") stage.ended = seconds; +} + +/** Folds one lifecycle event into its root run's card. */ +export function foldLifecycle( + cards: Map, + event: WorkflowLifecycleEvent, + name?: string, +): void { + let seconds = Math.floor(event.occurredAt / 1000); + let card = cards.get(event.rootRunId); + if (!card) { + for (let [id, other] of cards) if (other.ended !== undefined) cards.delete(id); + card = { + id: event.rootRunId, + name: name ?? "workflow", + status: "running", + started: seconds, + updated: seconds, + stages: [], + waiting: 0, + stageIndex: new Map(), + prompts: new Set(), + }; + cards.set(event.rootRunId, card); + } + if (name) card.name = name; + card.updated = Math.max(card.updated, seconds); + let target = event.target; + if (target.kind === "run" && target.runId === event.rootRunId) { + card.status = ENDED[target.status] ?? (target.status === "paused" ? "paused" : "running"); + if ( + card.status === "finished" || card.status === "blocked" || card.status === "failed" + || card.status === "stopped" + ) { + card.ended = seconds; + } else delete card.ended; + } else if (target.kind === "stage") { + foldStep(card, `${target.runId}:${target.stageId}`, target.stageName, target.status, seconds); + } else if (target.kind === "tool") { + foldStep( + card, + `${target.runId}:${target.toolNodeId}`, + target.toolName, + TOOL_STEP[target.status], + seconds, + "tool", + ); + } else if (target.kind === "prompt") { + if (target.status === "opened") card.prompts.add(target.promptId); + else card.prompts.delete(target.promptId); + card.waiting = card.prompts.size; + } +} + +/** The run cards to show, with the live/paused split from activity taking precedence over lifecycle. */ +export function runCards(cards: Map, runs: Omit): Wire.Run[] { + return [...cards.values()].map(({ stageIndex: _index, prompts: _prompts, ...card }) => { + let waiting = Math.max( + card.waiting, + card.stages.filter(stage => stage.status === "awaiting_input").length, + ); + let status: Wire.Run["status"] = runs.paused.includes(card.id) + ? "paused" + : runs.active.includes(card.id) + ? waiting > 0 ? "waiting" : "running" + : card.status === "running" || card.status === "waiting" + ? "finished" + : card.status; + let earlierStages = Math.max(0, card.stages.length - STORED_STAGES); + return { + ...card, + status, + waiting, + stages: card.stages.slice(earlierStages).map(stage => ({ ...stage })), + ...(earlierStages ? { earlierStages } : {}), + }; + }); +} + +/** The workflow run id reported in a `workflow` tool result, when the call started one. */ +function launchedRunId(content: readonly { type: string; text?: string }[]): string | undefined { + let text = content.map(part => part.type === "text" ? part.text ?? "" : "").join("\n"); + return /"runId"\s*:\s*"([0-9a-f-]{36})"/.exec(text)?.[1]; +} + +/** + * Keeps `planner.runs` in step with the workflows this session owns, through + * Atomic's public activity stream and lifecycle hooks. Recovering or + * unavailable snapshots are unknown rather than empty, so they leave the last + * known runs in place. + */ +export function workflowRuns(planner: FullPlanner): ExtensionFactory { + return (atomic: ExtensionAPI): void => { + let lease: WorkflowActivitySubscription | undefined; + let roots = new Map(); + let cards = new Map(); + let names = new Map(); + let parents = new Map(); + let seen = new Set(); + planner.rootOf = runId => parents.get(runId) ?? runId; + let ready = false; + let publish = () => { + if (!ready) return; + for (let [id, card] of cards) { + if (seen.has(id) && !roots.has(id) && card.ended === undefined) { + card.status = "finished"; + card.ended = card.updated; + } + } + let ended = new Set( + [...cards.values()].filter(card => card.ended !== undefined).map(card => card.id), + ); + let split = classifyRuns(roots.values(), ended); + planner.runs = { ...split, cards: runCards(cards, split) }; + planner.onRuns?.(planner.runs); + for (let wake of planner.waiters ?? []) wake(); + }; + atomic.on("session_start", (_event, ctx) => { + lease?.dispose(); + lease = ctx.observeWorkflowActivity(frame => { + if (frame.kind === "snapshot") { + if (frame.availability !== "ready") return; + ready = true; + roots = new Map(frame.roots.map(root => [root.rootRunId, root])); + } else if (frame.kind === "changed") roots.set(frame.root.rootRunId, frame.root); + else roots.delete(frame.rootRunId); + for (let id of roots.keys()) seen.add(id); + publish(); + }); + }); + atomic.on("tool_result", event => { + if (event.toolName !== "workflow" || event.input.action !== "run") return; + let runId = launchedRunId(event.content); + let name = typeof event.input.workflow === "string" ? event.input.workflow : undefined; + if (!runId || !name) return; + names.set(runId, name); + let card = cards.get(runId); + if (card) { + card.name = name; + publish(); + } + }); + atomic.on("workflow_lifecycle", event => { + if (event.target.kind !== "prompt") parents.set(event.target.runId, event.rootRunId); + foldLifecycle(cards, event, names.get(event.rootRunId)); + publish(); + }); + atomic.on("session_shutdown", () => lease?.dispose()); + }; +} diff --git a/apps/server/src/harness/atomic/human-input.test.ts b/apps/server/src/harness/atomic/human-input.test.ts index ddd0f666..7dfb32c0 100644 --- a/apps/server/src/harness/atomic/human-input.test.ts +++ b/apps/server/src/harness/atomic/human-input.test.ts @@ -554,3 +554,31 @@ test("an answer that would make the document too large leaves it unchanged and t f.controller.abort(); expect((await response).cancelled).toBe(true); }); + +test("a workflow question's answer is returned only after its paused run is released", async () => { + let release!: () => void; + let gate = new Promise(resolve => { + release = resolve; + }); + let held: string[] = []; + let f = await hostInputRoom(undefined, async runId => { + held.push(runId); + await gate; + }); + cleanups.push(f.close); + let settled = false; + let response = f.input.questionnaire( + { questions: [params.questions[0]!] }, + { ...f.options, workflowRunId: "run-1", workflowStageId: "stage-1" }, + ).then(result => { + settled = true; + return result; + }); + let [card] = await f.cards(1); + await f.answer(card!.id, [0]); + await Bun.sleep(20); + expect(held).toEqual(["run-1"]); + expect(settled).toBe(false); + release(); + expect((await response).answers).toMatchObject([{ kind: "option", answer: " A " }]); +}); diff --git a/apps/server/src/harness/atomic/human-input.ts b/apps/server/src/harness/atomic/human-input.ts index f64ed743..4cb395e8 100644 --- a/apps/server/src/harness/atomic/human-input.ts +++ b/apps/server/src/harness/atomic/human-input.ts @@ -15,11 +15,13 @@ import type { DocumentRoom } from "../../agent/tools"; * Atomic owns input requests; Chopin owns their shared decision records. * * A request nobody answers within `expiresInMs` expires: its cards stay in - * Decisions marked expired, and Atomic gets no answer. + * Decisions marked expired, and Atomic gets no answer. An answer to a workflow + * question is returned only once `hold` allows it, so a paused run stays stopped. */ export function createHumanInput( room: DocumentRoom, expiresInMs = limits.INPUT_EXPIRY_MS, + hold?: (workflowRunId: string, signal: AbortSignal) => Promise, ): HostInput { async function questionnaire( params: QuestionParams, @@ -58,6 +60,10 @@ export function createHumanInput( if (options.signal.aborted || ended.some(outcome => outcome.status === "expired")) { return { answers: [], cancelled: true }; } + if (hold && options.workflowRunId !== undefined) { + await hold(options.workflowRunId, options.signal); + if (options.signal.aborted) return { answers: [], cancelled: true }; + } let answers: QuestionAnswer[] = []; ended.forEach((outcome, questionIndex) => { if (outcome.status !== "answered") return; diff --git a/apps/server/src/harness/atomic/workspace.ts b/apps/server/src/harness/atomic/workspace.ts index a7e65772..b8053a18 100644 --- a/apps/server/src/harness/atomic/workspace.ts +++ b/apps/server/src/harness/atomic/workspace.ts @@ -4,24 +4,24 @@ * A checkout is only ever supplied through `invoke_planner`, verified against * the channel's repository, and remembered for the channel until the process * exits. Every later Planner session for that channel re-verifies it before use. - * Without one, the session gets a directory Chopin created empty for that - * channel alone; it keeps whatever the Planner writes there until shutdown. + * Without one, the session works in a directory Chopin keeps for that channel + * alone under its per-user state directory, so files a paused or resumed + * workflow wrote survive server restarts instead of filling the shared tmp. */ -import { mkdtemp, rm, stat } from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { chmod, lstat, mkdir } from "node:fs/promises"; +import { homedir } from "node:os"; import { join } from "node:path"; import { verifiedCheckout } from "./checkout"; export type PlannerWorkspace = { cwd: string; - /** True when `cwd` is a verified checkout rather than the empty channel directory. */ + /** True when `cwd` is a verified checkout rather than the channel's own directory. */ checkout: boolean; }; let checkouts = new Map(); -let directories = new Map>(); export function rememberCheckout(channelId: string, checkout: string): void { checkouts.set(channelId, checkout); @@ -39,29 +39,35 @@ export async function plannerWorkspace( return { cwd: await directory(channelId), checkout: false }; } +/** The platform's per-user state directory for Chopin. */ +export function stateDirectory(): string { + if (process.platform === "win32") { + return join(process.env.LOCALAPPDATA || join(homedir(), "AppData", "Local"), "Chopin"); + } + if (process.platform === "darwin") { + return join(homedir(), "Library", "Application Support", "Chopin"); + } + return join(process.env.XDG_STATE_HOME || join(homedir(), ".local", "state"), "chopin"); +} + /** - * Chained per channel so concurrent sessions share one directory. `mkdtemp` - * creates it with mode 0700; one removed from under a running server is - * replaced rather than recreated at a name somebody else could have taken. + * A stable directory per channel, private to the server's user. It lives in a + * per-user directory rather than the shared tmp, so a predictable name is safe; + * a symlink or file planted at that name is refused rather than followed. */ -function directory(channelId: string): Promise { - let previous = directories.get(channelId); - let next = (async () => { - let path = await previous?.catch(() => undefined); - if (path && await stat(path).then(found => found.isDirectory(), () => false)) return path; - return mkdtemp(join(tmpdir(), "chopin-planner-")); - })(); - directories.set(channelId, next); - return next; +async function directory(channelId: string): Promise { + if (!/^[\w-]+$/.test(channelId)) throw new Error("Planner workspace needs a plain channel id"); + let path = join(stateDirectory(), "planner", channelId); + await mkdir(path, { recursive: true, mode: 0o700 }); + let found = await lstat(path); + if (!found.isDirectory() || found.isSymbolicLink()) { + throw new Error(`Planner workspace ${path} is not a directory`); + } + await chmod(path, 0o700); + return path; } -export async function removeWorkspaces(): Promise { - let paths = await Promise.all( - [...directories.values()].map(path => path.catch(() => undefined)), - ); - directories.clear(); +/** Forget remembered checkouts. Channel directories persist across restarts. */ +export function forgetWorkspaces(): void { checkouts.clear(); - await Promise.all( - paths.map(path => path && rm(path, { recursive: true, force: true }).catch(() => {})), - ); } diff --git a/apps/server/src/harness/harnesses.ts b/apps/server/src/harness/harnesses.ts index 44228f29..d8928c26 100644 --- a/apps/server/src/harness/harnesses.ts +++ b/apps/server/src/harness/harnesses.ts @@ -1,5 +1,5 @@ import { ATOMIC_AUTH_MODES, createAtomicAdapter } from "./atomic/adapter"; -import { removeWorkspaces } from "./atomic/workspace"; +import { forgetWorkspaces } from "./atomic/workspace"; import { createCopilotSdk } from "./copilot-sdk/adapter"; import { createPiAdapter } from "./pi/adapter"; @@ -113,5 +113,5 @@ export async function shutdownHarnesses(): Promise { credentials.clear(); await selected?.shutdown(); selected = undefined; - await removeWorkspaces(); + forgetWorkspaces(); } diff --git a/apps/server/src/harness/session.test.ts b/apps/server/src/harness/session.test.ts index f9573744..1575f27e 100644 --- a/apps/server/src/harness/session.test.ts +++ b/apps/server/src/harness/session.test.ts @@ -1,10 +1,10 @@ -import { afterEach, describe, expect, it } from "bun:test"; -import { mkdtemp, readdir, realpath, rm, stat } from "node:fs/promises"; +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import { mkdir, mkdtemp, readdir, realpath, rm, stat, symlink } from "node:fs/promises"; import { join } from "node:path"; import { tmpdir } from "node:os"; import { openPlannerSession } from "./session"; import { fullPlanner } from "./atomic/full"; -import { rememberCheckout, removeWorkspaces } from "./atomic/workspace"; +import { forgetWorkspaces, rememberCheckout, stateDirectory } from "./atomic/workspace"; import { plannerInstructions } from "../agent/planner"; import type { ActiveOwnerBinding } from "../agent/active-owner"; @@ -120,8 +120,16 @@ describe("openPlannerSession", () => { describe("Planner workspaces", () => { let roots: string[] = []; + let previousState = process.env.XDG_STATE_HOME; + beforeEach(async () => { + let state = await mkdtemp(join(tmpdir(), "chopin-planner-state-")); + roots.push(state); + process.env.XDG_STATE_HOME = state; + }); afterEach(async () => { - await removeWorkspaces(); + forgetWorkspaces(); + if (previousState === undefined) delete process.env.XDG_STATE_HOME; + else process.env.XDG_STATE_HOME = previousState; for (let root of roots.splice(0)) await rm(root, { recursive: true, force: true }); }); @@ -175,6 +183,8 @@ describe("Planner workspaces", () => { expect(registered?.humanInput.questionnaire).toBeFunction(); expect(instructions).toContain(`${path}, is a local checkout of owner/repo`); expect(instructions).toContain("proceed on your best judgement"); + expect(instructions).toContain("never tell them to use `/workflow connect`"); + expect(instructions).toContain("Chat shows a card for each run"); expect(instructions).not.toContain("You have no shell"); } let git = Bun.spawn(["git", "-C", path, "remote", "set-url", "origin", "workgit:other/repo"], { @@ -185,27 +195,42 @@ describe("Planner workspaces", () => { let fallback = await open("channel", "atomic"); expect(fallback.registered?.cwd).not.toBe(path); expect(await readdir(fallback.registered!.cwd)).toEqual([]); - expect(fallback.instructions).toContain("scratch directory Chopin created"); + expect(fallback.instructions).toContain("scratch directory Chopin keeps"); }); - it("gives each atomic channel without a checkout its own empty, private directory", async () => { + it("gives each atomic channel without a checkout its own private directory that outlives shutdown", async () => { let first = await open("first", "atomic"); let again = await open("first", "atomic"); let second = await open("second", "atomic"); expect(again.registered?.cwd).toBe(first.registered!.cwd); expect(second.registered?.cwd).not.toBe(first.registered!.cwd); - for (let { registered, instructions } of [first, second]) { + for ( + let [name, { registered, instructions }] of [["first", first], ["second", second]] as const + ) { let cwd = registered!.cwd; + if (process.platform === "linux") expect(cwd).toBe(join(stateDirectory(), "planner", name)); expect((await stat(cwd)).mode & 0o777).toBe(0o700); expect(await readdir(cwd)).toEqual([]); expect(registered?.humanInput.questionnaire).toBeFunction(); - expect(instructions).toContain(`${cwd}, is a scratch directory Chopin created`); + expect(instructions).toContain(`${cwd}, is a scratch directory Chopin keeps`); expect(instructions).toContain("holds no repository files"); expect(instructions).toContain("`read_repository_file`"); expect(instructions).toContain("proceed on your best judgement"); + expect(instructions).toContain("never tell them to use `/workflow connect`"); + expect(instructions).toContain("Chat shows a card for each run"); } - await removeWorkspaces(); - expect(await stat(first.registered!.cwd).catch(() => undefined)).toBeUndefined(); + forgetWorkspaces(); + expect((await stat(first.registered!.cwd)).isDirectory()).toBe(true); + expect((await open("first", "atomic")).registered?.cwd).toBe(first.registered!.cwd); + }); + + it("refuses a symlink planted where a channel's directory belongs", async () => { + if (process.platform !== "linux") return; + let target = await mkdtemp(join(tmpdir(), "chopin-planner-target-")); + roots.push(target); + await mkdir(join(stateDirectory(), "planner"), { recursive: true }); + await symlink(target, join(stateDirectory(), "planner", "planted")); + await expect(open("planted", "atomic")).rejects.toThrow(); }); it("keeps copilot-sdk and pi Planner sessions isolated, whatever the channel remembers", async () => { @@ -215,6 +240,7 @@ describe("Planner workspaces", () => { expect(registered).toBeUndefined(); expect(instructions).toContain("You have no shell"); expect(instructions).not.toContain("proceed on your best judgement"); + expect(instructions).not.toContain("/workflow connect"); } }); }); diff --git a/apps/server/src/harness/session.ts b/apps/server/src/harness/session.ts index d02751a9..5804d871 100644 --- a/apps/server/src/harness/session.ts +++ b/apps/server/src/harness/session.ts @@ -2,7 +2,14 @@ import { createJustBashNetworkSandboxSession } from "@ai-sdk/sandbox-just-bash"; import { plannerAgent } from "./agents"; import { githubTools, type GitHubToolsError, type Result } from "./github-tools"; import { registerCredential } from "./harnesses"; -import { registerFullPlanner } from "./atomic/full"; +import { + type FullPlanner, + pauseOwnedRuns, + type PlannerRuns, + registerFullPlanner, + resumeOwnedRuns, + untilUnpaused, +} from "./atomic/full"; import { createHumanInput } from "./atomic/human-input"; import { type PlannerWorkspace, plannerWorkspace } from "./atomic/workspace"; @@ -30,9 +37,26 @@ export type PlannerChannel = { harness?: string; }; +/** Atomic answers an action it could not carry out with `noop` or `cancelled`; report it rather than claim success. */ +function applied(outcome: { status: string; message: string }): void { + if (outcome.status === "noop" || outcome.status === "cancelled") throw new Error(outcome.message); +} + export type PlannerSession = { stream: (prompt: string, abortSignal: AbortSignal) => ReturnType; destroy: () => Promise; + /** Workflow runs this session owns; only atomic Planner sessions report them. */ + runs?: () => PlannerRuns | undefined; + /** Subscribes to run changes; returns the unsubscribe function. */ + watchRuns?: (listener: (runs: PlannerRuns) => void) => () => void; + /** Pauses every workflow run this session owns, resumably. */ + pauseRuns?: () => Promise; + /** Resumes the runs this session paused. */ + resumeRuns?: () => Promise; + /** Pauses one run this session owns. */ + pauseRun?: (runId: string) => Promise; + /** Resumes one paused run this session owns. */ + resumeRun?: (runId: string) => Promise; }; export type PlannerSessionDependencies = { @@ -82,11 +106,21 @@ export async function openPlannerSession( let workspace = channel.harness === "atomic" ? await plannerWorkspace(channel.room.id, channel.repository) : undefined; + let planner: FullPlanner | undefined; + let listeners = new Set<(runs: PlannerRuns) => void>(); if (workspace) { - unregisterFull = registerFullPlanner(sessionId, { + planner = { cwd: workspace.cwd, - humanInput: createHumanInput(channel.room), - }); + humanInput: createHumanInput( + channel.room, + undefined, + (runId, signal) => untilUnpaused(planner!, runId, signal), + ), + onRuns: runs => { + for (let listener of listeners) listener(runs); + }, + }; + unregisterFull = registerFullPlanner(sessionId, planner); } let instructions = typeof channel.instructions === "function" ? channel.instructions(workspace) @@ -109,9 +143,29 @@ export async function openPlannerSession( let active = session; let release = unregister; let stopped: Promise | undefined; + let workflows = () => { + if (!planner?.workflows) { + throw new Error("The Planner has no live Atomic session to control."); + } + return planner.workflows; + }; + let runControl = planner + ? { + runs: () => planner.runs, + watchRuns: (listener: (runs: PlannerRuns) => void) => { + listeners.add(listener); + return () => listeners.delete(listener); + }, + pauseRuns: () => pauseOwnedRuns(workflows()), + resumeRuns: () => resumeOwnedRuns(workflows()), + pauseRun: async (runId: string) => applied(await workflows().pause(runId)), + resumeRun: async (runId: string) => applied(await workflows().resume(runId)), + } + : {}; return { ok: true, value: { + ...runControl, stream: (prompt, abortSignal) => agent.stream({ session: active, diff --git a/apps/server/src/main.ts b/apps/server/src/main.ts index bf4b4c77..04344c49 100644 --- a/apps/server/src/main.ts +++ b/apps/server/src/main.ts @@ -190,6 +190,13 @@ function conversation( ownerAvailable: () => jobRunner?.ownerAvailable(room.id) ?? Promise.resolve(), jobs: config.backgroundJobs ? jobService : undefined, references: referenceService, + hold: () => { + let held = Rooms.hold(room.id); + return () => { + held.release(); + evict(held.room); + }; + }, createResearch: config.agent ? async request => { let service = researchService; @@ -306,6 +313,15 @@ async function receive(ws: Socket, raw: string): Promise { if (room.plan) await Chat.abort(chat(room, ws), ws); return; + case "chat:resume": + if (room.plan) await Chat.resume(chat(room, ws), ws); + return; + + case "chat:pause-run": + case "chat:resume-run": + if (room.plan) await Chat.controlRun(chat(room, ws), ws, frame); + return; + case "chat:unqueue": if (room.plan) Chat.unqueue(chat(room, ws), ws, frame); return; diff --git a/apps/server/src/mcp.ts b/apps/server/src/mcp.ts index d0ec91a9..0a5faf0c 100644 --- a/apps/server/src/mcp.ts +++ b/apps/server/src/mcp.ts @@ -104,6 +104,7 @@ export type CreateDocument = { | { kind: "created"; document: CreatedDocument } | { kind: "replayed"; document: CreatedDocument } | { kind: "conflict" } + | { kind: "title-taken" } | { kind: "forbidden" } | { kind: "unavailable" } >; @@ -213,6 +214,28 @@ const ARCHIVED_DOCUMENT = { required: ["id", "title", "archivedAt"], }; +/** Validation failures of a submitted plan, returned instead of a document. */ +const ISSUES = { + type: "object", + properties: { + issues: { + type: "array", + items: { + type: "object", + properties: { + code: { type: "string" }, + message: { type: "string" }, + path: { type: "string" }, + }, + required: ["code", "message", "path"], + additionalProperties: true, + }, + }, + }, + required: ["issues"], + additionalProperties: false, +}; + function outcome(codes: string[]) { return { type: "object", @@ -321,15 +344,22 @@ export const TOOLS: Tool[] = [ }, outputSchema: { type: "object", - properties: { - ...DOCUMENT.properties, - brief: BRIEF, - source: { type: "string" }, - revision: { type: "integer", minimum: 0 }, - url: { type: "string" }, - }, - required: ["id", "title", "brief", "source", "revision", "url"], - additionalProperties: false, + oneOf: [ + { + type: "object", + properties: { + ...DOCUMENT.properties, + brief: BRIEF, + source: { type: "string" }, + revision: { type: "integer", minimum: 0 }, + url: { type: "string" }, + }, + required: ["id", "title", "brief", "source", "revision", "url"], + additionalProperties: false, + }, + outcome(["idempotency-conflict", "title-taken", "document-unavailable"]), + ISSUES, + ], }, }, { @@ -383,26 +413,7 @@ export const TOOLS: Tool[] = [ "repository-forbidden", "document-unavailable", ]), - { - type: "object", - properties: { - issues: { - type: "array", - items: { - type: "object", - properties: { - code: { type: "string" }, - message: { type: "string" }, - path: { type: "string" }, - }, - required: ["code", "message", "path"], - additionalProperties: true, - }, - }, - }, - required: ["issues"], - additionalProperties: false, - }, + ISSUES, ], }, }, @@ -899,6 +910,9 @@ export function handler( if (outcome.kind === "conflict") { return respond(text({ code: "idempotency-conflict" }, true)); } + if (outcome.kind === "title-taken") { + return respond(text({ code: "title-taken" }, true)); + } if (outcome.kind === "unavailable") { return respond(text({ code: "document-unavailable" }, true)); } diff --git a/apps/server/src/mcp/hosted.test.ts b/apps/server/src/mcp/hosted.test.ts index 236e9df0..5138db72 100644 --- a/apps/server/src/mcp/hosted.test.ts +++ b/apps/server/src/mcp/hosted.test.ts @@ -716,6 +716,28 @@ describe("the hosted MCP adapter", () => { ).toEqual({ kind: "conflict" }); }); + it("reports a title the repository already uses as title-taken, not an idempotency conflict", async () => { + let context = setup(); + context.github.repositoryValue = { + ...context.github.repositoryValue, + permissions: { pull: true, push: true, admin: false }, + }; + let adapter = hosted(context.auth); + let caller = await adapter.caller(request("Bearer allowed")); + if (!caller) throw new Error("test caller was not authenticated"); + if (!adapter.create) throw new Error("hosted creation adapter is unavailable"); + await adapter.create.create(caller, creation); + + expect( + await adapter.create.create(caller, { + ...creation, + idempotencyKey: "another-request", + fingerprint: "another-request", + title: creation.title.toUpperCase(), + }), + ).toEqual({ kind: "title-taken" }); + }); + it("reconstructs a stored checkpoint and journal with the validated plan revision", async () => { let context = setup(); let opened = await plan(context); diff --git a/apps/server/src/mcp/hosted.ts b/apps/server/src/mcp/hosted.ts index 29b97940..3df9adcc 100644 --- a/apps/server/src/mcp/hosted.ts +++ b/apps/server/src/mcp/hosted.ts @@ -508,9 +508,9 @@ export function hosted( if (!(err instanceof StorageError) || err.failure !== "conflict") throw err; if (callbacks.isChannelDeleting?.(id)) return { kind: "unavailable" }; let stored = await auth.storage.collaboration.load(id, auth.clock()); - if (!stored || stored.channel.repositoryId !== repository.id) { - return { kind: "conflict" }; - } + // Nothing stored under this key's id: the repository already has a document with this title. + if (!stored) return { kind: "title-taken" }; + if (stored.channel.repositoryId !== repository.id) return { kind: "conflict" }; let restored = await Plan.readStored(stored); if ( !restored.creation diff --git a/apps/server/src/plan/persistence.test.ts b/apps/server/src/plan/persistence.test.ts index 87c15c1c..43fdf069 100644 --- a/apps/server/src/plan/persistence.test.ts +++ b/apps/server/src/plan/persistence.test.ts @@ -168,6 +168,52 @@ describe("hosted plan persistence", () => { await Service.close(restored); }); + it("restores workflow run cards, with runs still live then stopped", async () => { + let context = await hosted(); + let plan = await Service.open(context.channel.id, context.backend, context.server); + let stage = { id: "run-1:draft", name: "draft-1", status: "running" as const, started: 100 }; + plan.chat.runs = [ + { + id: "run-1", + name: "plan-review", + status: "running", + started: 100, + updated: 160, + stages: [stage], + waiting: 1, + }, + { + id: "run-0", + name: "docs-pass", + status: "finished", + started: 10, + updated: 90, + ended: 90, + stages: [{ ...stage, id: "run-0:draft", status: "completed", ended: 90 }], + earlierStages: 3, + waiting: 0, + }, + ]; + await Service.persist(plan); + await Service.close(plan); + + let restored = await Service.open(context.channel.id, context.backend, context.server); + expect(restored.chat.runs?.map(run => [run.id, run.status, run.waiting])).toEqual([ + ["run-1", "stopped", 0], + ["run-0", "finished", 0], + ]); + expect(restored.chat.runs?.[0]?.ended).toBeNumber(); + expect(restored.chat.runs?.[1]).toMatchObject({ ended: 90, earlierStages: 3 }); + await Service.close(restored); + + let empty = await Service.open(context.channel.id, context.backend, context.server); + empty.chat.runs = undefined; + await Service.persist(empty); + await Service.close(empty); + expect((await Service.open(context.channel.id, context.backend, context.server)).chat.runs) + .toBeUndefined(); + }); + it("restores typed chat references with their observed target state", async () => { let context = await hosted(); let plan = await Service.open(context.channel.id, context.backend, context.server); diff --git a/apps/server/src/plan/service.ts b/apps/server/src/plan/service.ts index 78588607..c724d12c 100644 --- a/apps/server/src/plan/service.ts +++ b/apps/server/src/plan/service.ts @@ -207,6 +207,7 @@ type Sidecar = { openQuestions: Questions.StoredOpen[]; threads: Comments.Record[]; transcript: Chat.Chat["entries"]; + workflowRuns?: NonNullable; }; function state(plan: Plan): Sidecar { @@ -225,6 +226,7 @@ function state(plan: Plan): Sidecar { openQuestions: Questions.dump(plan.questions), threads: [...plan.threads.values()], transcript: plan.chat.entries, + ...(plan.chat.runs?.length ? { workflowRuns: plan.chat.runs } : {}), }; } @@ -385,6 +387,7 @@ function restoredState( if (Object.hasOwn(item, "execution")) expected.push("execution"); if (Object.hasOwn(item, "lifecycle")) expected.push("lifecycle"); if (Object.hasOwn(item, "mcpUpdates")) expected.push("mcpUpdates"); + if (Object.hasOwn(item, "workflowRuns")) expected.push("workflowRuns"); expected.sort(); if ( keys.length !== expected.length @@ -472,6 +475,7 @@ function restoredState( } } let mcpUpdates = restoreMcpUpdates(item.mcpUpdates); + let workflowRuns = restoreWorkflowRuns(item.workflowRuns); return { version: 1, revision: item.revision, @@ -485,9 +489,63 @@ function restoredState( openQuestions: openQuestions as unknown as Questions.StoredOpen[], threads: threads as never[], transcript: transcript as unknown as Chat.Chat["entries"], + ...(workflowRuns.length > 0 ? { workflowRuns } : {}), }; } +const RUN_STATUSES = new Set([ + "running", + "waiting", + "paused", + "finished", + "blocked", + "failed", + "stopped", +]); +const RUN_STAGE_STATUSES = new Set([ + "pending", + "running", + "awaiting_input", + "paused", + "blocked", + "completed", + "failed", + "skipped", +]); + +function seconds(value: JsonValue | undefined, optional = false): boolean { + if (value === undefined) return optional; + return typeof value === "number" && Number.isSafeInteger(value) && value >= 0; +} + +/** Workflow run cards Chat showed, as `Chat.Run` values. */ +function restoreWorkflowRuns(value: JsonValue | undefined): NonNullable { + if (value === undefined) return []; + if (!Array.isArray(value)) throw new Error("hosted channel has invalid workflow runs"); + let runs = objects(value, "workflow run"); + for (let run of runs) { + let stages = run.stages; + if ( + typeof run.name !== "string" + || !RUN_STATUSES.has(run.status as string) + || !seconds(run.started) + || !seconds(run.updated) + || !seconds(run.ended, true) + || !seconds(run.earlierStages, true) + || !seconds(run.waiting) + || !Array.isArray(stages) + || stages.some(stage => + !stage || typeof stage !== "object" || Array.isArray(stage) + || typeof stage.id !== "string" || typeof stage.name !== "string" + || stage.kind !== undefined && stage.kind !== "tool" + || !RUN_STAGE_STATUSES.has(stage.status as string) + || !seconds(stage.started, true) || !seconds(stage.ended, true) + ) + ) throw new Error("hosted channel has invalid workflow runs"); + } + return runs as unknown as NonNullable; +} + function restoreMcpUpdates(value: JsonValue | undefined): McpUpdateRecord[] { if (value === undefined) return []; if (!Array.isArray(value)) throw new Error("hosted channel has invalid MCP updates"); @@ -941,7 +999,7 @@ export async function open( presence: presence.create(), questions: Questions.restore(sidecar.openQuestions), comments: Comments.create(), - chat: Chat.restore(sidecar.transcript), + chat: Chat.restore(sidecar.transcript, sidecar.workflowRuns), outlines: new Map(), records: new Map(sidecar.questions.map(record => [record.id, record])), threads: new Map(sidecar.threads.map(record => [record.id, record])), diff --git a/apps/server/src/testing/decisions.ts b/apps/server/src/testing/decisions.ts index c8fe6082..da63f903 100644 --- a/apps/server/src/testing/decisions.ts +++ b/apps/server/src/testing/decisions.ts @@ -12,7 +12,10 @@ import type { DocumentRoom } from "../agent/tools"; import type { Socket, SocketData } from "../wire"; /** A room whose Decisions answer Atomic HostInput requests, with a member to answer them. */ -export async function hostInputRoom(expiresInMs?: number) { +export async function hostInputRoom( + expiresInMs?: number, + hold?: (workflowRunId: string, signal: AbortSignal) => Promise, +) { let storage = new MemoryStorage(); let now = new Date(); await storage.users.put({ id: "user", login: "reader", avatarUrl: "", now }); @@ -60,7 +63,7 @@ export async function hostInputRoom(expiresInMs?: number) { sessionId: "session", signal: controller.signal, }; - let input = createHumanInput(room, expiresInMs); + let input = createHumanInput(room, expiresInMs, hold); async function cards(count: number) { for (let deadline = Date.now() + 4_000; Date.now() < deadline; await Bun.sleep(5)) { if (Store.outstanding(plan.questions).length === count) { diff --git a/apps/web/src/assets/icons/planner-resume.svg b/apps/web/src/assets/icons/planner-resume.svg new file mode 100644 index 00000000..2443ffab --- /dev/null +++ b/apps/web/src/assets/icons/planner-resume.svg @@ -0,0 +1,7 @@ + + circle-play + diff --git a/apps/web/src/chat/chat.tsx b/apps/web/src/chat/chat.tsx index dd851054..70f2cc72 100644 --- a/apps/web/src/chat/chat.tsx +++ b/apps/web/src/chat/chat.tsx @@ -34,9 +34,11 @@ import { referenceTriggerKey, reviseComposerDraft, } from "./references"; +import { RunStack } from "./run-card"; import { Transcript } from "./transcript"; import { TerminalAlert } from "../terminal-alert"; import plannerStop from "../assets/icons/planner-stop.svg"; +import plannerResume from "../assets/icons/planner-resume.svg"; import type { Chat as Wire } from "@chopin/protocol"; import type { Repository } from "../api"; @@ -55,8 +57,19 @@ export type ChatProps = { agent?: boolean; active?: boolean; onActivity?: (event: { type: "message" | "working"; busy: boolean }) => void; + /** Opens Decisions, where a waiting workflow's questions are. */ + onShowDecisions?: () => void; }; +/** How many runs are still live, and how many are paused and resumable. */ +export function runCounts(runs: Wire.Runs | undefined): { active: number; paused: number } { + let list = runs ?? []; + return { + active: list.filter(run => run.status === "running" || run.status === "waiting").length, + paused: list.filter(run => run.status === "paused").length, + }; +} + export function Chat( { active = true, @@ -64,6 +77,7 @@ export function Chat( connected, handle, onActivity, + onShowDecisions, referencesEnabled, repository, room, @@ -75,6 +89,8 @@ export function Chat( let [queue, setQueue] = useState([]); let [busy, setBusy] = useState(false); let [turn, setTurn] = useState(); + let [runs, setRuns] = useState(); + let counts = runCounts(runs); let [draft, setDraft] = useState({ text: "", references: [], @@ -142,6 +158,7 @@ export function Chat( setBusy(frame.busy); reportedBusy.current = frame.busy; setTurn(frame.turn); + setRuns(frame.runs); // History is not unread, but a turn already in progress still needs // a signal outside a closed Chat destination. activity.current?.({ type: "working", busy: frame.busy }); @@ -192,6 +209,7 @@ export function Chat( reportedBusy.current = frame.busy; setBusy(frame.busy); setTurn(frame.turn); + setRuns(frame.runs); }), wire.on("chat:queue", frame => setQueue(frame.waiting)), ]; @@ -284,6 +302,19 @@ export function Chat( : undefined} /> + {agent && !!runs?.length && ( +
+ + wire?.send("chat:pause-run", { runId })} + onResume={runId => + wire?.send("chat:resume-run", { runId })} + onShowDecisions={onShowDecisions} + runs={runs} + /> +
+ )} +
{pickerOpen && trigger && ( )} - {agent && busy && ( + {agent && (busy || counts.active > 0) && ( )} + {agent && !busy && !counts.active && counts.paused > 0 && ( + + )} { + let runs = [ + run("old-done", "finished", 100), + run("live", "running", 300), + run("held", "paused", 400), + run("asking", "waiting", 200, 1), + run("new-done", "stopped", 500), + run("newer-live", "running", 600), + ]; + expect(orderRuns(runs).map(item => item.id)).toEqual([ + "asking", + "newer-live", + "live", + "held", + "new-done", + "old-done", + ]); + expect(defaultOpen(orderRuns(runs))).toBe("asking"); + expect(defaultOpen(orderRuns(runs.filter(item => item.status !== "waiting")))).toBe("newer-live"); + expect(defaultOpen(orderRuns([run("done", "finished", 1)]))).toBeUndefined(); +}); + +test("folding keeps at least three rows and never hides a run waiting on people", () => { + let many = orderRuns([ + run("a", "waiting", 1, 1), + run("b", "waiting", 2, 1), + run("c", "waiting", 3, 1), + run("d", "waiting", 4, 1), + run("e", "running", 5), + ]); + expect(visibleRuns(many, false).map(item => item.id)).toEqual(["d", "c", "b", "a"]); + expect(visibleRuns(many, true)).toHaveLength(5); + expect(visibleRuns(orderRuns([run("x", "running", 1), run("y", "paused", 2)]), false)) + .toHaveLength(2); +}); + +test("a concurrent stack opens only the waiting run, offers per-run controls, and folds the rest", () => { + let markup = renderToStaticMarkup(createElement(RunStack, { + onPause: () => {}, + onResume: () => {}, + onShowDecisions: () => {}, + runs: [ + run("drafting", "running", 300), + run("held", "paused", 200), + run("asking", "waiting", 100, 2), + run("done", "finished", 50), + ], + })); + expect(markup.match(/data-run-stage=/g)).toHaveLength(1); + expect(markup).toContain("asking-stage"); + expect(markup).not.toContain("drafting-stage"); + expect(markup).toContain('aria-label="Pause asking"'); + expect(markup).toContain('aria-label="Pause drafting"'); + expect(markup).toContain('aria-label="Resume held"'); + expect(markup).toContain("Waiting on 2 Decisions"); + expect(markup).toContain("+1 more"); + expect(markup).not.toContain(">done<"); + expect(markup.indexOf('data-run-status="waiting"')).toBeLessThan( + markup.indexOf('data-run-status="running"'), + ); +}); + +test("a tool step carries the wrench and says it is a tool step", () => { + let tooled = run("demo", "running", 10); + tooled.stages = [{ + id: "demo:t1", + name: "prepare", + kind: "tool", + status: "running", + started: 10, + }]; + let markup = renderToStaticMarkup(createElement(RunStack, { runs: [tooled] })); + expect(markup).toContain('tool step: prepare'); + expect(markup).toContain("data-nucleo-icon"); + let plain = renderToStaticMarkup( + createElement(RunStack, { runs: [run("agent", "running", 10)] }), + ); + expect(plain).not.toContain("tool step"); +}); + +test("runs of the same workflow are told apart by a short run id", () => { + let first = { ...run("aaaaaaaa-1111", "running", 10), name: "demo-wait" }; + let second = { ...run("bbbbbbbb-2222", "running", 20), name: "demo-wait" }; + let markup = renderToStaticMarkup( + createElement(RunStack, { onPause: () => {}, runs: [first, second] }), + ); + expect(markup).toContain('aria-label="Pause demo-wait aaaaaaaa"'); + expect(markup).toContain('aria-label="Pause demo-wait bbbbbbbb"'); + let single = renderToStaticMarkup(createElement(RunStack, { onPause: () => {}, runs: [first] })); + expect(single).toContain('aria-label="Pause demo-wait"'); +}); diff --git a/apps/web/src/chat/run-card.tsx b/apps/web/src/chat/run-card.tsx new file mode 100644 index 00000000..1e738ea3 --- /dev/null +++ b/apps/web/src/chat/run-card.tsx @@ -0,0 +1,306 @@ +/** Planner-launched workflow runs: which are alive, which need someone, and how the last ones ended. */ + +import { useEffect, useState } from "react"; +import { CheckIcon, ChevronIcon, CloseIcon, WarningIcon, WrenchIcon } from "@chopin/icons"; + +import plannerResume from "../assets/icons/planner-resume.svg"; +import plannerStop from "../assets/icons/planner-stop.svg"; + +import type { Chat as Wire } from "@chopin/protocol"; + +const SHOWN_STAGES = 6; +const SHOWN_RUNS = 3; + +const STATUS: Record = { + running: "running", + waiting: "waiting on Decisions", + paused: "paused", + finished: "finished", + blocked: "blocked", + failed: "failed", + stopped: "stopped", +}; + +const STAGE: Record = { + pending: "pending", + running: "running", + awaiting_input: "waiting on Decisions", + paused: "paused", + blocked: "blocked", + completed: "done", + failed: "failed", + skipped: "skipped", +}; + +/** Waiting runs first, then running, paused, and ended; newest first within each. */ +const GROUP: Record = { + waiting: 0, + running: 1, + paused: 2, + finished: 3, + blocked: 3, + failed: 3, + stopped: 3, +}; + +/** "45s", "6m", "1h 4m". */ +export function elapsed(seconds: number): string { + let s = Math.max(0, Math.floor(seconds)); + if (s < 60) return `${s}s`; + let m = Math.floor(s / 60); + if (m < 60) return `${m}m`; + return `${Math.floor(m / 60)}h ${m % 60}m`; +} + +export function orderRuns(runs: readonly Wire.Run[]): Wire.Run[] { + return [...runs].sort((a, b) => GROUP[a.status] - GROUP[b.status] || b.started - a.started); +} + +/** The run whose stages show without a click: the first one waiting on people, else the newest live one. */ +export function defaultOpen(ordered: readonly Wire.Run[]): string | undefined { + return ordered.find(run => run.status === "waiting")?.id + ?? ordered.find(run => run.status === "running")?.id; +} + +/** The rows shown before "more": at least the first few, and never fewer than every waiting run. */ +export function visibleRuns(ordered: readonly Wire.Run[], showAll: boolean): Wire.Run[] { + if (showAll) return [...ordered]; + let waiting = ordered.filter(run => run.status === "waiting").length; + return ordered.slice(0, Math.max(SHOWN_RUNS, waiting)); +} + +function useNow(live: boolean): number { + let [now, setNow] = useState(() => Date.now() / 1000); + useEffect(() => { + if (!live) return; + let timer = setInterval(() => setNow(Date.now() / 1000), 5_000); + return () => clearInterval(timer); + }, [live]); + return now; +} + +/** A dot that pulses while work is live; it stays still under reduced motion. */ +function RunPulse() { + return