/** * Integration tests for the `stop_all_background` kill switch against REAL detached runner processes (mock pi child): * the stop request is only a file in the run's control inbox, the runner consumes it by itself, ends with status `stopped` * and exits. Runs of other sessions are untouched. * * Requires pi packages to be importable. Skips gracefully if unavailable. */ import { describe, it } from "node:test"; import assert from "node:assert/strict"; import * as fs from "node:fs"; import * as path from "node:path"; import { createEventBus } from "../../../../core/event-bus.ts"; import { RpcReleaseBridge } from "../../../../modes/rpc/release-bridge.ts"; import { registerSubagentRpcBridge } from "../../src/extension/rpc.ts"; import { registerReleaseReadiness } from "../../src/extension/release-readiness.ts"; import { makeAgent, events } from "../support/helpers.ts"; import type { AsyncJobState, SubagentState } from "../../src/shared/types.ts"; import type { ExtensionContext } from "@selesai/code"; import { installAsyncExecutionHooks, available, isAsyncAvailable, executeAsyncSingle, ASYNC_DIR, RESULTS_DIR, waitForAsyncResultFile, waitForAsyncState, tempDir, mockPi, } from "../support/async-execution-fixture.ts"; const skip = !isAsyncAvailable() ? "jiti not available" : undefined; const SESSION = "session-1"; function launch(id: string, sessionId: string) { const receipt = executeAsyncSingle(id, { agent: "worker", task: "Do the work", agentConfig: makeAgent("worker"), ctx: { pi: { events: { emit() {} } }, cwd: tempDir, currentSessionId: sessionId }, artifactConfig: { enabled: false, includeInput: false, includeOutput: false, includeJsonl: false, includeMetadata: false, cleanupDays: 7 }, shareEnabled: false, sessionRoot: path.join(tempDir, "sessions"), maxSubagentDepth: 2, }); assert.equal(receipt.isError === true, false, "launch must succeed"); } const readRunStatus = (id: string) => JSON.parse(fs.readFileSync(path.join(ASYNC_DIR, id, "status.json"), "utf-8")) as { state: string; pid?: number }; const job = (id: string, sessionId: string, pid?: number): AsyncJobState => ({ asyncId: id, asyncDir: path.join(ASYNC_DIR, id), status: "running", startedAt: Date.now(), updatedAt: Date.now(), sessionId, pid } as AsyncJobState); const alive = (pid: number) => { try { process.kill(pid, 0); return true; } catch { return false; } }; async function waitFor(check: () => boolean, ms: number, what: string): Promise { const until = Date.now() + ms; while (!check()) { if (Date.now() > until) assert.fail(`timed out waiting for ${what}`); await new Promise((resolve) => setTimeout(resolve, 25)); } } describe("stop_all_background against real detached runners", { skip: !available ? "pi packages not available" : undefined }, () => { installAsyncExecutionHooks(); it("stops this session's runner (status stopped, process gone) and leaves another session's runner running", { timeout: 60_000, skip }, async () => { const mine = `stopall-mine-${Date.now().toString(36)}`; const theirs = `stopall-theirs-${Date.now().toString(36)}`; const releaseTheirs = path.join(tempDir, `${theirs}.release`); const bash = (n: number) => ({ type: "tool_execution_start", toolCallId: `bash-${n}`, toolName: "bash", args: { command: "sleep 600" } }); mockPi.onCall({ steps: [{ jsonl: [bash(1)] }, { waitForPath: path.join(tempDir, `${mine}.never`), jsonl: [] }] }); launch(mine, SESSION); await waitForAsyncState(mine, (status) => status.state === "running" && status.steps?.[0]?.status === "running"); mockPi.onCall({ steps: [{ waitForPath: releaseTheirs, jsonl: [events.assistantMessage("Done")] }] }); launch(theirs, "session-2"); await waitForAsyncState(theirs, (status) => status.state === "running" && status.steps?.[0]?.status === "running"); const minePid = readRunStatus(mine).pid!; const theirsPid = readRunStatus(theirs).pid!; assert.ok(alive(minePid) && alive(theirsPid)); // Session-1's extension state tracks its own run, and (adversarially) also lists session-2's run. const state = { baseCwd: tempDir, currentSessionId: SESSION, supervisorOwnerSessionId: "host", asyncJobs: new Map([[mine, job(mine, SESSION, minePid)], [theirs, job(theirs, "session-2", theirsPid)]]), foregroundControls: new Map(), lastForegroundControlId: null } as unknown as SubagentState; const bus = createEventBus(); const ctx = { sessionManager: { getSessionFile: () => SESSION, getSessionId: () => SESSION } } as unknown as ExtensionContext; const rpc = registerSubagentRpcBridge({ events: bus, getContext: () => ctx, execute: async () => { throw new Error("unused"); }, asyncDirRoot: ASYNC_DIR, resultsDir: RESULTS_DIR, state }); registerReleaseReadiness({ events: bus, getHostSessionId: () => "host", subagents: { state, hasPendingDelivery: () => false, listUndeliveredResults: () => [], resultsDir: RESULTS_DIR }, scheduler: { armedSchedules: () => [], observedCompletionRunIds: () => [] }, stopAll: { subagents: { state, resultsDir: RESULTS_DIR, asyncDirRoot: ASYNC_DIR, observedCompletionRunIds: () => [], stopRun: (target) => rpc.stopRun(target) }, scheduler: { disarmSessionSchedules: () => [] }, }, }); const bridge = new RpcReleaseBridge(bus, { version: 1, sessionId: "host", epoch: "e1" }); const reply = await bridge.stopAll({ type: "stop_all_background", version: 1 }, () => [], () => []); assert.ok(reply.success && reply.command === "stop_all_background", JSON.stringify(reply)); assert.deepEqual(reply.data.stopped.map((item) => [item.source, item.kind, item.id, item.ok, item.error]), [ ["subagents", "async_run", theirs, false, "foreign_session"], ["subagents", "async_run", mine, true, undefined], ]); // The runner ended with status stopped (written by the detached runner itself) ... assert.equal(readRunStatus(mine).state, "stopped"); const result = JSON.parse(fs.readFileSync(await waitForAsyncResultFile(mine), "utf-8")) as { state?: string; success?: boolean; stopped?: boolean }; assert.equal(result.state, "stopped"); assert.equal(result.success, false); // ... and the process is gone, not merely idle. await waitFor(() => !alive(minePid), 15_000, "the stopped runner process to exit"); // The other session's runner is untouched and still completes normally. assert.equal(readRunStatus(theirs).state, "running"); assert.ok(alive(theirsPid)); assert.equal(reply.data.remaining.blockers.some((b) => b.id === mine && b.kind === "async_run"), false, "the stopped run is no longer an async_run blocker"); // Idempotent: the run is terminal now, nothing is stopped again (the foreign one is reported again, never touched). const again = await bridge.stopAll({ type: "stop_all_background", version: 1 }, () => [], () => []); assert.ok(again.success && again.command === "stop_all_background"); assert.deepEqual(again.data.stopped.map((item) => [item.id, item.ok]), [[theirs, false]]); fs.writeFileSync(releaseTheirs, ""); const done = JSON.parse(fs.readFileSync(await waitForAsyncResultFile(theirs), "utf-8")) as { state?: string; success?: boolean }; assert.equal(done.success, true); }); it("stops a run that is not tracked in memory but awaited by a scheduled run (scheduled_run_pending)", { timeout: 60_000, skip }, async () => { const id = `stopall-sched-${Date.now().toString(36)}`; mockPi.onCall({ steps: [{ jsonl: [{ type: "tool_execution_start", toolCallId: "bash-1", toolName: "bash", args: { command: "sleep 600" } }] }, { waitForPath: path.join(tempDir, `${id}.never`), jsonl: [] }] }); launch(id, SESSION); await waitForAsyncState(id, (status) => status.state === "running" && status.steps?.[0]?.status === "running"); const pid = readRunStatus(id).pid!; const state = { baseCwd: tempDir, currentSessionId: SESSION, supervisorOwnerSessionId: "host", asyncJobs: new Map(), foregroundControls: new Map(), lastForegroundControlId: null } as unknown as SubagentState; const bus = createEventBus(); const ctx = { sessionManager: { getSessionFile: () => SESSION, getSessionId: () => SESSION } } as unknown as ExtensionContext; const rpc = registerSubagentRpcBridge({ events: bus, getContext: () => ctx, execute: async () => { throw new Error("unused"); }, asyncDirRoot: ASYNC_DIR, resultsDir: RESULTS_DIR, state }); registerReleaseReadiness({ events: bus, getHostSessionId: () => "host", subagents: { state, hasPendingDelivery: () => false, listUndeliveredResults: () => [], resultsDir: RESULTS_DIR }, scheduler: { armedSchedules: () => [], observedCompletionRunIds: () => [id] }, stopAll: { subagents: { state, resultsDir: RESULTS_DIR, asyncDirRoot: ASYNC_DIR, observedCompletionRunIds: () => [id], stopRun: (target) => rpc.stopRun(target) }, scheduler: { disarmSessionSchedules: () => [] }, }, }); const reply = await new RpcReleaseBridge(bus, { version: 1, sessionId: "host", epoch: "e1" }).stopAll({ type: "stop_all_background", version: 1 }, () => [], () => []); assert.ok(reply.success && reply.command === "stop_all_background", JSON.stringify(reply)); assert.deepEqual(reply.data.stopped.map((item) => [item.kind, item.id, item.ok]), [["scheduled_run", id, true]]); assert.equal(readRunStatus(id).state, "stopped"); await waitFor(() => !alive(pid), 15_000, "the stopped runner process to exit"); }); });