From 2671c29c601636675023eb77ee7e1b523513dc51 Mon Sep 17 00:00:00 2001 From: Alex Lavaee Date: Wed, 30 Sep 2026 02:13:49 -0700 Subject: [PATCH 01/13] Keep the Planner while its workflows run, and let Stop pause them Chopin destroyed every Planner session when its turn ended. A workflow the atomic Planner starts runs in the background under that session, so it was orphaned seconds after launch and never progressed. Keep an atomic Planner session, its owner binding, and the loaded document while the session still owns live or paused workflow runs, tracked through Atomic's workflow activity stream; the next turn reuses it so the Planner can still steer its runs. The session is let go when the runs finish, the owner binding ends, or the document closes. Stop Planner now also pauses the session's live runs, resumably, and a new Resume Planner control (`chat:resume`) resumes them. Chat state carries the running and paused counts, and the composer shows them. Run control goes through the session's own workflow commands until the SDK exposes run-control primitives. The document's working directory without a checkout moves from the shared temporary directory to a per-channel directory under Chopin's per-user state directory, so workflow files survive restarts. Assistant-model: Claude Opus 5.5 Assistant-workflow: inline Assistant-verification: bun test passed: 1751 pass, 2 PostgreSQL skips, 0 fail, including new retention, pause, resume, release, and run-classification tests Assistant-verification: bun run types passed: all workspaces Assistant-verification: bun run ci passed: dprint, oxlint, tokens, type scale, Impeccable (no new findings) Assistant-verification: source check passed: DBOS records showed a Planner-launched run with no progress after its turn ended, matching the per-turn session disposal User-preference: A pause in the Planner pauses its workflows and resume resumes everything; steering goes to the main chat, which decides whether to steer a workflow over intercom User-preference: Keep long-lived Planner working directories out of the shared tmp Co-authored-by: Alex Lavaee --- apps/server/src/agent/planner.ts | 4 +- apps/server/src/chat/invoke.test.ts | 141 ++++++++++++++++++- apps/server/src/chat/service.ts | 141 +++++++++++++++++-- apps/server/src/harness/atomic/adapter.ts | 12 +- apps/server/src/harness/atomic/full.test.ts | 22 ++- apps/server/src/harness/atomic/full.ts | 60 +++++++- apps/server/src/harness/atomic/workspace.ts | 58 ++++---- apps/server/src/harness/harnesses.ts | 4 +- apps/server/src/harness/session.test.ts | 41 ++++-- apps/server/src/harness/session.ts | 38 ++++- apps/server/src/main.ts | 11 ++ apps/web/src/assets/icons/planner-resume.svg | 7 + apps/web/src/chat/chat.tsx | 45 +++++- apps/web/src/design-audit/icons.tsx | 2 + apps/web/src/tokens.test.ts | 6 + docs/hosted-agent.md | 26 +++- docs/local-agent-mcp.md | 4 +- docs/self-hosting.md | 12 +- packages/protocol/chat.d.ts | 16 ++- 19 files changed, 568 insertions(+), 82 deletions(-) create mode 100644 apps/web/src/assets/icons/planner-resume.svg diff --git a/apps/server/src/agent/planner.ts b/apps/server/src/agent/planner.ts index bb2aba7d..20b3a293 100644 --- a/apps/server/src/agent/planner.ts +++ b/apps/server/src/agent/planner.ts @@ -249,8 +249,8 @@ 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 diff --git a/apps/server/src/chat/invoke.test.ts b/apps/server/src/chat/invoke.test.ts index a20c02a5..804eacae 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,7 @@ 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 { HostedAuth } from "../auth/routes"; import type { Config } from "../config"; import type { GitHub } from "../github/client"; @@ -23,9 +24,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 +418,129 @@ test("HARNESS=atomic runs every Planner session full in hosted and local configu } } }); + +/** A Planner session that owns workflow runs, driven by the test. */ +function runningSession() { + let runs = { active: ["run-1"], paused: [] as string[] }; + let listeners = new Set<(runs: { active: string[]; paused: string[] }) => 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: { active: string[]; paused: string[] }) => void) { + listeners.add(listener); + return () => listeners.delete(listener); + }, + async pauseRuns() { + calls.push("pause"); + }, + async resumeRuns() { + calls.push("resume"); + }, + }; + return { + session, + calls, + set(next: { active: string[]; paused: string[] }) { + 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; + + 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({ active: 1, paused: 0 }); + expect(runs()).toEqual({ active: 1, paused: 0 }); + 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"] }); + expect(runs()).toEqual({ active: 0, paused: 1 }); + 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: [] }); + + 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: [] }); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(planner.calls.at(-1)).toBe("destroy"); + expect(context.chat.retained).toBeUndefined(); + expect(context.chat.runs).toBeUndefined(); + expect(runs()).toBeUndefined(); + 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: () => ({ active: [], paused: [] }) }; + 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(); +}); diff --git a/apps/server/src/chat/service.ts b/apps/server/src/chat/service.ts index c88749af..d4e9cd45 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; @@ -343,6 +347,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 +400,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 +433,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 +788,111 @@ 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() }); +} + +function publishRuns( + context: Room, + runs: { active: string[]; paused: string[] } | undefined, +): void { + let { chat, room, server } = context; + chat.runs = runs && (runs.active.length || runs.paused.length) + ? { active: runs.active.length, paused: runs.paused.length } + : undefined; + state(chat, server, room); +} + +/** + * 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 +1044,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 +1267,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 +1562,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..b0d26233 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,11 @@ export function createAtomicAdapter( customTools: turn.tools.map(hostTool), }); let session = created.session; + if (full) { + full.command = async text => { + await session.prompt(text); + }; + } 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..42a8e23c 100644 --- a/apps/server/src/harness/atomic/full.test.ts +++ b/apps/server/src/harness/atomic/full.test.ts @@ -10,7 +10,7 @@ import { ATOMIC_RESULT_TOOL_NAME, createAtomicAdapter, } from "./adapter"; -import { registerFullPlanner } from "./full"; +import { classifyRuns, registerFullPlanner } from "./full"; import { startStubModelServer } from "../pi/model-stub"; import { hostInputRoom } from "../../testing/decisions"; import type { HostInput, QuestionParams } from "@bastani/atomic"; @@ -306,3 +306,23 @@ 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("working and blocked runs are live, paused idle runs are paused, and finished runs are neither", () => { + 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("held", "idle", "paused"), + root("done", "idle", "quiescent"), + ])).toEqual({ active: ["drafting", "asking"], paused: ["held"] }); + expect(classifyRuns([])).toEqual({ active: [], paused: [] }); +}); diff --git a/apps/server/src/harness/atomic/full.ts b/apps/server/src/harness/atomic/full.ts index fc79c5a5..2f78a8cb 100644 --- a/apps/server/src/harness/atomic/full.ts +++ b/apps/server/src/harness/atomic/full.ts @@ -1,6 +1,24 @@ -import type { HostInput } from "@bastani/atomic"; +import type { + ExtensionAPI, + ExtensionFactory, + HostInput, + WorkflowActivitySubscription, + WorkflowRootActivity, +} from "@bastani/atomic"; -export type FullPlanner = { cwd: string; humanInput: HostInput }; +/** Root workflow runs a Planner session owns: still live, or paused and resumable. */ +export type PlannerRuns = { active: string[]; paused: string[] }; + +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; + /** Runs a slash command in the live session without a model turn; set once the session exists. */ + command?: (text: string) => Promise; +}; let planners = new Map(); /** Every Atomic Planner session is registered; background workers never are. */ @@ -15,3 +33,41 @@ export function registerFullPlanner(sessionId: string, planner: FullPlanner): () export function fullPlanner(sessionId: string): FullPlanner | undefined { return planners.get(sessionId); } + +/** Working and blocked roots are live; an idle root is paused only when Atomic says so. */ +export function classifyRuns(roots: Iterable): PlannerRuns { + let runs: PlannerRuns = { active: [], paused: [] }; + for (let root of roots) { + if (root.state === "working" || root.state === "blocked") runs.active.push(root.rootRunId); + else if (root.reason === "paused") runs.paused.push(root.rootRunId); + } + return runs; +} + +/** + * Keeps `planner.runs` in step with the workflows this session owns, through + * Atomic's public activity stream. 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 publish = () => { + planner.runs = classifyRuns(roots.values()); + planner.onRuns?.(planner.runs); + }; + atomic.on("session_start", (_event, ctx) => { + lease?.dispose(); + lease = ctx.observeWorkflowActivity(frame => { + if (frame.kind === "snapshot") { + if (frame.availability !== "ready") return; + 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); + publish(); + }); + }); + atomic.on("session_shutdown", () => lease?.dispose()); + }; +} 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..4074c8ee 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 }); }); @@ -185,27 +193,40 @@ 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"); } - 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 () => { diff --git a/apps/server/src/harness/session.ts b/apps/server/src/harness/session.ts index d02751a9..0de0f870 100644 --- a/apps/server/src/harness/session.ts +++ b/apps/server/src/harness/session.ts @@ -2,7 +2,7 @@ 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, type PlannerRuns, registerFullPlanner } from "./atomic/full"; import { createHumanInput } from "./atomic/human-input"; import { type PlannerWorkspace, plannerWorkspace } from "./atomic/workspace"; @@ -33,6 +33,14 @@ export type PlannerChannel = { 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; }; export type PlannerSessionDependencies = { @@ -82,11 +90,17 @@ 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), - }); + onRuns: runs => { + for (let listener of listeners) listener(runs); + }, + }; + unregisterFull = registerFullPlanner(sessionId, planner); } let instructions = typeof channel.instructions === "function" ? channel.instructions(workspace) @@ -109,9 +123,27 @@ export async function openPlannerSession( let active = session; let release = unregister; let stopped: Promise | undefined; + let command = async (text: string) => { + if (!planner?.command) throw new Error("The Planner has no live Atomic session to control."); + await planner.command(text); + }; + let control = planner + ? { + runs: () => planner.runs, + watchRuns: (listener: (runs: PlannerRuns) => void) => { + listeners.add(listener); + return () => listeners.delete(listener); + }, + pauseRuns: () => command("/workflow pause --all"), + resumeRuns: async () => { + for (let runId of planner.runs?.paused ?? []) await command(`/workflow resume ${runId}`); + }, + } + : {}; return { ok: true, value: { + ...control, stream: (prompt, abortSignal) => agent.stream({ session: active, diff --git a/apps/server/src/main.ts b/apps/server/src/main.ts index bf4b4c77..573b69d7 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,10 @@ 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:unqueue": if (room.plan) Chat.unqueue(chat(room, ws), ws, frame); return; 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..e73af295 100644 --- a/apps/web/src/chat/chat.tsx +++ b/apps/web/src/chat/chat.tsx @@ -37,6 +37,7 @@ import { 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"; @@ -57,6 +58,15 @@ export type ChatProps = { onActivity?: (event: { type: "message" | "working"; busy: boolean }) => void; }; +/** "1 workflow running", "2 workflows paused", or both. */ +export function runsLabel(runs: Wire.Runs): string { + let count = (n: number, word: string) => `${n} ${n === 1 ? "workflow" : "workflows"} ${word}`; + return [ + runs.active ? count(runs.active, "running") : "", + runs.paused ? count(runs.paused, "paused") : "", + ].filter(Boolean).join(" · "); +} + export function Chat( { active = true, @@ -75,6 +85,7 @@ export function Chat( let [queue, setQueue] = useState([]); let [busy, setBusy] = useState(false); let [turn, setTurn] = useState(); + let [runs, setRuns] = useState(); let [draft, setDraft] = useState({ text: "", references: [], @@ -142,6 +153,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 +204,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)), ]; @@ -407,12 +420,21 @@ export function Chat( />
- {sendError && ( - - {sendError} - - )} - {agent && busy && ( + {sendError + ? ( + + {sendError} + + ) + : agent && runs && ( + + {runsLabel(runs)} + + )} + {agent && (busy || !!runs?.active) && ( )} + {agent && !busy && !runs?.active && !!runs?.paused && ( + + )} { size: "btn-icon", tiers: ["btn-secondary"], }], + ["apps/web/src/chat/chat.tsx", { + action: "Resume Planner", + marker: 'wire?.send("chat:resume")', + size: "btn-icon", + tiers: ["btn-secondary"], + }], ["packages/editor/src/send-action.tsx", { action: "shared send action", marker: "aria-label={label}", diff --git a/docs/hosted-agent.md b/docs/hosted-agent.md index e69e329c..cf41e518 100644 --- a/docs/hosted-agent.md +++ b/docs/hosted-agent.md @@ -131,7 +131,7 @@ agent directory (including `ATOMIC_CODING_AGENT_DIR` or the legacy `PI_CODING_AGENT_DIR` override) supplies extensions, skills, prompt templates, context files, and a read-only copy of its settings. A verified checkout's `.atomic/settings.json` is read the same way, as trusted project settings, so -packages installed for that project load too; the empty per-channel directory +packages installed for that project load too; the document's own directory has none. Chopin's compaction, summary, and cache overrides still apply on top. Chopin appends its Planner instructions to Atomic's assembled prompt. Session, settings, and model @@ -152,10 +152,13 @@ and every later Planner session for the document, browser-started or MCP-started re-verifies it before using it as its working directory. Without a remembered checkout that still verifies, the session runs with the -same tools in an empty working directory Chopin creates for that document alone, -under the operating system's temporary directory with mode `0700`. It is never -shared between documents, keeps the Planner's own files for later sessions of the -same document, and is removed when the server shuts down cleanly. The Planner's +same tools in a directory Chopin keeps for that document alone under its per-user state +directory (`$XDG_STATE_HOME/chopin/planner/`, defaulting to +`~/.local/state`; `~/Library/Application Support/Chopin/planner` on macOS; +`%LOCALAPPDATA%\Chopin\planner` on Windows), with mode `0700`. It is never shared +with another document, keeps what the Planner and its workflows write there +across later sessions and server restarts, and is not placed in the shared +temporary directory. A symlink or file at that path is refused. The Planner's instructions state which case applies: a verified checkout, or an empty directory with no repository files, in which case it reads the repository through Chopin's repository tools. @@ -264,6 +267,19 @@ A harness session is disposable. A process restart, credential rotation, logout, or ownership reset discards it. A later turn bootstraps from the bounded transcript and reads the current document. +The one exception is an atomic Planner session that still owns Atomic workflow +runs when its turn ends. A workflow the Planner starts outlives the turn that +launched it, and its runs belong to that session, so Chopin keeps the session +and its owner binding, and keeps the document loaded, until every run has +finished. The next turn reuses that session, so the Planner can still see and +steer its runs over Intercom; people steer through Chat and Decisions, never a +stage directly. **Stop Planner** also pauses the session's live runs, resumably, +and **Resume Planner** resumes the runs it paused. Chat shows how many runs are +running or paused. The session is let go when its runs finish, when its owner +binding ends, or when the document closes; run state is durable in Atomic's +workflow store, so runs interrupted that way can be resumed from a later +session. + An interrupted turn is visible and is never replayed automatically because it may already have made durable document or question changes. `HarnessAgent.stream()` returns an AI SDK stream that `chat/service.ts` consumes part by part; the Chat diff --git a/docs/local-agent-mcp.md b/docs/local-agent-mcp.md index 97296170..cff41ec9 100644 --- a/docs/local-agent-mcp.md +++ b/docs/local-agent-mcp.md @@ -115,8 +115,8 @@ server's machine. Chopin verifies that its `origin` matches the document repository and refuses the request before posting if it does not. A verified path is remembered for the document until the server restarts, and every later Planner session for it, from the browser or through MCP, works there after -checking the path again. Without one, the Planner works in an empty directory -Chopin keeps for that document. Other harnesses ignore `checkout` entirely: it is +checking the path again. Without one, the Planner works in a directory +Chopin keeps for that document under its per-user state directory. Other harnesses ignore `checkout` entirely: it is neither verified, used, nor remembered. The result contains `id`, `title`, and the canonical `url` (plus any generated diff --git a/docs/self-hosting.md b/docs/self-hosting.md index 4bcdde7d..40567a2d 100644 --- a/docs/self-hosting.md +++ b/docs/self-hosting.md @@ -215,11 +215,13 @@ for the document, whether started from the browser or through MCP, re-verifies the remembered path before using it. The check establishes repository coordinates, not remote-host authenticity, and it does not confine the shell. -Without a remembered checkout that still verifies, the session runs in an empty -working directory that Chopin creates for that document alone under the operating -system's temporary directory, with mode `0700`. It is never shared with another -document, keeps what the Planner writes there for later sessions of the same -document, and is removed when the server shuts down cleanly. The session has the +Without a remembered checkout that still verifies, the session runs in a directory Chopin keeps for that document alone under its per-user state +directory (`$XDG_STATE_HOME/chopin/planner/`, defaulting to +`~/.local/state`; `~/Library/Application Support/Chopin/planner` on macOS; +`%LOCALAPPDATA%\Chopin\planner` on Windows), with mode `0700`. It is never shared +with another document, keeps what the Planner and its workflows write there +across later sessions and server restarts, and is not placed in the shared +temporary directory. A symlink or file at that path is refused. The session has the same tools either way; the Planner is told whether it is in a checkout or in an empty directory without repository files, and its repository tools remain available. diff --git a/packages/protocol/chat.d.ts b/packages/protocol/chat.d.ts index 4ae29cd8..dd4af91e 100644 --- a/packages/protocol/chat.d.ts +++ b/packages/protocol/chat.d.ts @@ -14,6 +14,7 @@ export declare namespace Chat { export type Incoming = | Request | Request + | Request | Request; export type Outgoing = History | Message | Delta | Tool | State | Queue | Sent; @@ -87,11 +88,15 @@ export declare namespace Chat { responded: boolean; }; + /** Workflow runs the Planner is still carrying after its turn: live, and paused. */ + export type Runs = { active: number; paused: number }; + /** Everything said so far, sent on join. */ export type History = KIND<"chat:history"> & { entries: Entry[]; busy: boolean; turn?: Turn; + runs?: Runs; queued: Waiting[]; }; @@ -104,10 +109,11 @@ export declare namespace Chat { /** A tool call starting, or finishing. */ export type Tool = KIND<"chat:tool"> & { entry: string; activity: Activity }; - /** Whether a turn is running, and for whom. */ + /** Whether a turn is running, and for whom, and any workflow runs outliving it. */ export type State = KIND<"chat:state"> & { busy: boolean; turn?: Turn; + runs?: Runs; }; /** A message waiting for the current turn to end. */ @@ -138,9 +144,15 @@ export declare namespace Chat { /** The member message or queue entry is accepted by the server. */ export type Sent = KIND<"chat:send"> & { id: string; queued: boolean }; - /** Stop the running turn. Anyone may; the transcript records who did. */ + /** + * Stop the running turn and pause the Planner's workflow runs. Anyone may; + * the transcript records who did. + */ export type Abort = KIND<"chat:abort">; + /** Resume the workflow runs the Planner paused. Anyone may; the transcript records who did. */ + export type Resume = KIND<"chat:resume">; + /** Withdraw a queued message. Only its author may. */ export type Unqueue = KIND<"chat:unqueue"> & { id: string }; } From ec6645a07857f47cf35ead93ba8f9e0363c77443 Mon Sep 17 00:00:00 2001 From: Alex Lavaee Date: Wed, 30 Sep 2026 08:29:21 -0700 Subject: [PATCH 02/13] Report a taken title from create_document and declare its outcomes Document titles are unique per repository, so creating a second document with an existing title failed in storage and came back as idempotency-conflict, which reads as a reused key. Report it as title-taken instead. create_document's output schema also described only success, so strict MCP clients rejected its idempotency-conflict, document-unavailable, and validation-issue responses as schema violations and callers never saw the real outcome; declare them the way update_document does. Assistant-model: Claude Opus 5.5 Assistant-workflow: inline Assistant-verification: bun test passed: 1752 pass, 2 PostgreSQL skips, 0 fail, including a new hosted test for title-taken Assistant-verification: bun run types passed: all workspaces Assistant-verification: bun run ci passed: dprint, oxlint, tokens, type scale, Impeccable (no new findings) Assistant-verification: raw JSON-RPC probe passed: a fresh key with an existing title returned idempotency-conflict before this change, and a new title succeeded Co-authored-by: Alex Lavaee --- apps/server/src/mcp.ts | 72 ++++++++++++++++++------------ apps/server/src/mcp/hosted.test.ts | 22 +++++++++ apps/server/src/mcp/hosted.ts | 6 +-- docs/local-agent-mcp.md | 6 +++ 4 files changed, 74 insertions(+), 32 deletions(-) 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/docs/local-agent-mcp.md b/docs/local-agent-mcp.md index cff41ec9..a9f7f3da 100644 --- a/docs/local-agent-mcp.md +++ b/docs/local-agent-mcp.md @@ -141,6 +141,12 @@ Refusals return an object with `code`: - `document-unavailable`, `document-archived`, and `repository-forbidden` retain their ordinary document/access meanings. +Document titles are unique within a repository, ignoring case. `create_document` +returns `title-taken` when the repository already has a document with that title; +choose another title and retry under a new idempotency key. `idempotency-conflict` +means the same key was reused with a different request, and validation failures +return `issues`. Every one of these outcomes matches the tool's output schema. + ## Document URLs and IDs The `url` returned by `create_document` is the readable canonical route: From bc123227e04fc08a295f74067f5d9d469ab11d48 Mon Sep 17 00:00:00 2001 From: Alex Lavaee Date: Wed, 30 Sep 2026 08:45:14 -0700 Subject: [PATCH 03/13] Show each Planner workflow run as a card in Chat, and pause without a prompt Replace the running/paused count with a card per run, built from Atomic's workflow lifecycle hooks: the run's name and status (running, waiting on Decisions, paused, or ended), elapsed time, each stage it reached with its status and duration, and a link to the Decisions it waits on. Chat records a line when a run ends. Stop now passes --yes to the pause command: without it, Atomic asked the host to confirm, which Chopin routes to Decisions, so the pause waited on a question instead of pausing. A design-audit specimen covers the running, waiting, paused, and finished states. Assistant-model: Claude Opus 5.5 Assistant-workflow: inline Assistant-verification: bun test passed: 1753 pass, 2 PostgreSQL skips, 0 fail, including run-card folding, paused/waiting/finished status, and the run-ended transcript line Assistant-verification: bun run types passed: all workspaces Assistant-verification: bun run ci passed: dprint, oxlint, tokens, type scale, Impeccable (no new findings) Assistant-verification: agent-browser E2E failed before this change: Stop Planner on a run waiting in Decisions left "1 workflow running" because the pause waited on its own confirmation User-preference: Show the workflow's current stages in Chopin, and have pause or quit change what Chat shows immediately Co-authored-by: Alex Lavaee --- apps/server/src/chat/invoke.test.ts | 58 ++++++-- apps/server/src/chat/service.ts | 30 ++++- apps/server/src/harness/atomic/full.test.ts | 93 +++++++++++++ apps/server/src/harness/atomic/full.ts | 132 ++++++++++++++++++- apps/server/src/harness/session.ts | 2 +- apps/web/src/chat/chat.tsx | 48 +++---- apps/web/src/chat/run-card.tsx | 139 ++++++++++++++++++++ apps/web/src/design-audit/inventory.ts | 6 + apps/web/src/design-audit/surfaces.tsx | 64 +++++++++ apps/web/src/room-workspace.tsx | 1 + docs/hosted-agent.md | 7 +- packages/protocol/chat.d.ts | 35 ++++- 12 files changed, 564 insertions(+), 51 deletions(-) create mode 100644 apps/web/src/chat/run-card.tsx diff --git a/apps/server/src/chat/invoke.test.ts b/apps/server/src/chat/invoke.test.ts index 804eacae..0fba6018 100644 --- a/apps/server/src/chat/invoke.test.ts +++ b/apps/server/src/chat/invoke.test.ts @@ -16,6 +16,7 @@ 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"; @@ -419,10 +420,31 @@ 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 = { active: ["run-1"], paused: [] as string[] }; - let listeners = new Set<(runs: { active: string[]; paused: string[] }) => void>(); + 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) { @@ -437,7 +459,7 @@ function runningSession() { calls.push("destroy"); }, runs: () => runs, - watchRuns(listener: (runs: { active: string[]; paused: string[] }) => void) { + watchRuns(listener: (runs: Runs) => void) { listeners.add(listener); return () => listeners.delete(listener); }, @@ -451,7 +473,7 @@ function runningSession() { return { session, calls, - set(next: { active: string[]; paused: string[] }) { + set(next: Runs) { runs = next; for (let listener of listeners) listener(next); }, @@ -473,27 +495,34 @@ test("a Planner that still owns workflow runs outlives its turn, pauses, resumes }; 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; + 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({ active: 1, paused: 0 }); - expect(runs()).toEqual({ active: 1, paused: 0 }); + 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"] }); - expect(runs()).toEqual({ active: 0, paused: 1 }); + 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: [] }); + 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; @@ -501,7 +530,9 @@ test("a Planner that still owns workflow runs outlives its turn, pauses, resumes expect(planner.calls.filter(call => call.startsWith("stream"))).toHaveLength(2); expect(planner.calls).not.toContain("destroy"); - planner.set({ active: [], paused: [] }); + 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(); @@ -533,7 +564,10 @@ test("a retained Planner is let go when its owner's binding ends, and sessions w expect(context.chat.retained).toBeUndefined(); expect(planner.calls.at(-1)).toBe("destroy"); - let quiet = { ...runningSession().session, runs: () => ({ active: [], paused: [] }) }; + let quiet = { + ...runningSession().session, + runs: (): Runs => ({ active: [], paused: [], cards: [] }), + }; let destroyed = 0; quiet.destroy = async () => { destroyed++; diff --git a/apps/server/src/chat/service.ts b/apps/server/src/chat/service.ts index d4e9cd45..fef62d39 100644 --- a/apps/server/src/chat/service.ts +++ b/apps/server/src/chat/service.ts @@ -833,14 +833,36 @@ export async function resume(context: Room, ws: Socket): Promise { say(chat, server, room, { id: ulid(), author: { kind: "system" }, text, ts: now() }); } +const ENDED_RUN: Partial> = { + finished: "finished", + failed: "failed", + stopped: "was stopped", +}; + +/** + * Show the retained session's runs, and say once in the transcript when one + * ends, so the record outlives the card. + */ function publishRuns( context: Room, - runs: { active: string[]; paused: string[] } | undefined, + runs: { active: string[]; paused: string[]; cards?: Wire.Run[] } | undefined, ): void { let { chat, room, server } = context; - chat.runs = runs && (runs.active.length || runs.paused.length) - ? { active: runs.active.length, paused: runs.paused.length } - : undefined; + let before = new Map((chat.runs ?? []).map(run => [run.id, run.status])); + let cards = runs?.cards ?? []; + 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 = runs && (runs.active.length || runs.paused.length) ? cards : undefined; state(chat, server, room); } diff --git a/apps/server/src/harness/atomic/full.test.ts b/apps/server/src/harness/atomic/full.test.ts index 42a8e23c..6663cc36 100644 --- a/apps/server/src/harness/atomic/full.test.ts +++ b/apps/server/src/harness/atomic/full.test.ts @@ -326,3 +326,96 @@ test("working and blocked runs are live, paused idle runs are paused, and finish ])).toEqual({ active: ["drafting", "asking"], 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, + }); +}); diff --git a/apps/server/src/harness/atomic/full.ts b/apps/server/src/harness/atomic/full.ts index 2f78a8cb..434dc643 100644 --- a/apps/server/src/harness/atomic/full.ts +++ b/apps/server/src/harness/atomic/full.ts @@ -3,11 +3,16 @@ import type { ExtensionFactory, HostInput, WorkflowActivitySubscription, + WorkflowLifecycleEvent, WorkflowRootActivity, } from "@bastani/atomic"; +import type { Chat as Wire } from "@chopin/protocol"; -/** Root workflow runs a Planner session owns: still live, or paused and resumable. */ -export type PlannerRuns = { active: string[]; paused: string[] }; +/** + * 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; @@ -35,8 +40,8 @@ export function fullPlanner(sessionId: string): FullPlanner | undefined { } /** Working and blocked roots are live; an idle root is paused only when Atomic says so. */ -export function classifyRuns(roots: Iterable): PlannerRuns { - let runs: PlannerRuns = { active: [], paused: [] }; +export function classifyRuns(roots: Iterable): Omit { + let runs = { active: [] as string[], paused: [] as string[] }; for (let root of roots) { if (root.state === "working" || root.state === "blocked") runs.active.push(root.rootRunId); else if (root.reason === "paused") runs.paused.push(root.rootRunId); @@ -44,17 +49,113 @@ export function classifyRuns(roots: Iterable): PlannerRuns return runs; } +type Card = Wire.Run & { stageIndex: Map; prompts: Set }; + +const ENDED: Record = { + completed: "finished", + skipped: "finished", + failed: "failed", + blocked: "failed", + killed: "stopped", + cancelled: "stopped", +}; + +/** 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) { + 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 === "failed" || card.status === "stopped") { + card.ended = seconds; + } else delete card.ended; + } else if (target.kind === "stage") { + let key = `${target.runId}:${target.stageId}`; + let index = card.stageIndex.get(key); + if (index === undefined) { + index = card.stages.length; + card.stageIndex.set(key, index); + card.stages.push({ id: key, name: target.stageName, status: target.status }); + } + let stage = card.stages[index]!; + stage.status = target.status; + if (target.status === "running" && stage.started === undefined) stage.started = seconds; + if ( + target.status === "completed" || target.status === "failed" || target.status === "skipped" + ) { + stage.ended = seconds; + } + } 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 status: Wire.Run["status"] = runs.paused.includes(card.id) + ? "paused" + : runs.active.includes(card.id) + ? card.waiting > 0 ? "waiting" : "running" + : card.status === "running" || card.status === "waiting" + ? "finished" + : card.status; + return { ...card, status, stages: card.stages.map(stage => ({ ...stage })) }; + }); +} + +/** 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. Recovering or unavailable snapshots are - * unknown rather than empty, so they leave the last known runs in place. + * 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 ready = false; let publish = () => { - planner.runs = classifyRuns(roots.values()); + if (!ready) return; + let split = classifyRuns(roots.values()); + for (let id of cards.keys()) { + if (!split.active.includes(id) && !split.paused.includes(id) && !roots.has(id)) { + let card = cards.get(id)!; + if (card.status === "running" || card.status === "waiting") card.status = "finished"; + } + } + planner.runs = { ...split, cards: runCards(cards, split) }; planner.onRuns?.(planner.runs); }; atomic.on("session_start", (_event, ctx) => { @@ -62,12 +163,29 @@ export function workflowRuns(planner: FullPlanner): ExtensionFactory { 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); 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 => { + foldLifecycle(cards, event, names.get(event.rootRunId)); + publish(); + }); atomic.on("session_shutdown", () => lease?.dispose()); }; } diff --git a/apps/server/src/harness/session.ts b/apps/server/src/harness/session.ts index 0de0f870..0271868e 100644 --- a/apps/server/src/harness/session.ts +++ b/apps/server/src/harness/session.ts @@ -134,7 +134,7 @@ export async function openPlannerSession( listeners.add(listener); return () => listeners.delete(listener); }, - pauseRuns: () => command("/workflow pause --all"), + pauseRuns: () => command("/workflow pause --all --yes"), resumeRuns: async () => { for (let runId of planner.runs?.paused ?? []) await command(`/workflow resume ${runId}`); }, diff --git a/apps/web/src/chat/chat.tsx b/apps/web/src/chat/chat.tsx index e73af295..f38f4d4e 100644 --- a/apps/web/src/chat/chat.tsx +++ b/apps/web/src/chat/chat.tsx @@ -34,6 +34,7 @@ import { referenceTriggerKey, reviseComposerDraft, } from "./references"; +import { RunCard } from "./run-card"; import { Transcript } from "./transcript"; import { TerminalAlert } from "../terminal-alert"; import plannerStop from "../assets/icons/planner-stop.svg"; @@ -56,15 +57,17 @@ 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; }; -/** "1 workflow running", "2 workflows paused", or both. */ -export function runsLabel(runs: Wire.Runs): string { - let count = (n: number, word: string) => `${n} ${n === 1 ? "workflow" : "workflows"} ${word}`; - return [ - runs.active ? count(runs.active, "running") : "", - runs.paused ? count(runs.paused, "paused") : "", - ].filter(Boolean).join(" · "); +/** 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( @@ -74,6 +77,7 @@ export function Chat( connected, handle, onActivity, + onShowDecisions, referencesEnabled, repository, room, @@ -86,6 +90,7 @@ export function Chat( let [busy, setBusy] = useState(false); let [turn, setTurn] = useState(); let [runs, setRuns] = useState(); + let counts = runCounts(runs); let [draft, setDraft] = useState({ text: "", references: [], @@ -297,6 +302,12 @@ export function Chat( : undefined} /> + {agent && !!runs?.length && ( +
+ {runs.map(run => )} +
+ )} +
{pickerOpen && trigger && (
- {sendError - ? ( - - {sendError} - - ) - : agent && runs && ( - - {runsLabel(runs)} - - )} - {agent && (busy || !!runs?.active) && ( + {sendError && ( + + {sendError} + + )} + {agent && (busy || counts.active > 0) && ( )} - {agent && !busy && !runs?.active && !!runs?.paused && ( + {agent && !busy && !counts.active && counts.paused > 0 && ( + )} + + ); +} diff --git a/apps/web/src/design-audit/inventory.ts b/apps/web/src/design-audit/inventory.ts index 25863b65..28e8c325 100644 --- a/apps/web/src/design-audit/inventory.ts +++ b/apps/web/src/design-audit/inventory.ts @@ -160,6 +160,12 @@ export const AUDIT_INVENTORY: readonly AuditGroup[] = [ source: "apps/web/src/chat/chat.tsx", states: ["member", "planner", "tool", "busy", "error"], }, + { + id: "workflow-runs", + label: "Workflow runs", + source: "apps/web/src/chat/run-card.tsx", + states: ["running", "waiting", "paused", "finished"], + }, { id: "identity", label: "Identity and collaboration", diff --git a/apps/web/src/design-audit/surfaces.tsx b/apps/web/src/design-audit/surfaces.tsx index 85de8217..7828246d 100644 --- a/apps/web/src/design-audit/surfaces.tsx +++ b/apps/web/src/design-audit/surfaces.tsx @@ -10,6 +10,7 @@ import { } from "@chopin/icons"; import { DecisionCard, PlanStatus, SendAction, SidecarCard } from "@chopin/editor"; +import { RunCard } from "../chat/run-card"; import { Transcript } from "../chat/transcript"; import { TerminalAlert } from "../terminal-alert"; import { AuditPlate, StateLabel } from "./frame"; @@ -213,6 +214,68 @@ function Conversation() { ); } +const RUN_STAGES: Chat.RunStage[] = [ + { id: "r:preflight", name: "preflight", status: "completed", started: 1_000, ended: 1_002 }, + { id: "r:draft-1", name: "draft-1", status: "completed", started: 1_002, ended: 1_380 }, + { id: "r:reviewer-a-1", name: "reviewer-a-1", status: "completed", started: 1_380, ended: 1_620 }, + { id: "r:draft-2", name: "draft-2", status: "running", started: 1_620 }, +]; + +function run(status: Chat.Run["status"], stages: Chat.RunStage[], waiting = 0): Chat.Run { + let ended = status === "finished" ? { ended: 2_100 } : {}; + return { + id: `audit-${status}`, + name: "plan-review", + status, + started: 1_000, + updated: 2_100, + ...ended, + stages, + waiting, + }; +} + +function WorkflowRuns() { + return ( + +
+ Running + + Waiting + {}} + run={run("waiting", [ + ...RUN_STAGES.slice(0, 3), + { id: "r:draft-2", name: "draft-2", status: "awaiting_input", started: 1_620 }, + ], 2)} + /> + Paused + + Finished + ({ + ...stage, + status: "completed", + ended: stage.ended ?? 2_100, + })), + )} + /> +
+
+ ); +} + function Decisions() { return ( <> @@ -394,6 +457,7 @@ export function Surfaces() { + diff --git a/apps/web/src/room-workspace.tsx b/apps/web/src/room-workspace.tsx index 6ea1a406..216f22ab 100644 --- a/apps/web/src/room-workspace.tsx +++ b/apps/web/src/room-workspace.tsx @@ -448,6 +448,7 @@ export function RoomWorkspace( connected={status === "connected" && workspaceCanEdit} handle={handle} onActivity={onChatActivity} + onShowDecisions={() => selectDestination("decisions")} referencesEnabled={chatReferences.wire === wire && chatReferences.enabled} repository={repository} room={room} diff --git a/docs/hosted-agent.md b/docs/hosted-agent.md index cf41e518..e71054b8 100644 --- a/docs/hosted-agent.md +++ b/docs/hosted-agent.md @@ -274,8 +274,11 @@ and its owner binding, and keeps the document loaded, until every run has finished. The next turn reuses that session, so the Planner can still see and steer its runs over Intercom; people steer through Chat and Decisions, never a stage directly. **Stop Planner** also pauses the session's live runs, resumably, -and **Resume Planner** resumes the runs it paused. Chat shows how many runs are -running or paused. The session is let go when its runs finish, when its owner +and **Resume Planner** resumes the runs it paused. Chat shows a card for each +run: its name, status (running, waiting on Decisions, paused, or ended), elapsed +time, every stage it has reached with its status and duration, and a link to the +Decisions it is waiting on. When a run ends, Chat records a line saying how it +ended and how long it took. The session is let go when its runs finish, when its owner binding ends, or when the document closes; run state is durable in Atomic's workflow store, so runs interrupted that way can be resumed from a later session. diff --git a/packages/protocol/chat.d.ts b/packages/protocol/chat.d.ts index dd4af91e..242a2709 100644 --- a/packages/protocol/chat.d.ts +++ b/packages/protocol/chat.d.ts @@ -88,8 +88,39 @@ export declare namespace Chat { responded: boolean; }; - /** Workflow runs the Planner is still carrying after its turn: live, and paused. */ - export type Runs = { active: number; paused: number }; + /** One stage of a workflow run, in the order the run reached it. */ + export type RunStage = { + id: string; + name: string; + status: + | "pending" + | "running" + | "awaiting_input" + | "paused" + | "blocked" + | "completed" + | "failed" + | "skipped"; + /** Seconds since the epoch. */ + started?: number; + ended?: number; + }; + + /** An Atomic workflow run the Planner started and still carries after its turn. */ + export type Run = { + id: string; + name: string; + status: "running" | "waiting" | "paused" | "finished" | "failed" | "stopped"; + /** Seconds since the epoch. */ + started: number; + updated: number; + ended?: number; + stages: RunStage[]; + /** Decisions questions the run is waiting on. */ + waiting: number; + }; + + export type Runs = Run[]; /** Everything said so far, sent on join. */ export type History = KIND<"chat:history"> & { From 4bdc2b5add60bcf7774d65cbcfcb364e04edd387 Mon Sep 17 00:00:00 2001 From: Alex Lavaee Date: Wed, 30 Sep 2026 08:59:17 -0700 Subject: [PATCH 04/13] Pause and resume Planner workflows through the workflow tool, and pulse live runs Stop Planner reported that it paused the workflows, but they kept running: the /workflow slash command reports its outcome only through ui.notify, which a headless session drops, so a refused or no-op pause looked like success. Run pause and resume through the session's own workflow tool with a real tool context instead, so the outcome returns as a result and a failure reaches the chat transcript. The run card's spinner never animated because chat-tool-loader has no animation in the web app. Live runs and running stages now show a pulsing dot that stays still under reduced motion. Assistant-model: Claude Opus 5.5 Assistant-workflow: inline Assistant-verification: bun test passed: 1754 pass, 2 PostgreSQL skips, 0 fail, including run control through the workflow tool and its failure path Assistant-verification: bun run types passed: all workspaces Assistant-verification: bun run ci passed: dprint, oxlint, tokens, type scale, Impeccable (no new findings) Assistant-verification: agent-browser E2E failed before this change: Stop Planner posted "paused its workflows" while the run card stayed running and no pause was recorded User-preference: Live workflow runs in Chopin show a pulsing dot rather than a spinner Co-authored-by: Alex Lavaee --- apps/server/src/harness/atomic/adapter.ts | 6 +-- apps/server/src/harness/atomic/full.test.ts | 49 +++++++++++++++++++- apps/server/src/harness/atomic/full.ts | 50 ++++++++++++++++++++- apps/server/src/harness/session.ts | 21 +++++---- apps/web/src/chat/run-card.css | 37 +++++++++++++++ apps/web/src/chat/run-card.tsx | 11 +++-- apps/web/src/main.tsx | 1 + 7 files changed, 157 insertions(+), 18 deletions(-) create mode 100644 apps/web/src/chat/run-card.css diff --git a/apps/server/src/harness/atomic/adapter.ts b/apps/server/src/harness/atomic/adapter.ts index b0d26233..3d79b4b0 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, workflowRuns } from "./full"; +import { controlWorkflows, fullPlanner, workflowRuns } from "./full"; import type { HarnessV1, @@ -506,9 +506,7 @@ export function createAtomicAdapter( }); let session = created.session; if (full) { - full.command = async text => { - await session.prompt(text); - }; + full.control = params => controlWorkflows(session, params); } let leak = full ? undefined diff --git a/apps/server/src/harness/atomic/full.test.ts b/apps/server/src/harness/atomic/full.test.ts index 6663cc36..92fbb1e2 100644 --- a/apps/server/src/harness/atomic/full.test.ts +++ b/apps/server/src/harness/atomic/full.test.ts @@ -10,7 +10,7 @@ import { ATOMIC_RESULT_TOOL_NAME, createAtomicAdapter, } from "./adapter"; -import { classifyRuns, registerFullPlanner } from "./full"; +import { classifyRuns, controlWorkflows, registerFullPlanner } from "./full"; import { startStubModelServer } from "../pi/model-stub"; import { hostInputRoom } from "../../testing/decisions"; import type { HostInput, QuestionParams } from "@bastani/atomic"; @@ -419,3 +419,50 @@ test("run cards fold lifecycle events into ordered stages, Decisions waits, and ended: 900, }); }); + +test("run control calls the session's workflow tool with a real tool context and surfaces failures", async () => { + let calls: { params: unknown; ctx: unknown; aborted: boolean }[] = []; + let reply: { content: { type: string; text: string }[]; isError?: boolean } = { + content: [{ type: "text", text: "Paused 1 run(s)." }], + }; + let session = { + getToolDefinition: (name: string) => + name === "workflow" + ? { + execute: (async ( + _id: string, + params: unknown, + signal: AbortSignal, + _update: unknown, + ctx: unknown, + ) => { + calls.push({ params, ctx, aborted: signal.aborted }); + return reply; + }) as never, + } + : undefined, + extensionRunner: { createToolContext: (id: string) => ({ toolCallId: id }) }, + }; + await controlWorkflows(session, { action: "pause", all: true }); + await controlWorkflows(session, { action: "resume", runId: "run-1" }); + expect(calls.map(call => call.params)).toEqual([{ action: "pause", all: true }, { + action: "resume", + runId: "run-1", + }]); + expect( + calls.every(call => + !call.aborted && typeof (call.ctx as { toolCallId: string }).toolCallId === "string" + ), + ).toBe(true); + reply = { content: [{ type: "text", text: "Run not found" }], isError: true }; + await expect(controlWorkflows(session, { action: "resume", runId: "gone" })).rejects.toThrow( + "Run not found", + ); + await expect( + controlWorkflows({ ...session, getToolDefinition: () => undefined }, { + action: "pause", + all: true, + }), + ) + .rejects.toThrow("no workflow tool"); +}); diff --git a/apps/server/src/harness/atomic/full.ts b/apps/server/src/harness/atomic/full.ts index 434dc643..24b9c872 100644 --- a/apps/server/src/harness/atomic/full.ts +++ b/apps/server/src/harness/atomic/full.ts @@ -21,11 +21,57 @@ export type FullPlanner = { runs?: PlannerRuns; /** Called with every change to `runs`. */ onRuns?: (runs: PlannerRuns) => void; - /** Runs a slash command in the live session without a model turn; set once the session exists. */ - command?: (text: string) => Promise; + /** Runs the session's own `workflow` tool without a model turn; set once the session exists. */ + control?: (params: WorkflowControl) => Promise; }; +export type WorkflowControl = { action: "pause"; all: true } | { action: "resume"; runId: string }; + let planners = new Map(); +/** The text of a tool result, for errors. */ +function resultText(content: readonly { type: string; text?: string }[]): string { + return content.map(part => part.type === "text" ? part.text ?? "" : "").join("\n").trim(); +} + +/** + * Runs one run-control action through the session's `workflow` tool with a real + * tool context, so the outcome comes back as a result rather than a UI notice. + */ +export async function controlWorkflows( + session: { + getToolDefinition( + name: string, + ): { execute: (...args: never[]) => Promise } | undefined; + extensionRunner: { createToolContext(id: string, signal: AbortSignal | undefined): unknown }; + }, + params: WorkflowControl, +): Promise { + let tool = session.getToolDefinition("workflow"); + if (!tool) throw new Error("The Planner session has no workflow tool."); + let id = `chopin-${params.action}-${crypto.randomUUID()}`; + let signal = AbortSignal.timeout(30_000); + let execute = tool.execute as ( + id: string, + params: WorkflowControl, + signal: AbortSignal, + onUpdate: undefined, + ctx: unknown, + ) => Promise< + { content?: { type: string; text?: string }[]; isError?: boolean; details?: unknown } + >; + let result = await execute( + id, + params, + signal, + undefined, + session.extensionRunner.createToolContext(id, signal), + ); + let text = resultText(result.content ?? []); + let failed = result.isError === true + || /"ok"\s*:\s*false|"status"\s*:\s*"failed"/.test(text); + if (failed) throw new Error(`workflow ${params.action} failed: ${text.slice(0, 500)}`); +} + /** 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"); diff --git a/apps/server/src/harness/session.ts b/apps/server/src/harness/session.ts index 0271868e..c28d14ea 100644 --- a/apps/server/src/harness/session.ts +++ b/apps/server/src/harness/session.ts @@ -2,7 +2,12 @@ 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 { type FullPlanner, type PlannerRuns, registerFullPlanner } from "./atomic/full"; +import { + type FullPlanner, + type PlannerRuns, + registerFullPlanner, + type WorkflowControl, +} from "./atomic/full"; import { createHumanInput } from "./atomic/human-input"; import { type PlannerWorkspace, plannerWorkspace } from "./atomic/workspace"; @@ -123,27 +128,27 @@ export async function openPlannerSession( let active = session; let release = unregister; let stopped: Promise | undefined; - let command = async (text: string) => { - if (!planner?.command) throw new Error("The Planner has no live Atomic session to control."); - await planner.command(text); + let control = async (params: WorkflowControl) => { + if (!planner?.control) throw new Error("The Planner has no live Atomic session to control."); + await planner.control(params); }; - let control = planner + let runControl = planner ? { runs: () => planner.runs, watchRuns: (listener: (runs: PlannerRuns) => void) => { listeners.add(listener); return () => listeners.delete(listener); }, - pauseRuns: () => command("/workflow pause --all --yes"), + pauseRuns: () => control({ action: "pause", all: true }), resumeRuns: async () => { - for (let runId of planner.runs?.paused ?? []) await command(`/workflow resume ${runId}`); + for (let runId of planner.runs?.paused ?? []) await control({ action: "resume", runId }); }, } : {}; return { ok: true, value: { - ...control, + ...runControl, stream: (prompt, abortSignal) => agent.stream({ session: active, diff --git a/apps/web/src/chat/run-card.css b/apps/web/src/chat/run-card.css new file mode 100644 index 00000000..baddc153 --- /dev/null +++ b/apps/web/src/chat/run-card.css @@ -0,0 +1,37 @@ +.run-pulse { + position: relative; + display: inline-block; + flex-shrink: 0; + width: 8px; + height: 8px; + margin: 3px; + border-radius: 9999px; + background-color: var(--color-brand); + + &::after { + position: absolute; + inset: 0; + border-radius: inherit; + background-color: inherit; + content: ""; + animation: run-pulse 1.6s cubic-bezier(0.16, 1, 0.3, 1) infinite; + } +} + +@keyframes run-pulse { + from { + opacity: 0.6; + transform: scale(1); + } + + to { + opacity: 0; + transform: scale(2.5); + } +} + +@media (prefers-reduced-motion: reduce) { + .run-pulse::after { + animation: none; + } +} diff --git a/apps/web/src/chat/run-card.tsx b/apps/web/src/chat/run-card.tsx index 7f9ab60d..64166047 100644 --- a/apps/web/src/chat/run-card.tsx +++ b/apps/web/src/chat/run-card.tsx @@ -1,7 +1,7 @@ /** A Planner-launched workflow run: is it alive, what is it doing, does it need anyone. */ import { useEffect, useState } from "react"; -import { CheckIcon, CloseIcon, LoaderIcon, WarningIcon } from "@chopin/icons"; +import { CheckIcon, CloseIcon, WarningIcon } from "@chopin/icons"; import type { Chat as Wire } from "@chopin/protocol"; @@ -46,9 +46,14 @@ function useNow(live: boolean): number { return now; } +/** A dot that pulses while work is live; it stays still under reduced motion. */ +function RunPulse() { + return
+ {(!ended || expanded) && ( +
    + {hidden > 0 && ( +
  1. +{hidden} earlier
  2. + )} + {stages.map(stage => ( +
  3. - {STAGE[stage.status]} - - {stage.started !== undefined && ( - - {elapsed((stage.ended ?? (live ? now : run.updated)) - stage.started)} + + {stage.name} + + {STAGE[stage.status]} - )} -
  4. - ))} -
+ {stage.started !== undefined && ( + + {elapsed((stage.ended ?? (live ? now : run.updated)) - stage.started)} + + )} + + ))} + + )} {run.waiting > 0 && ( + {control && ( )}
- {(!ended || expanded) && ( -
    - {hidden > 0 && ( -
  1. +{hidden} earlier
  2. - )} - {stages.map(stage => ( -
  3. - - {stage.name} - +
      + {earlier > 0 && ( +
    1. +{earlier} earlier
    2. + )} + {stages.map(stage => ( +
    3. - {STAGE[stage.status]} - - {stage.started !== undefined && ( - - {elapsed((stage.ended ?? (live ? now : run.updated)) - stage.started)} + + {stage.name} + + {STAGE[stage.status]} - )} -
    4. - ))} -
    + {stage.started !== undefined && ( + + {elapsed((stage.ended ?? (live ? now : run.updated)) - stage.started)} + + )} +
  4. + ))} +
+
)} {run.waiting > 0 && ( + )} ); } diff --git a/apps/web/src/design-audit/inventory.ts b/apps/web/src/design-audit/inventory.ts index 07ff88be..c5c14e36 100644 --- a/apps/web/src/design-audit/inventory.ts +++ b/apps/web/src/design-audit/inventory.ts @@ -164,7 +164,7 @@ export const AUDIT_INVENTORY: readonly AuditGroup[] = [ id: "workflow-runs", label: "Workflow runs", source: "apps/web/src/chat/run-card.tsx", - states: ["running", "waiting", "paused", "blocked", "finished"], + states: ["running", "concurrent", "ended"], }, { id: "identity", diff --git a/apps/web/src/design-audit/surfaces.tsx b/apps/web/src/design-audit/surfaces.tsx index 79819efb..28493577 100644 --- a/apps/web/src/design-audit/surfaces.tsx +++ b/apps/web/src/design-audit/surfaces.tsx @@ -10,7 +10,7 @@ import { } from "@chopin/icons"; import { DecisionCard, PlanStatus, SendAction, SidecarCard } from "@chopin/editor"; -import { RunCard } from "../chat/run-card"; +import { RunStack } from "../chat/run-card"; import { Transcript } from "../chat/transcript"; import { TerminalAlert } from "../terminal-alert"; import { AuditPlate, StateLabel } from "./frame"; @@ -221,13 +221,21 @@ const RUN_STAGES: Chat.RunStage[] = [ { id: "r:draft-2", name: "draft-2", status: "running", started: 1_620 }, ]; -function run(status: Chat.Run["status"], stages: Chat.RunStage[], waiting = 0): Chat.Run { - let ended = status === "finished" ? { ended: 2_100 } : {}; +function run( + status: Chat.Run["status"], + stages: Chat.RunStage[], + { name = "plan-review", started = 1_000, waiting = 0 }: { + name?: string; + started?: number; + waiting?: number; + } = {}, +): Chat.Run { + let ended = ["finished", "blocked", "failed", "stopped"].includes(status) ? { ended: 2_100 } : {}; return { - id: `audit-${status}`, - name: "plan-review", + id: `audit-${name}-${status}`, + name, status, - started: 1_000, + started, updated: 2_100, ...ended, stages, @@ -235,55 +243,47 @@ function run(status: Chat.Run["status"], stages: Chat.RunStage[], waiting = 0): }; } +const DONE_STAGES = RUN_STAGES.map(stage => ({ + ...stage, + status: "completed" as const, + ended: stage.ended ?? 2_100, +})); + function WorkflowRuns() { return (
Running - - Waiting - {}} runs={[run("running", RUN_STAGES)]} /> + Concurrent + {}} + onResume={() => {}} onShowDecisions={() => {}} - run={run("waiting", [ - ...RUN_STAGES.slice(0, 3), - { id: "r:draft-2", name: "draft-2", status: "awaiting_input", started: 1_620 }, - ], 2)} - /> - Paused - - Blocked - ({ - ...stage, - status: "completed", - ended: stage.ended ?? 2_100, - })), - ), - ended: 2_100, - }} + runs={[ + run("running", RUN_STAGES, { started: 1_400 }), + run("finished", DONE_STAGES, { name: "docs-pass", started: 600 }), + run("paused", [ + ...RUN_STAGES.slice(0, 3), + { id: "r:draft-2", name: "draft-2", status: "paused", started: 1_620 }, + ], { name: "api-audit", started: 1_200 }), + run("waiting", [ + ...RUN_STAGES.slice(0, 3), + { id: "r:draft-2", name: "draft-2", status: "awaiting_input", started: 1_620 }, + ], { name: "review-b", started: 1_100, waiting: 2 }), + run("blocked", DONE_STAGES, { name: "spec-check", started: 400 }), + ]} /> - Finished - ({ - ...stage, - status: "completed", - ended: stage.ended ?? 2_100, - })), - )} + Ended +
diff --git a/docs/hosted-agent.md b/docs/hosted-agent.md index c920455a..c144df19 100644 --- a/docs/hosted-agent.md +++ b/docs/hosted-agent.md @@ -278,11 +278,19 @@ through Atomic's session run control, and **Resume Planner** resumes the runs it paused. A run waiting on a question pauses too: Atomic withdraws the question from Decisions while the run is paused and presents it again on Resume, and an answer that arrives during the pause reaches the run only after Resume, so a -paused run never advances. Chat shows a card for each -run: its name, status (running, waiting on Decisions, paused, or ended), elapsed -time, every stage it has reached with its status and duration, and a link to the -Decisions it is waiting on. When a run ends, Chat records a line saying how it -ended and how long it took. The session is let go when its runs finish, when its owner +paused run never advances. Chat shows the session's runs as one stack, one row +per run: its name, status (running, waiting on Decisions, paused, or ended), and +elapsed time. Several runs can be live at once. Runs waiting on Decisions come +first, then running, paused, and ended runs, newest first within each, and more +than three rows fold behind "more" without ever hiding a waiting run. Only the +first waiting run, or else the newest live one, shows its stages; any row opens +on click to show its most recent stages with their status and duration, and a +waiting row links to its Decisions. Each live row can be paused, and each paused +row resumed, on its own; Chat records who did. The stack is stored with the +document, so a reload or restart keeps it; a run that was live when the server +stopped comes back stopped. Ended rows stay until the next workflow starts in the +document. When a run ends, Chat also records a line saying how it ended and how +long it took. The session is let go when its runs finish, when its owner binding ends, or when the document closes; run state is durable in Atomic's workflow store, so runs interrupted that way can be resumed from a later session. diff --git a/packages/protocol/chat.d.ts b/packages/protocol/chat.d.ts index 71e65c4c..7c8d5288 100644 --- a/packages/protocol/chat.d.ts +++ b/packages/protocol/chat.d.ts @@ -15,6 +15,8 @@ export declare namespace Chat { | Request | Request | Request + | Request + | Request | Request; export type Outgoing = History | Message | Delta | Tool | State | Queue | Sent; @@ -115,7 +117,10 @@ export declare namespace Chat { started: number; updated: number; ended?: number; + /** The most recent stages, in the order the run reached them. */ stages: RunStage[]; + /** Stages reached before the ones listed. */ + earlierStages?: number; /** Decisions questions the run is waiting on. */ waiting: number; }; @@ -184,6 +189,12 @@ export declare namespace Chat { /** Resume the workflow runs the Planner paused. Anyone may; the transcript records who did. */ export type Resume = KIND<"chat:resume">; + /** Pause one workflow run. Anyone may; the transcript records who did. */ + export type PauseRun = KIND<"chat:pause-run"> & { runId: string }; + + /** Resume one paused workflow run. Anyone may; the transcript records who did. */ + export type ResumeRun = KIND<"chat:resume-run"> & { runId: string }; + /** Withdraw a queued message. Only its author may. */ export type Unqueue = KIND<"chat:unqueue"> & { id: string }; } From 09eacad63c06d51c2aab99cf884bf085c81256ea Mon Sep 17 00:00:00 2001 From: Alex Lavaee Date: Thu, 1 Oct 2026 01:10:06 -0700 Subject: [PATCH 12/13] Show ctx.tool steps with a wrench and report a pause Atomic refused A workflow built from ctx.tool steps showed no steps in its run row, because only stage lifecycle events were folded. Tool steps now appear in order beside stages, marked with a new Nucleo-style wrench icon and "tool step" for screen readers; cached counts as done and cancelled as skipped. Per-run pause and resume ignored Atomic's outcome, so a pause Atomic could not apply, such as during a running ctx.tool call, was reported as done. A noop or cancelled outcome is now a failure with Atomic's message. Runs of the same workflow are told apart by a short run id in their rows and control labels. Assistant-model: Claude Opus 5.5 Assistant-workflow: inline Assistant-verification: bun test passed: 1767 pass, 2 PostgreSQL skips, 0 fail, including tool steps, the wrench and screen-reader label, and duplicate-name labels Assistant-verification: bun run types passed: all workspaces Assistant-verification: bun run ci passed: dprint, oxlint, tokens, type scale, Impeccable (no new findings) Assistant-verification: agent-browser E2E failed before this change: two concurrent ctx.tool-only runs showed no steps, and pausing one reported "paused" while it kept running User-preference: Tool stages show a tool icon Co-authored-by: Alex Lavaee --- apps/server/src/harness/atomic/full.test.ts | 31 +++++++++++ apps/server/src/harness/atomic/full.ts | 57 +++++++++++++++------ apps/server/src/harness/session.ts | 13 ++--- apps/server/src/plan/service.ts | 1 + apps/web/src/chat/run-card.test.tsx | 30 +++++++++++ apps/web/src/chat/run-card.tsx | 31 ++++++++--- apps/web/src/design-audit/icons.tsx | 2 + apps/web/src/design-audit/surfaces.tsx | 9 +++- packages/icons/src/index.ts | 1 + packages/icons/src/line.tsx | 8 +++ packages/protocol/chat.d.ts | 2 + 11 files changed, 157 insertions(+), 28 deletions(-) diff --git a/apps/server/src/harness/atomic/full.test.ts b/apps/server/src/harness/atomic/full.test.ts index 7f9b205d..8a755590 100644 --- a/apps/server/src/harness/atomic/full.test.ts +++ b/apps/server/src/harness/atomic/full.test.ts @@ -561,3 +561,34 @@ test("a new run clears ended cards, keeps live ones, and lists only the twelve m 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 1938ee64..241ea259 100644 --- a/apps/server/src/harness/atomic/full.ts +++ b/apps/server/src/harness/atomic/full.ts @@ -6,6 +6,7 @@ import type { WorkflowActivitySubscription, WorkflowLifecycleEvent, WorkflowRootActivity, + WorkflowToolNodeStatus, } from "@bastani/atomic"; import type { Chat as Wire } from "@chopin/protocol"; @@ -115,6 +116,37 @@ const ENDED: Record = { 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, @@ -150,21 +182,16 @@ export function foldLifecycle( card.ended = seconds; } else delete card.ended; } else if (target.kind === "stage") { - let key = `${target.runId}:${target.stageId}`; - let index = card.stageIndex.get(key); - if (index === undefined) { - index = card.stages.length; - card.stageIndex.set(key, index); - card.stages.push({ id: key, name: target.stageName, status: target.status }); - } - let stage = card.stages[index]!; - stage.status = target.status; - if (target.status === "running" && stage.started === undefined) stage.started = seconds; - if ( - target.status === "completed" || target.status === "failed" || target.status === "skipped" - ) { - stage.ended = seconds; - } + 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); diff --git a/apps/server/src/harness/session.ts b/apps/server/src/harness/session.ts index 0efbc4a0..5804d871 100644 --- a/apps/server/src/harness/session.ts +++ b/apps/server/src/harness/session.ts @@ -37,6 +37,11 @@ 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; @@ -153,12 +158,8 @@ export async function openPlannerSession( }, pauseRuns: () => pauseOwnedRuns(workflows()), resumeRuns: () => resumeOwnedRuns(workflows()), - pauseRun: async (runId: string) => { - await workflows().pause(runId); - }, - resumeRun: async (runId: string) => { - await workflows().resume(runId); - }, + pauseRun: async (runId: string) => applied(await workflows().pause(runId)), + resumeRun: async (runId: string) => applied(await workflows().resume(runId)), } : {}; return { diff --git a/apps/server/src/plan/service.ts b/apps/server/src/plan/service.ts index caa3ff4c..c724d12c 100644 --- a/apps/server/src/plan/service.ts +++ b/apps/server/src/plan/service.ts @@ -537,6 +537,7 @@ function restoreWorkflowRuns(value: JsonValue | undefined): NonNullable !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) ) diff --git a/apps/web/src/chat/run-card.test.tsx b/apps/web/src/chat/run-card.test.tsx index d385fb1b..003d0bde 100644 --- a/apps/web/src/chat/run-card.test.tsx +++ b/apps/web/src/chat/run-card.test.tsx @@ -87,3 +87,33 @@ test("a concurrent stack opens only the waiting run, offers per-run controls, an 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 index a94a5ecd..136e6def 100644 --- a/apps/web/src/chat/run-card.tsx +++ b/apps/web/src/chat/run-card.tsx @@ -1,7 +1,7 @@ /** 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 } from "@chopin/icons"; +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"; @@ -123,8 +123,10 @@ type RunControls = { onResume?: (runId: string) => void; }; -function RunRow({ run, open, onToggle, onShowDecisions, onPause, onResume }: RunControls & { +function RunRow({ run, tag, open, onToggle, onShowDecisions, onPause, onResume }: RunControls & { run: Wire.Run; + /** Tells apart runs of the same workflow. */ + tag?: string; open: boolean; onToggle: () => void; }) { @@ -135,11 +137,12 @@ function RunRow({ run, open, onToggle, onShowDecisions, onPause, onResume }: Run let stages = run.stages.slice(hidden); let earlier = hidden + (run.earlierStages ?? 0); let stagesId = `run-stages-${run.id}`; + let title = tag ? `${run.name} ${tag}` : run.name; let control = live && onPause - ? { label: `Pause ${run.name}`, icon: plannerStop, act: () => onPause(run.id), kind: "pause" } + ? { label: `Pause ${title}`, icon: plannerStop, act: () => onPause(run.id), kind: "pause" } : run.status === "paused" && onResume ? { - label: `Resume ${run.name}`, + label: `Resume ${title}`, icon: plannerResume, act: () => onResume(run.id), kind: "resume", @@ -152,7 +155,7 @@ function RunRow({ run, open, onToggle, onShowDecisions, onPause, onResume }: Run