diff --git a/extensions/subagents/src/backend.ts b/extensions/subagents/src/backend.ts index b0ce0f2a..e232cd95 100644 --- a/extensions/subagents/src/backend.ts +++ b/extensions/subagents/src/backend.ts @@ -19,6 +19,13 @@ import type { SubagentMeta, } from "./domain.ts"; +export interface SubagentCleanupReceipt { + /** True when the session/worktree state could not be proven quiescent. */ + readonly uncertain: boolean; + /** Bounded operator-facing explanation of the cleanup result. */ + readonly message: string; +} + /** * A live subagent session. The manager is the single consumer of `events`; * it folds them into the `SubagentSnapshot` everything else reads. @@ -44,6 +51,8 @@ export interface SubagentSession { readonly interrupt: Effect.Effect; /** Ephemeral Pi-native tool projection for the operator-facing child page. */ readonly toolRenderer?: AgentToolRenderer; + /** Cleanup receipt populated when scope teardown cannot prove safety. */ + readonly cleanupReceipt?: () => SubagentCleanupReceipt | undefined; } export interface SubagentBackend { diff --git a/extensions/subagents/src/backends/pi.ts b/extensions/subagents/src/backends/pi.ts index 641e6b89..b50db916 100644 --- a/extensions/subagents/src/backends/pi.ts +++ b/extensions/subagents/src/backends/pi.ts @@ -24,7 +24,11 @@ import { import type { Cause, Scope } from "effect"; import { Effect, Queue, Stream } from "effect"; import { resolveAgentModel } from "../agent-types.ts"; -import type { SubagentBackend, SubagentSession } from "../backend.ts"; +import type { + SubagentBackend, + SubagentCleanupReceipt, + SubagentSession, +} from "../backend.ts"; import type { SpawnTask, SubagentEvent, @@ -40,10 +44,14 @@ import { createChildResources, shutdownAndDisposeChildSession, } from "../../../shared/child-session.ts"; -import { reclaimWorktree } from "../../../shared/worktree.ts"; +import { + formatWorktreeCleanupWarning, + reclaimWorktree, +} from "../../../shared/worktree.ts"; import { AgentToolRenderLedger } from "../../../shared/agent-tool-renderer.ts"; const DIRECT_WORKTREE_CLEANUP_TIMEOUT_MS = 4_000; +const PARTIAL_TEXT_MAX_LENGTH = 128 * 1_024; // --- Model + effort resolution ----------------------------------------------- @@ -51,6 +59,19 @@ type ThinkingLevel = NonNullable< NonNullable[0]>["thinkingLevel"] >; +export type PiAgentSessionFactory = ( + options: Parameters[0], +) => Promise<{ session: AgentSession }>; + +export interface PiBackendOptions { + /** Test seam for production-adapter lifecycle coverage. */ + readonly sessionFactory?: PiAgentSessionFactory; + /** Test-only override for bounded child shutdown. */ + readonly shutdownTimeoutMs?: number; + /** Test seam for observing whether unsafe worktree reclamation is attempted. */ + readonly worktreeCleanup?: typeof reclaimWorktree; +} + // --- Event translation ---------------------------------------------------------- function messageRole(msg: unknown): Message["role"] | undefined { @@ -71,20 +92,12 @@ function lastAssistantMessage( return undefined; } -/** Final assistant text output (last assistant message with text), v1 semantics. */ -function finalOutput(session: AgentSession): string { - const messages = session.messages; - for (let i = messages.length - 1; i >= 0; i--) { - const msg = messages[i]; - if (messageRole(msg) !== "assistant") continue; - const text = (msg as AssistantMessage).content - .filter((part) => part.type === "text") - .map((part) => part.text) - .join("\n") - .trim(); - if (text) return text; - } - return ""; +function assistantText(message: AssistantMessage) { + return message.content + .filter((part) => part.type === "text") + .map((part) => part.text) + .join("\n") + .trim(); } function safeJson(value: unknown): string | undefined { @@ -166,6 +179,7 @@ function boundedError(error: unknown) { const makePiSession = ( task: SpawnTask, + options: PiBackendOptions, ): Effect.Effect => Effect.gen(function* () { const registry = task.parent.modelRegistry; @@ -193,7 +207,9 @@ const makePiSession = ( ? { appendSystemPrompt: [...task.appendSystemPrompt] } : {}), }); - const { session } = await createAgentSession({ + const { session } = await ( + options.sessionFactory ?? createAgentSession + )({ cwd: task.cwd, sessionManager: SessionManager.create(task.cwd), settingsManager, @@ -208,7 +224,9 @@ const makePiSession = ( try { await bindChildSessionExtensions(session, task.tools); } catch (error) { - await shutdownAndDisposeChildSession(session); + await shutdownAndDisposeChildSession(session, { + timeoutMs: options.shutdownTimeoutMs, + }); throw error; } return session; @@ -216,14 +234,35 @@ const makePiSession = ( catch: (error) => new SpawnError({ message: boundedError(error) }), }); + interface ActivePrompt { + cancelled: boolean; + lifecycleStarted: boolean; + lifecycleSettled: boolean; + promptSettled: boolean; + promptError?: string; + lastAssistant?: AssistantMessage; + finalText: string; + liveText: string; + partialTextAtCancel?: string; + promise: Promise; + } + + interface PendingRestart { + cancelled: boolean; + terminalEmitted: boolean; + } + const state = { closed: false, - /** prompt() rejection for the active run; folded into RunSettled. */ - runError: undefined as string | undefined, /** One terminal event per run: lifecycle, prompt-rejection, and abort * fallbacks can all race to settle; the first wins. */ settled: false, + /** Pi events have no run identity, so the exact prompt Promise is the + * only safe lease for deciding when another run may start. */ + activePrompt: undefined as ActivePrompt | undefined, + pendingRestart: undefined as PendingRestart | undefined, }; + let cleanupReceipt: SubagentCleanupReceipt | undefined; const events = yield* Queue.make(); const emit = (event: SubagentEvent) => { @@ -271,20 +310,13 @@ const makePiSession = ( }); }; - const settle = () => { - if (state.settled) return; + const settle = (prompt = state.activePrompt) => { + if (!prompt || state.activePrompt !== prompt || state.settled) return; state.settled = true; - const last = lastAssistantMessage(session); - const partialText = finalOutput(session) || undefined; - if (last?.stopReason === "aborted") { - emit({ - _tag: "RunSettled", - outcome: { _tag: "Interrupted", partialText }, - }); - return; - } + const last = prompt.lastAssistant; + const partialText = prompt.liveText || prompt.finalText || undefined; const errorText = - state.runError ?? + prompt.promptError ?? (last?.stopReason === "error" ? (last.errorMessage ?? "Run failed") : undefined); @@ -299,24 +331,86 @@ const makePiSession = ( }); return; } + if (last?.stopReason === "aborted") { + emit({ + _tag: "RunSettled", + outcome: { _tag: "Interrupted", partialText }, + }); + return; + } emit({ _tag: "RunSettled", - outcome: { _tag: "Completed", finalText: finalOutput(session) }, + outcome: { + _tag: "Completed", + finalText: prompt.finalText, + }, }); }; + /** + * A Pi lifecycle event is evidence, not the lease itself. The exact + * prompt() Promise must also settle before a terminal outcome is trusted; + * otherwise agent_settled can publish Completed and hide a later reject. + */ + const maybeSettle = (prompt: ActivePrompt) => { + if ( + state.closed || + state.activePrompt !== prompt || + prompt.cancelled || + state.settled || + !prompt.promptSettled + ) { + return; + } + if (prompt.promptError !== undefined) { + settle(prompt); + } else if (!prompt.lifecycleStarted) { + prompt.promptError = + "Pi prompt completed without lifecycle events; no agent response was observed."; + settle(prompt); + } else if (!prompt.lifecycleSettled) { + prompt.promptError = + "Pi prompt completed without an agent_settled lifecycle event."; + settle(prompt); + } else { + settle(prompt); + } + if (state.activePrompt === prompt && prompt.promptSettled) { + state.activePrompt = undefined; + } + }; + const handleEvent = (event: AgentSessionEvent) => { if (state.closed) return; + const activePrompt = state.activePrompt; + if (!activePrompt) return; + if (activePrompt.cancelled) { + if (event.type === "agent_start") { + // abort() can return during Pi's preflight while prompt() remains + // pending. If that preflight later enters the agent lifecycle, + // abort again before it can reach a provider or tool. Never reopen + // the event gate: native events carry no generation identity. + try { + void session.abort().catch(() => undefined); + } catch { + // The original interrupt still owns bounded cleanup. + } + } + return; + } switch (event.type) { case "agent_start": // Extensions may register tools between runs; guard new ones too. toolTimeout.apply(session); - state.settled = false; + activePrompt.lifecycleStarted = true; emit({ _tag: "RunStarted" }); break; case "message_update": { const streamEvent = event.assistantMessageEvent; if (streamEvent.type === "text_delta") { + activePrompt.liveText = ( + activePrompt.liveText + streamEvent.delta + ).slice(-PARTIAL_TEXT_MAX_LENGTH); emit({ _tag: "AssistantDelta", kind: "text", @@ -337,9 +431,14 @@ const makePiSession = ( const text = userText(event.message as Message); if (text.trim()) emit({ _tag: "UserMessage", text }); } else if (role === "assistant") { + const assistant = event.message as AssistantMessage; + activePrompt.lastAssistant = assistant; + const text = assistantText(assistant); + if (text) activePrompt.finalText = text; + activePrompt.liveText = ""; emit({ _tag: "AssistantMessage", - parts: assistantParts(event.message as AssistantMessage), + parts: assistantParts(assistant), }); emitUsage(); emit({ _tag: "MetaChanged", meta: currentMeta() }); @@ -405,7 +504,8 @@ const makePiSession = ( }); break; case "agent_settled": - settle(); + activePrompt.lifecycleSettled = true; + maybeSettle(activePrompt); break; } }; @@ -420,17 +520,61 @@ const makePiSession = ( } catch { // Continue with abort/dispose. } - await shutdownAndDisposeChildSession(session, { + const shutdown = await shutdownAndDisposeChildSession(session, { abort: true, - timeoutMs: CHILD_SHUTDOWN_TIMEOUT_MS, + timeoutMs: options.shutdownTimeoutMs ?? CHILD_SHUTDOWN_TIMEOUT_MS, }); - // A direct child can run multiple turns, so its worktree lives until - // this session scope closes. Cleanup is fail-closed and shares a hard - // deadline; unknown/dirty/ignored/detached work is preserved. - if (task.worktree) { - await reclaimWorktree(task.worktree.repoCwd, task.worktree, { - timeoutMs: DIRECT_WORKTREE_CLEANUP_TIMEOUT_MS, - }).catch(() => {}); + const promptQuiescent = + state.activePrompt === undefined || state.activePrompt.promptSettled; + const cleanupMessages = [...shutdown.errors]; + let cleanupUncertain = + shutdown.errors.length > 0 || shutdown.timedOut || !promptQuiescent; + if (shutdown.timedOut) + cleanupMessages.push("cleanup deadline exceeded"); + if (!promptQuiescent) { + cleanupMessages.push( + "active prompt did not quiesce before session shutdown", + ); + } + + // A preflight hook can outlive abort() and continue touching the child + // cwd. Never reclaim a worktree until both the prompt and session + // cleanup have reached a known-safe terminal state. + if ( + task.worktree && + promptQuiescent && + shutdown.ok && + !shutdown.timedOut + ) { + try { + const cleanup = await (options.worktreeCleanup ?? reclaimWorktree)( + task.worktree.repoCwd, + task.worktree, + { + timeoutMs: DIRECT_WORKTREE_CLEANUP_TIMEOUT_MS, + }, + ); + const warning = formatWorktreeCleanupWarning( + cleanup, + task.worktree.path, + ); + if (warning) cleanupMessages.push(`worktree cleanup: ${warning}`); + } catch (error) { + cleanupUncertain = true; + cleanupMessages.push( + `worktree cleanup failed: ${boundedError(error)}; checkout preserved at ${task.worktree.path}`, + ); + } + } else if (task.worktree) { + cleanupMessages.push( + `worktree cleanup skipped; checkout preserved at ${task.worktree.path}`, + ); + } + if (cleanupMessages.length > 0) { + cleanupReceipt = { + uncertain: cleanupUncertain, + message: `${cleanupUncertain ? "Subagent cleanup is uncertain" : "Subagent cleanup"}: ${cleanupMessages.join("; ")}`, + }; } Queue.endUnsafe(events); }), @@ -438,15 +582,63 @@ const makePiSession = ( /** Start a fresh run (v1 manager.run): fire-and-forget, errors -> events. */ const startRun = (text: string) => { - state.runError = undefined; + if (state.activePrompt) { + throw new Error( + "Cannot start a Pi prompt while another prompt is active.", + ); + } + const activePrompt: ActivePrompt = { + cancelled: false, + lifecycleStarted: false, + lifecycleSettled: false, + promptSettled: false, + finalText: "", + liveText: "", + promise: Promise.resolve(), + }; + state.activePrompt = activePrompt; state.settled = false; emit({ _tag: "RunStarted" }); - void session.prompt(text).catch((error) => { - state.runError = boundedError(error); - // Preflight failures may never start the agent lifecycle, so no - // agent_settled will arrive for them. - if (!session.isStreaming) settle(); - }); + let prompt: Promise; + try { + prompt = session.prompt(text, { + preflightResult: (accepted) => { + if (accepted && (activePrompt.cancelled || state.closed)) { + // This callback runs immediately before Pi enters + // _runAgentPrompt. Throwing here is the only synchronous seam + // that prevents a cancelled preflight from reviving after the + // adapter subscription and worktree ownership are gone. + throw new Error( + "Subagent prompt was cancelled during preflight.", + ); + } + }, + }); + } catch (error) { + prompt = Promise.reject(error); + } + activePrompt.promise = prompt + .then( + () => undefined, + (error) => { + if ( + !state.closed && + state.activePrompt === activePrompt && + !activePrompt.cancelled + ) { + activePrompt.promptError = boundedError(error); + } + }, + ) + .finally(() => { + activePrompt.promptSettled = true; + maybeSettle(activePrompt); + if (state.activePrompt === activePrompt) { + if (activePrompt.cancelled || state.settled) { + state.activePrompt = undefined; + } + } + }); }; // Session naming is best-effort. @@ -465,6 +657,7 @@ const makePiSession = ( return { toolRenderer, + cleanupReceipt: () => cleanupReceipt, meta: Effect.sync(currentMeta), events: Stream.fromQueue(events), send: (text) => @@ -472,7 +665,21 @@ const makePiSession = ( if (state.closed) { return new SendError({ message: "Subagent session is closed." }); } - if (session.isStreaming) { + const pendingRestart = state.pendingRestart; + if (pendingRestart) { + return new SendError({ + message: pendingRestart.cancelled + ? "Subagent cancellation is still in progress." + : "A subagent restart is already pending.", + }); + } + const activePrompt = state.activePrompt; + if (activePrompt?.cancelled) { + return new SendError({ + message: "Subagent cancellation is still in progress.", + }); + } + if (activePrompt && session.isStreaming) { // Steer the active run via the SDK's queue; queue_update events // render it, message_end(user) lands it in the transcript. A // rejected steer is a real send failure, not a diagnostic. @@ -481,34 +688,115 @@ const makePiSession = ( catch: (error) => new SendError({ message: boundedError(error) }), }).pipe(Effect.asVoid); } + if ( + activePrompt && + !activePrompt.lifecycleSettled && + !activePrompt.promptSettled + ) { + return new SendError({ + message: + "Subagent prompt is still starting; wait for it to run or settle before sending.", + }); + } + if (activePrompt) { + // Pi emits agent_settled before prompt() resolves. Preserve the + // same-session restart contract while closing that microtask gap + // without ever overlapping two native prompts. + const restart: PendingRestart = { + cancelled: false, + terminalEmitted: false, + }; + state.pendingRestart = restart; + return Effect.tryPromise({ + try: async (signal) => { + const cancelPendingRestart = () => { + restart.cancelled = true; + if (state.pendingRestart === restart) { + state.pendingRestart = undefined; + } + }; + signal.addEventListener("abort", cancelPendingRestart, { + once: true, + }); + if (signal.aborted) cancelPendingRestart(); + try { + await activePrompt.promise; + if (state.closed) { + throw new Error("Subagent session is closed."); + } + if (restart.cancelled) { + // interrupt emits the terminal event for this accepted + // restart. Resolve the send without starting a prompt so + // the manager keeps its restarting marker until that + // terminal event is folded. + return; + } + state.pendingRestart = undefined; + startRun(text); + } finally { + signal.removeEventListener("abort", cancelPendingRestart); + if (state.pendingRestart === restart) { + state.pendingRestart = undefined; + } + } + }, + catch: (error) => new SendError({ message: boundedError(error) }), + }).pipe(Effect.asVoid); + } return Effect.sync(() => startRun(text)); }), interrupt: Effect.promise(async () => { if (state.closed) return; + const pendingRestart = state.pendingRestart; + if (pendingRestart) pendingRestart.cancelled = true; + const activePrompt = state.activePrompt; + if (activePrompt) { + activePrompt.cancelled = true; + activePrompt.partialTextAtCancel = pendingRestart + ? undefined + : activePrompt.liveText || activePrompt.finalText || undefined; + } try { session.clearQueue(); } catch { // Abort regardless. } await session.abort().catch(() => undefined); - // Only resolve once streaming has actually stopped: reporting the - // interrupt as complete while the run keeps working would let the - // manager settle a run that is still mutating the workspace. The - // manager bounds this effect at 5s and force-disposes on timeout. - while (!state.closed && session.isStreaming) { - await new Promise((resolve) => setTimeout(resolve, 50)); + // isStreaming is false during Pi preflight, so it cannot prove the + // invocation is quiescent. Wait for the exact prompt Promise instead; + // the manager bounds this effect at 5s and force-disposes on timeout. + await activePrompt?.promise; + if ( + !state.closed && + pendingRestart && + !pendingRestart.terminalEmitted + ) { + pendingRestart.terminalEmitted = true; + if (state.pendingRestart === pendingRestart) { + state.pendingRestart = undefined; + } + state.settled = true; + emit({ _tag: "RunSettled", outcome: { _tag: "Interrupted" } }); + return; } - // No streaming run means no agent_settled will arrive; emit the - // terminal event (once) so the run cannot look running forever. if (!state.closed && !state.settled) { state.settled = true; - emit({ _tag: "RunSettled", outcome: { _tag: "Interrupted" } }); + emit({ + _tag: "RunSettled", + outcome: { + _tag: "Interrupted", + partialText: activePrompt?.partialTextAtCancel, + }, + }); } }), } satisfies SubagentSession; }); -export const piBackend: SubagentBackend = { - name: "pi", - spawn: makePiSession, -}; +export const makePiBackend = (options: PiBackendOptions = {}) => + ({ + name: "pi", + spawn: (task) => makePiSession(task, options), + }) satisfies SubagentBackend; + +export const piBackend = makePiBackend(); diff --git a/extensions/subagents/src/manager.ts b/extensions/subagents/src/manager.ts index 64c6fccd..033bcd30 100644 --- a/extensions/subagents/src/manager.ts +++ b/extensions/subagents/src/manager.ts @@ -309,7 +309,20 @@ const makeManager = (config: SubagentManagerConfig = {}) => }; const closeEntryScope = (entry: Entry) => - Scope.close(entry.scope, Exit.void).pipe(Effect.ignore); + Scope.close(entry.scope, Exit.void).pipe( + Effect.tap(() => + Effect.sync(() => { + const receipt = entry.session.cleanupReceipt?.(); + if (!receipt?.uncertain) return; + const current = entry.snapshot.errorText; + entry.snapshot.errorText = current + ? `${current}; ${receipt.message}` + : receipt.message; + notify(entry.snapshot.id); + }), + ), + Effect.ignore, + ); const pruneSettled = () => { if (entries.size <= MAX_TRACKED) return; diff --git a/tests/extensions/subagents/pi-backend-lifecycle.test.ts b/tests/extensions/subagents/pi-backend-lifecycle.test.ts new file mode 100644 index 00000000..8a39f828 --- /dev/null +++ b/tests/extensions/subagents/pi-backend-lifecycle.test.ts @@ -0,0 +1,1234 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type { + AgentSession, + CreateAgentSessionOptions, + ModelRegistry, +} from "@earendil-works/pi-coding-agent"; +import { Effect, Layer, ManagedRuntime, Stream } from "effect"; +import type { + SubagentCleanupReceipt, + SubagentBackend, + SubagentSession, +} from "../../../extensions/subagents/src/backend.ts"; +import { BackendRegistry } from "../../../extensions/subagents/src/backend.ts"; +import { + makePiBackend, + type PiAgentSessionFactory, +} from "../../../extensions/subagents/src/backends/pi.ts"; +import type { + BackendName, + SpawnTask, + SubagentEvent, +} from "../../../extensions/subagents/src/domain.ts"; +import { + makeSubagentManagerLayer, + SubagentManager, +} from "../../../extensions/subagents/src/manager.ts"; +import { + createPiAgentSessionHarness, + type PiAgentSessionHarness, + type PiAgentSessionHarnessOptions, +} from "../../support/pi-agent-session-harness.ts"; + +const FIXTURE_MODEL = { + provider: "fixture", + id: "fixture-model", + name: "Fixture Model", + api: "openai-completions", + baseUrl: "http://127.0.0.1:1", + reasoning: true, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 8_192, + maxTokens: 1_024, +}; + +const MODEL_REGISTRY = { + find: (provider: string, id: string) => + provider === FIXTURE_MODEL.provider && id === FIXTURE_MODEL.id + ? FIXTURE_MODEL + : undefined, + getAll: () => [FIXTURE_MODEL], +} as unknown as ModelRegistry; + +type SessionCreationOptions = NonNullable; + +function deferred() { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((next, fail) => { + resolve = next; + reject = fail; + }); + return { promise, resolve, reject }; +} + +async function waitFor( + predicate: () => boolean, + label: string, + timeoutMs = 2_000, +) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) return; + await new Promise((resolve) => setTimeout(resolve, 5)); + } + assert.ok(predicate(), `timed out waiting for ${label}`); +} + +function task(prompt: string, overrides: Partial = {}): SpawnTask { + return { + prompt, + title: "lifecycle fixture", + cwd: process.cwd(), + model: `${FIXTURE_MODEL.provider}/${FIXTURE_MODEL.id}`, + reasoningEffort: "high", + tools: ["read", "subagent_spawn"], + parent: { + parentCwd: process.cwd(), + projectTrusted: false, + inheritedModel: { + provider: FIXTURE_MODEL.provider, + id: FIXTURE_MODEL.id, + }, + inheritedThinkingLevel: "medium", + modelRegistry: MODEL_REGISTRY, + }, + ...overrides, + }; +} + +function harnessFactory( + configure: ( + options: SessionCreationOptions, + index: number, + ) => PiAgentSessionHarnessOptions = () => ({}), +) { + const creations: SessionCreationOptions[] = []; + const harnesses: PiAgentSessionHarness[] = []; + const factory: PiAgentSessionFactory = async (options) => { + assert.ok(options); + const index = creations.length; + creations.push(options); + const harness = createPiAgentSessionHarness({ + model: options.model, + activeTools: options.tools, + ...configure(options, index), + }); + harnesses.push(harness); + return { session: harness.session }; + }; + return { factory, creations, harnesses }; +} + +function createManagerRuntime( + backend: SubagentBackend, + firstResponseTimeoutMs = 500, +) { + const registry = Layer.succeed( + BackendRegistry, + new Map([["pi", backend]]), + ); + return ManagedRuntime.make( + makeSubagentManagerLayer({ firstResponseTimeoutMs }).pipe( + Layer.provide(registry), + ), + ); +} + +async function spawnDirect( + backend: SubagentBackend, + spawnTask: SpawnTask, + drive: (session: SubagentSession) => void, +) { + const events: SubagentEvent[] = []; + const firstSettlement = deferred(); + await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const session = yield* backend.spawn(spawnTask); + yield* Effect.forkScoped( + Stream.runForEach(session.events, (event) => + Effect.sync(() => { + events.push(event); + if (event._tag === "RunSettled") firstSettlement.resolve(); + }), + ), + ); + drive(session); + yield* Effect.promise(() => firstSettlement.promise); + }), + ), + ); + return events; +} + +test("the production Pi adapter bridges startup, first response, tools, usage, and one settlement", async () => { + const fixtures = harnessFactory(); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + + const events = await spawnDirect( + backend, + task("inspect the repository"), + () => { + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitUser("inspect the repository"); + harness.emitTextDelta("Looking"); + harness.emit({ + type: "tool_execution_start", + toolCallId: "tool-1", + toolName: "read", + args: { path: "README.md" }, + }); + harness.emit({ + type: "tool_execution_update", + toolCallId: "tool-1", + toolName: "read", + args: { path: "README.md" }, + partialResult: { content: [{ type: "text", text: "partial\nrest" }] }, + }); + harness.emit({ + type: "tool_execution_end", + toolCallId: "tool-1", + toolName: "read", + result: { content: [{ type: "text", text: "complete\nrest" }] }, + isError: false, + }); + harness.setContextUsage({ tokens: 321, contextWindow: 8_192 }); + harness.emitAssistant("Repository inspected"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + harness.emit({ type: "agent_settled" }); + }, + ); + + const creation = fixtures.creations[0]; + const harness = fixtures.harnesses[0]; + assert.ok(creation); + assert.ok(harness); + assert.equal(creation.cwd, process.cwd()); + assert.equal(creation.model, FIXTURE_MODEL); + assert.equal(creation.thinkingLevel, "high"); + assert.deepEqual(creation.tools, ["read"]); + assert.ok(creation.excludeTools?.includes("subagent_spawn")); + assert.ok(creation.resourceLoader); + assert.ok(creation.settingsManager); + assert.ok(creation.sessionManager); + assert.deepEqual(harness.activeTools(), ["read"]); + assert.deepEqual(harness.calls.bindings, [{ mode: "print" }]); + assert.deepEqual(harness.calls.prompts, ["inspect the repository"]); + assert.deepEqual(harness.calls.sessionNames, ["subagent: lifecycle fixture"]); + + assert.ok(events.some((event) => event._tag === "AssistantDelta")); + assert.ok(events.some((event) => event._tag === "ToolStart")); + assert.ok(events.some((event) => event._tag === "ToolUpdate")); + assert.ok(events.some((event) => event._tag === "ToolEnd")); + assert.ok( + events.some( + (event) => + event._tag === "UsageChanged" && + event.tokens === 321 && + event.contextWindow === 8_192, + ), + ); + assert.ok( + events.some( + (event) => + event._tag === "MetaChanged" && + event.meta.modelLabel === "fixture/fixture-model", + ), + ); + assert.deepEqual( + events + .filter((event) => event._tag === "RunSettled") + .map((event) => event.outcome), + [{ _tag: "Completed", finalText: "Repository inspected" }], + ); + assert.equal(harness.calls.clearQueues, 1); + assert.equal(harness.calls.aborts, 1); + assert.equal(harness.calls.disposals, 1); +}); + +test("prompt rejection wins over an earlier agent_settled event", async () => { + const prompt = deferred(); + const fixtures = harnessFactory(() => ({ + prompt: () => prompt.promise, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend, 10_000); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("settled event before prompt rejection")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("partial provider output"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(manager.view.get(spawned.id)?.status, "running"); + + prompt.reject(new Error("provider rejected after settlement event")); + await waitFor( + () => manager.view.get(spawned.id)?.status === "error", + "prompt rejection to settle", + ); + assert.equal( + manager.view.get(spawned.id)?.errorText, + "provider rejected after settlement event", + ); + } finally { + prompt.resolve(); + await runtime.dispose(); + } +}); + +test("a prompt that resolves without lifecycle events fails explicitly", async () => { + const fixtures = harnessFactory(() => ({ prompt: async () => {} })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend, 10_000); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("extension handled prompt")), + ); + await waitFor( + () => manager.view.get(spawned.id)?.status === "error", + "prompt without lifecycle to fail", + ); + assert.match( + manager.view.get(spawned.id)?.errorText ?? "", + /prompt completed without lifecycle events/i, + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("late lifecycle output"); + harness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(manager.view.get(spawned.id)?.finalText, ""); + } finally { + await runtime.dispose(); + } +}); + +test("a non-quiescent prompt preserves its worktree and exposes cleanup uncertainty", async () => { + const stuckPreflight = deferred(); + let cleanupCalls = 0; + const fixtures = harnessFactory(() => ({ + preflight: () => stuckPreflight.promise, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 25, + worktreeCleanup: async () => { + cleanupCalls++; + throw new Error("worktree cleanup should not run"); + }, + }); + let cleanupReceipt: (() => SubagentCleanupReceipt | undefined) | undefined; + await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const session = yield* backend.spawn( + task("stuck preflight worktree", { + worktree: { + path: "C:/fixture-worktree", + branch: "pi/fixture-worktree", + repoCwd: process.cwd(), + }, + }), + ); + cleanupReceipt = session.cleanupReceipt; + }), + ), + ); + assert.equal(cleanupCalls, 0); + assert.match(cleanupReceipt?.()?.message ?? "", /prompt did not quiesce/i); + stuckPreflight.resolve(); +}); + +test("streaming send steers while idle send starts a distinct run", async () => { + const fixtures = harnessFactory(); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + const settlements: string[] = []; + manager.view.setOnSettled((snapshot) => + settlements.push(`${snapshot.id}:${snapshot.finalText}`), + ); + + const first = await runtime.runPromise( + manager.spawn("pi", task("first turn")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitTextDelta("first response"); + + await runtime.runPromise(manager.send(first.id, "steer the active run")); + assert.deepEqual(harness.calls.steers, ["steer the active run"]); + assert.deepEqual(harness.calls.prompts, ["first turn"]); + + harness.emitAssistant("first result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await runtime.runPromise(manager.waitFor([first.id])); + assert.equal(manager.view.get(first.id)?.finalText, "first result"); + + await runtime.runPromise(manager.send(first.id, "second turn")); + assert.deepEqual(harness.calls.prompts, ["first turn", "second turn"]); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitTextDelta("second response"); + harness.emitAssistant("second result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await runtime.runPromise(manager.waitFor([first.id])); + + const afterRestart = manager.view.get(first.id); + assert.equal(afterRestart?.status, "done"); + assert.equal(afterRestart?.finalText, "second result"); + assert.equal(settlements.length, 2); + assert.deepEqual( + settlements.map((entry) => entry.split(":", 1)[0]), + [first.id, first.id], + ); + + harness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(settlements.length, 2); + } finally { + await runtime.dispose(); + } +}); + +test("cancelled active run restarts only after the old prompt is quiescent", async () => { + const fixtures = harnessFactory(); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + let settlements = 0; + manager.view.setOnSettled(() => settlements++); + + const spawned = await runtime.runPromise( + manager.spawn("pi", task("cancelled turn")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + + const cancellation = runtime.runPromise(manager.cancel([spawned.id])); + await waitFor(() => harness.calls.aborts === 1, "active run abort"); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + const report = await cancellation; + assert.equal(report[0]?.cancelled, true); + assert.equal(settlements, 1); + + await runtime.runPromise(manager.send(spawned.id, "restarted turn")); + assert.deepEqual(harness.calls.prompts, [ + "cancelled turn", + "restarted turn", + ]); + await waitFor( + () => manager.view.get(spawned.id)?.status === "running", + "restarted run to reach the manager", + ); + assert.equal(manager.view.get(spawned.id)?.status, "running"); + assert.equal(manager.view.get(spawned.id)?.finalText, ""); + assert.equal(settlements, 1); + + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("fresh restarted result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await runtime.runPromise(manager.waitFor([spawned.id])); + assert.equal(manager.view.get(spawned.id)?.status, "done"); + assert.equal( + manager.view.get(spawned.id)?.finalText, + "fresh restarted result", + ); + assert.equal(settlements, 2); + } finally { + await runtime.dispose(); + } +}); + +test("preflight cancellation stays single-flight until the exact prompt settles", async () => { + const firstPrompt = deferred(); + const secondPrompt = deferred(); + let promptIndex = 0; + let concurrentPrompts = 0; + let maxConcurrentPrompts = 0; + const fixtures = harnessFactory(() => ({ + prompt: async () => { + const index = promptIndex++; + concurrentPrompts++; + maxConcurrentPrompts = Math.max(maxConcurrentPrompts, concurrentPrompts); + try { + await (index === 0 ? firstPrompt.promise : secondPrompt.promise); + } finally { + concurrentPrompts--; + } + }, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + let settlements = 0; + manager.view.setOnSettled(() => settlements++); + + const spawned = await runtime.runPromise( + manager.spawn("pi", task("cancel during preflight")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + assert.deepEqual(harness.calls.prompts, ["cancel during preflight"]); + await assert.rejects( + runtime.runPromise( + manager.send(spawned.id, "must not overlap preflight"), + ), + /prompt is still starting/, + ); + assert.deepEqual(harness.calls.prompts, ["cancel during preflight"]); + + let cancelFinished = false; + const cancellation = runtime + .runPromise(manager.cancel([spawned.id])) + .then((result) => { + cancelFinished = true; + return result; + }); + await waitFor(() => harness.calls.aborts === 1, "initial preflight abort"); + await assert.rejects( + runtime.runPromise(manager.send(spawned.id, "must not overlap cancel")), + /cancellation is still in progress/, + ); + assert.deepEqual(harness.calls.prompts, ["cancel during preflight"]); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(cancelFinished, false); + assert.equal(settlements, 0); + + // A cancelled Pi preflight can resume after abort() returned, then enter + // the agent lifecycle. These events still belong to prompt #1. + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + await waitFor( + () => harness.calls.aborts === 2, + "late preflight agent_start to be aborted again", + ); + harness.emitTextDelta("stale delta"); + harness.emit({ + type: "tool_execution_start", + toolCallId: "stale-tool", + toolName: "read", + args: { path: "stale.txt" }, + }); + harness.emitAssistant("stale cancelled result"); + harness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(cancelFinished, false); + assert.equal(settlements, 0); + + firstPrompt.resolve(); + const report = await cancellation; + assert.equal(report[0]?.cancelled, true); + assert.equal(settlements, 1); + assert.equal(manager.view.get(spawned.id)?.finalText, ""); + assert.equal(manager.view.get(spawned.id)?.transcript.length, 0); + + await runtime.runPromise(manager.send(spawned.id, "fresh restart")); + assert.deepEqual(harness.calls.prompts, [ + "cancel during preflight", + "fresh restart", + ]); + assert.equal(maxConcurrentPrompts, 1); + + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitTextDelta("fresh delta"); + harness.emitAssistant("fresh result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + secondPrompt.resolve(); + await runtime.runPromise(manager.waitFor([spawned.id])); + + const restarted = manager.view.get(spawned.id); + assert.equal(restarted?.status, "done"); + assert.equal(restarted?.finalText, "fresh result"); + assert.equal(settlements, 2); + assert.equal( + restarted?.transcript.some((item) => + JSON.stringify(item).includes("stale"), + ), + false, + ); + } finally { + firstPrompt.resolve(); + secondPrompt.resolve(); + await runtime.dispose(); + } +}); + +test("settled restart waits for the previous prompt Promise microtask boundary", async () => { + const firstPrompt = deferred(); + const secondPrompt = deferred(); + let promptIndex = 0; + let concurrentPrompts = 0; + let maxConcurrentPrompts = 0; + const fixtures = harnessFactory(() => ({ + prompt: async () => { + const index = promptIndex++; + concurrentPrompts++; + maxConcurrentPrompts = Math.max(maxConcurrentPrompts, concurrentPrompts); + try { + await (index === 0 ? firstPrompt.promise : secondPrompt.promise); + } finally { + concurrentPrompts--; + } + }, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("first prompt")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("first result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(manager.view.get(spawned.id)?.status, "running"); + + let restartAccepted = false; + const restart = runtime + .runPromise(manager.send(spawned.id, "second prompt")) + .then(() => { + restartAccepted = true; + }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(restartAccepted, false); + assert.deepEqual(harness.calls.prompts, ["first prompt"]); + + firstPrompt.resolve(); + await restart; + assert.deepEqual(harness.calls.prompts, ["first prompt", "second prompt"]); + assert.equal(maxConcurrentPrompts, 1); + + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("second result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + secondPrompt.resolve(); + await runtime.runPromise(manager.waitFor([spawned.id])); + assert.equal(manager.view.get(spawned.id)?.finalText, "second result"); + } finally { + firstPrompt.resolve(); + secondPrompt.resolve(); + await runtime.dispose(); + } +}); + +test("cancel invalidates a restart waiting on the previous prompt Promise", async () => { + const firstPrompt = deferred(); + const secondPrompt = deferred(); + let promptIndex = 0; + const fixtures = harnessFactory(() => ({ + prompt: () => + promptIndex++ === 0 ? firstPrompt.promise : secondPrompt.promise, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + let settlements = 0; + manager.view.setOnSettled(() => settlements++); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("completed but prompt still unwinding")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("first result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(manager.view.get(spawned.id)?.status, "running"); + + const restart = runtime.runPromise( + manager.send(spawned.id, "restart that will be cancelled"), + ); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.deepEqual(harness.calls.prompts, [ + "completed but prompt still unwinding", + ]); + + const cancellation = runtime.runPromise(manager.cancel([spawned.id])); + await waitFor( + () => harness.calls.aborts === 1, + "pending restart cancellation", + ); + firstPrompt.resolve(); + await restart; + const report = await cancellation; + + assert.equal(report[0]?.cancelled, true); + assert.deepEqual(harness.calls.prompts, [ + "completed but prompt still unwinding", + ]); + assert.equal(manager.view.get(spawned.id)?.status, "error"); + assert.equal(manager.view.get(spawned.id)?.finalText, ""); + assert.equal(settlements, 1); + } finally { + firstPrompt.resolve(); + secondPrompt.resolve(); + await runtime.dispose(); + } +}); + +test("an aborted send cannot start its deferred restart later", async () => { + const firstPrompt = deferred(); + const secondPrompt = deferred(); + let promptIndex = 0; + const fixtures = harnessFactory(() => ({ + prompt: () => + promptIndex++ === 0 ? firstPrompt.promise : secondPrompt.promise, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("completed before aborted send")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("stable prior result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(manager.view.get(spawned.id)?.status, "running"); + + const controller = new AbortController(); + const send = runtime.runPromise( + manager.send(spawned.id, "must never start"), + { signal: controller.signal }, + ); + await new Promise((resolve) => setTimeout(resolve, 20)); + controller.abort(); + await assert.rejects(send); + + firstPrompt.resolve(); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.deepEqual(harness.calls.prompts, ["completed before aborted send"]); + assert.equal(manager.view.get(spawned.id)?.status, "done"); + assert.equal( + manager.view.get(spawned.id)?.finalText, + "stable prior result", + ); + } finally { + firstPrompt.resolve(); + secondPrompt.resolve(); + await runtime.dispose(); + } +}); + +test("run-local terminal evidence survives preflight and post-run compaction", async () => { + const initialMessages = Array.from({ length: 64 }, (_, index) => ({ + role: "user" as const, + content: `old context ${index}`, + timestamp: index, + })) as AgentSession["messages"]; + const fixtures = harnessFactory(() => ({ initialMessages })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("complete after compaction")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + + // Pi may replace the entire message array during prompt preflight. + harness.messages.splice(0, harness.messages.length, { + role: "user", + content: "compacted context", + timestamp: Date.now(), + }); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("fresh completed result"); + // Post-run compaction can replace it again before agent_settled. + harness.messages.splice(0, harness.messages.length); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await runtime.runPromise(manager.waitFor([spawned.id])); + assert.equal( + manager.view.get(spawned.id)?.finalText, + "fresh completed result", + ); + + await runtime.runPromise(manager.send(spawned.id, "fail after compaction")); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("fresh failed partial", { + stopReason: "error", + errorMessage: "fresh provider failure", + }); + harness.messages.splice(0, harness.messages.length); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await runtime.runPromise(manager.waitFor([spawned.id])); + assert.equal(manager.view.get(spawned.id)?.status, "error"); + assert.equal( + manager.view.get(spawned.id)?.errorText, + "fresh provider failure", + ); + assert.equal( + manager.view.get(spawned.id)?.finalText, + "fresh failed partial", + ); + + await runtime.runPromise( + manager.send(spawned.id, "interrupt after compaction"), + ); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitTextDelta("fresh interrupted partial"); + harness.messages.splice(0, harness.messages.length); + const cancellation = runtime.runPromise(manager.cancel([spawned.id])); + await waitFor(() => harness.calls.aborts === 1, "compacted run abort"); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await cancellation; + assert.equal(manager.view.get(spawned.id)?.status, "error"); + assert.equal( + manager.view.get(spawned.id)?.finalText, + "fresh interrupted partial", + ); + } finally { + await runtime.dispose(); + } +}); + +test("cancelled restart preflight never reuses the previous run output", async () => { + const restartPreflight = deferred(); + let preflightIndex = 0; + const fixtures = harnessFactory(() => ({ + preflight: () => + preflightIndex++ === 0 ? Promise.resolve() : restartPreflight.promise, + })); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn("pi", task("successful first run")), + ); + const harness = fixtures.harnesses[0]; + assert.ok(harness); + harness.setStreaming(true); + harness.emit({ type: "agent_start" }); + harness.emitAssistant("previous successful result"); + harness.setStreaming(false); + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + await runtime.runPromise(manager.waitFor([spawned.id])); + assert.equal( + manager.view.get(spawned.id)?.finalText, + "previous successful result", + ); + + await runtime.runPromise(manager.send(spawned.id, "cancel this preflight")); + await waitFor( + () => harness.calls.prompts.length === 2, + "restart preflight to begin", + ); + const cancellation = runtime.runPromise(manager.cancel([spawned.id])); + await waitFor(() => harness.calls.aborts === 1, "restart preflight abort"); + restartPreflight.resolve(); + const report = await cancellation; + + assert.equal(report[0]?.cancelled, true); + assert.equal(manager.view.get(spawned.id)?.status, "error"); + assert.equal(manager.view.get(spawned.id)?.finalText, ""); + } finally { + restartPreflight.resolve(); + await runtime.dispose(); + } +}); + +test("a preflight that ignores abort is force-closed without reopening the session", async () => { + const stuckPreflight = deferred(); + let providerRuns = 0; + let cleanupCalls = 0; + const fixtures = harnessFactory((_options, index) => + index === 0 + ? { + preflight: () => stuckPreflight.promise, + prompt: async () => { + providerRuns++; + }, + } + : {}, + ); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + worktreeCleanup: async () => { + cleanupCalls++; + throw new Error("unsafe worktree cleanup should not run"); + }, + }); + const runtime = createManagerRuntime(backend, 10_000); + try { + const manager = await runtime.runPromise(SubagentManager); + const spawned = await runtime.runPromise( + manager.spawn( + "pi", + task("never-ending preflight", { + worktree: { + path: "C:/fixture-worktree", + branch: "pi/fixture-worktree", + repoCwd: process.cwd(), + }, + }), + ), + ); + const startedAt = Date.now(); + const report = await runtime.runPromise(manager.cancel([spawned.id])); + assert.ok(Date.now() - startedAt >= 4_500); + assert.equal(report[0]?.cancelled, true); + assert.match( + manager.view.get(spawned.id)?.errorText ?? "", + /Abort deadline exceeded/, + ); + assert.match( + manager.view.get(spawned.id)?.errorText ?? "", + /worktree cleanup skipped; checkout preserved at C:\/fixture-worktree/, + ); + assert.equal(cleanupCalls, 0); + await waitFor( + () => fixtures.harnesses[0]?.calls.disposals === 1, + "stuck preflight session disposal", + ); + await assert.rejects( + runtime.runPromise(manager.send(spawned.id, "must stay closed")), + /session is closed/, + ); + assert.deepEqual(fixtures.harnesses[0]?.calls.prompts, [ + "never-ending preflight", + ]); + stuckPreflight.resolve(); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(providerRuns, 0); + + const fresh = await runtime.runPromise( + manager.spawn("pi", task("fresh after forced close")), + ); + assert.equal(fresh.status, "running"); + const freshCancellation = runtime.runPromise(manager.cancel([fresh.id])); + const freshHarness = fixtures.harnesses[1]; + assert.ok(freshHarness); + await waitFor(() => freshHarness.calls.aborts === 1, "fresh run abort"); + freshHarness.emit({ type: "agent_settled" }); + freshHarness.resolvePrompt(); + await freshCancellation; + } finally { + stuckPreflight.resolve(); + await runtime.dispose(); + } +}); + +test("prompt and tool cancellation ignore late completion and free capacity", async () => { + const latePrompt = deferred(); + const fixtures = harnessFactory((_options, index) => + index === 0 ? { prompt: () => latePrompt.promise } : {}, + ); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 50, + }); + const runtime = createManagerRuntime(backend); + try { + const manager = await runtime.runPromise(SubagentManager); + let settlements = 0; + manager.view.setOnSettled(() => settlements++); + const spawned = await Promise.all( + [0, 1, 2, 3].map((index) => + runtime.runPromise(manager.spawn("pi", task(`cancel ${index}`))), + ), + ); + for (const harness of fixtures.harnesses) harness.setStreaming(true); + for (const harness of fixtures.harnesses.slice(2)) { + harness.emit({ type: "agent_start" }); + harness.emitTextDelta("tool run started"); + harness.emit({ + type: "tool_execution_start", + toolCallId: `tool-${fixtures.harnesses.indexOf(harness)}`, + toolName: "read", + args: { path: "README.md" }, + }); + } + await waitFor( + () => + spawned + .slice(2) + .every( + (snapshot) => manager.view.get(snapshot.id)?.liveTools.length === 1, + ), + "tool starts to reach the manager", + ); + + let cancellationFinished = false; + const cancellation = runtime + .runPromise(manager.cancel(spawned.map((snapshot) => snapshot.id))) + .then((result) => { + cancellationFinished = true; + return result; + }); + await waitFor( + () => fixtures.harnesses.every((harness) => harness.calls.aborts === 1), + "all prompt and tool aborts", + ); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(cancellationFinished, false); + latePrompt.reject(new Error("late prompt rejection")); + for (const harness of fixtures.harnesses.slice(1)) { + harness.emit({ type: "agent_settled" }); + harness.resolvePrompt(); + } + const report = await cancellation; + assert.ok(report.every((result) => result.cancelled)); + assert.ok( + spawned.every( + (snapshot) => + manager.view.get(snapshot.id)?.errorText === "Run was aborted", + ), + ); + assert.equal(settlements, 4); + + const promptHarness = fixtures.harnesses[0]; + const toolHarness = fixtures.harnesses[2]; + assert.ok(promptHarness); + assert.ok(toolHarness); + promptHarness.emit({ type: "agent_start" }); + promptHarness.emitAssistant("late prompt completion"); + promptHarness.emit({ type: "agent_settled" }); + toolHarness.emit({ + type: "tool_execution_end", + toolCallId: "tool-2", + toolName: "read", + result: { content: [{ type: "text", text: "late tool result" }] }, + isError: false, + }); + toolHarness.emitAssistant("late tool completion"); + toolHarness.emit({ type: "agent_settled" }); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(settlements, 4); + assert.equal(manager.view.get(spawned[0]!.id)?.finalText, ""); + assert.equal( + manager.view.get(spawned[2]!.id)?.finalText, + "tool run started", + ); + + const afterCancellation = await runtime.runPromise( + manager.spawn("pi", task("capacity was released")), + ); + assert.equal(afterCancellation.status, "running"); + const afterCancellationStop = runtime.runPromise( + manager.cancel([afterCancellation.id]), + ); + const afterCancellationHarness = fixtures.harnesses[4]; + assert.ok(afterCancellationHarness); + await waitFor( + () => afterCancellationHarness.calls.aborts === 1, + "replacement run abort", + ); + afterCancellationHarness.emit({ type: "agent_settled" }); + afterCancellationHarness.resolvePrompt(); + await afterCancellationStop; + } finally { + await runtime.dispose(); + } +}); + +test("a silent provider is terminalized once, disposed, and releases all slots", async () => { + const never = new Promise(() => {}); + const fixtures = harnessFactory((_options, index) => { + if (index >= 4) return {}; + if (index % 2 === 0) return { shutdown: () => never }; + return { + shutdown: async () => { + throw new Error("shutdown fixture failed"); + }, + dispose: () => { + throw new Error("dispose fixture failed"); + }, + }; + }); + const backend = makePiBackend({ + sessionFactory: fixtures.factory, + shutdownTimeoutMs: 30, + }); + const runtime = createManagerRuntime(backend, 40); + try { + const manager = await runtime.runPromise(SubagentManager); + let settlements = 0; + manager.view.setOnSettled(() => settlements++); + const stalled = await Promise.all( + [0, 1, 2, 3].map((index) => + runtime.runPromise(manager.spawn("pi", task(`silent ${index}`))), + ), + ); + + await assert.rejects( + runtime.runPromise(manager.spawn("pi", task("over capacity"))), + /Max 4 subagent sessions/, + ); + await runtime.runPromise( + manager.waitFor(stalled.map((snapshot) => snapshot.id)), + ); + for (const snapshot of stalled) { + const failed = manager.view.get(snapshot.id); + assert.equal(failed?.status, "error"); + assert.match(failed?.errorText ?? "", /no assistant response event/); + } + assert.equal(settlements, 4); + await waitFor( + () => + fixtures.harnesses.every((harness) => harness.calls.disposals === 1), + "silent sessions to be disposed despite cleanup failures", + ); + + const late = fixtures.harnesses[0]; + assert.ok(late); + late.emitAssistant("late provider completion"); + late.emit({ type: "agent_settled" }); + late.resolvePrompt(); + await new Promise((resolve) => setTimeout(resolve, 20)); + assert.equal(settlements, 4); + assert.equal(manager.view.get(stalled[0]!.id)?.finalText, ""); + + const fresh = await runtime.runPromise( + manager.spawn("pi", task("fresh after watchdog")), + ); + assert.equal(fresh.status, "running"); + const freshCancellation = runtime.runPromise(manager.cancel([fresh.id])); + const freshHarness = fixtures.harnesses[4]; + assert.ok(freshHarness); + await waitFor(() => freshHarness.calls.aborts === 1, "fresh run abort"); + freshHarness.emit({ type: "agent_settled" }); + freshHarness.resolvePrompt(); + await freshCancellation; + } finally { + await runtime.dispose(); + } +}); + +test("adapter scope cleanup is bounded across shutdown timeouts and failures", async () => { + const never = new Promise(() => {}); + const timedOutFixtures = harnessFactory(() => ({ + prompt: async () => {}, + shutdown: () => never, + })); + const timedOutBackend = makePiBackend({ + sessionFactory: timedOutFixtures.factory, + shutdownTimeoutMs: 25, + }); + const timeoutStartedAt = Date.now(); + await Effect.runPromise( + Effect.scoped(timedOutBackend.spawn(task("timeout cleanup"))), + ); + assert.ok(Date.now() - timeoutStartedAt < 500); + const timedOutHarness = timedOutFixtures.harnesses[0]; + assert.ok(timedOutHarness); + assert.equal(timedOutHarness.calls.shutdowns, 1); + assert.equal(timedOutHarness.calls.disposals, 1); + + const failedFixtures = harnessFactory(() => ({ + prompt: async () => {}, + shutdown: async () => { + throw new Error("shutdown fixture failed"); + }, + dispose: () => { + throw new Error("dispose fixture failed"); + }, + })); + const failedBackend = makePiBackend({ + sessionFactory: failedFixtures.factory, + shutdownTimeoutMs: 25, + }); + await Effect.runPromise( + Effect.scoped(failedBackend.spawn(task("failed cleanup"))), + ); + const failedHarness = failedFixtures.harnesses[0]; + assert.ok(failedHarness); + assert.equal(failedHarness.calls.shutdowns, 1); + assert.equal(failedHarness.calls.disposals, 1); +}); diff --git a/tests/support/pi-agent-session-harness.ts b/tests/support/pi-agent-session-harness.ts new file mode 100644 index 00000000..8c874d0e --- /dev/null +++ b/tests/support/pi-agent-session-harness.ts @@ -0,0 +1,285 @@ +/** + * Minimal controllable Pi AgentSession for production-adapter integration + * tests. Tests drive native AgentSession events explicitly; this fake records + * SDK calls but deliberately does not reproduce Pi or subagent lifecycle state. + */ + +import type { AssistantMessage } from "@earendil-works/pi-ai"; +import type { + AgentSession, + AgentSessionEvent, + AgentSessionEventListener, +} from "@earendil-works/pi-coding-agent"; + +const ZERO_USAGE = { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + total: 0, + }, +}; + +export interface PiAgentSessionHarnessOptions { + readonly model?: AgentSession["model"]; + readonly initialMessages?: AgentSession["messages"]; + readonly activeTools?: readonly string[]; + readonly contextUsage?: { + readonly tokens: number | null; + readonly contextWindow: number; + }; + readonly sessionFile?: string; + readonly bind?: () => Promise; + readonly preflight?: ( + text: string, + harness: PiAgentSessionHarness, + ) => Promise; + readonly prompt?: ( + text: string, + harness: PiAgentSessionHarness, + ) => Promise; + readonly steer?: ( + text: string, + harness: PiAgentSessionHarness, + ) => Promise; + readonly abort?: (harness: PiAgentSessionHarness) => Promise; + readonly shutdown?: (harness: PiAgentSessionHarness) => Promise; + readonly dispose?: (harness: PiAgentSessionHarness) => void; +} + +export interface PiAgentSessionHarness { + readonly session: AgentSession; + readonly messages: AgentSession["messages"]; + readonly calls: { + readonly bindings: Array<{ mode: "print" }>; + readonly prompts: string[]; + readonly steers: string[]; + readonly sessionNames: string[]; + clearQueues: number; + aborts: number; + shutdowns: number; + disposals: number; + }; + emit(event: AgentSessionEvent): void; + /** Explicitly settle the oldest default prompt Promise. */ + resolvePrompt(): void; + emitUser(text: string): void; + emitAssistant( + text: string, + options?: { + readonly stopReason?: AssistantMessage["stopReason"]; + readonly errorMessage?: string; + readonly provider?: string; + readonly model?: string; + }, + ): AssistantMessage; + emitTextDelta(delta: string): void; + setStreaming(streaming: boolean): void; + setContextUsage( + usage: + | { readonly tokens: number | null; readonly contextWindow: number } + | undefined, + ): void; + activeTools(): string[]; +} + +export function createPiAgentSessionHarness( + options: PiAgentSessionHarnessOptions = {}, +) { + const listeners = new Set(); + const messages: AgentSession["messages"] = [ + ...(options.initialMessages ?? []), + ]; + const calls = { + bindings: [] as Array<{ mode: "print" }>, + prompts: [] as string[], + steers: [] as string[], + sessionNames: [] as string[], + clearQueues: 0, + aborts: 0, + shutdowns: 0, + disposals: 0, + }; + let streaming = false; + let contextUsage = options.contextUsage; + let activeTools = [...(options.activeTools ?? [])]; + let harness: PiAgentSessionHarness; + const pendingDefaultPrompts: Array<() => void> = []; + let queuedPromptResolutions = 0; + + const emit = (event: AgentSessionEvent) => { + for (const listener of [...listeners]) listener(event); + }; + + const makeAssistantMessage = ( + text: string, + messageOptions: { + readonly stopReason?: AssistantMessage["stopReason"]; + readonly errorMessage?: string; + readonly provider?: string; + readonly model?: string; + } = {}, + ) => + ({ + role: "assistant", + content: [{ type: "text", text }], + api: "openai-responses", + provider: messageOptions.provider ?? options.model?.provider ?? "fixture", + model: messageOptions.model ?? options.model?.id ?? "fixture-model", + usage: ZERO_USAGE, + stopReason: messageOptions.stopReason ?? "stop", + ...(messageOptions.errorMessage + ? { errorMessage: messageOptions.errorMessage } + : {}), + timestamp: Date.now(), + }) satisfies AssistantMessage; + + const session = { + messages, + get model() { + return options.model; + }, + get isStreaming() { + return streaming; + }, + get sessionFile() { + return options.sessionFile ?? "fixture-session.jsonl"; + }, + sessionManager: { + appendSessionInfo(name: string) { + calls.sessionNames.push(name); + }, + }, + extensionRunner: { + hasHandlers(eventType: string) { + return ( + eventType === "session_shutdown" && options.shutdown !== undefined + ); + }, + async emit() { + calls.shutdowns++; + await options.shutdown?.(harness); + }, + }, + async bindExtensions(bindings: { mode: "print" }) { + calls.bindings.push(bindings); + await options.bind?.(); + }, + subscribe(listener: AgentSessionEventListener) { + listeners.add(listener); + return () => listeners.delete(listener); + }, + async prompt( + text: string, + promptOptions?: Parameters[1], + ) { + calls.prompts.push(text); + try { + await options.preflight?.(text, harness); + } catch (error) { + promptOptions?.preflightResult?.(false); + throw error; + } + promptOptions?.preflightResult?.(true); + if (options.prompt) { + await options.prompt(text, harness); + return; + } + await new Promise((resolve) => { + if (queuedPromptResolutions > 0) { + queuedPromptResolutions--; + resolve(); + return; + } + pendingDefaultPrompts.push(resolve); + }); + }, + async steer(text: string) { + calls.steers.push(text); + await options.steer?.(text, harness); + }, + clearQueue() { + calls.clearQueues++; + return { steering: [], followUp: [] }; + }, + async abort() { + calls.aborts++; + await options.abort?.(harness); + streaming = false; + }, + dispose() { + calls.disposals++; + options.dispose?.(harness); + }, + getContextUsage() { + return contextUsage; + }, + getAllTools() { + return activeTools.map((name) => ({ name })); + }, + getToolDefinition() { + return undefined; + }, + getActiveToolNames() { + return [...activeTools]; + }, + setActiveToolsByName(toolNames: string[]) { + activeTools = [...toolNames]; + }, + } as unknown as AgentSession; + + harness = { + session, + messages, + calls, + emit, + resolvePrompt() { + const resolve = pendingDefaultPrompts.shift(); + if (resolve) resolve(); + else queuedPromptResolutions++; + }, + emitUser(text) { + const message = { + role: "user", + content: text, + timestamp: Date.now(), + } satisfies AgentSession["messages"][number]; + messages.push(message); + emit({ type: "message_end", message }); + }, + emitAssistant(text, messageOptions) { + const message = makeAssistantMessage(text, messageOptions); + messages.push(message); + emit({ type: "message_end", message }); + return message; + }, + emitTextDelta(delta) { + const partial = makeAssistantMessage(delta); + emit({ + type: "message_update", + message: partial, + assistantMessageEvent: { + type: "text_delta", + contentIndex: 0, + delta, + partial, + }, + }); + }, + setStreaming(next) { + streaming = next; + }, + setContextUsage(next) { + contextUsage = next; + }, + activeTools: () => [...activeTools], + }; + + return harness; +}