import assert from "node:assert/strict"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; import { afterEach, describe, it } from "node:test"; import type { ExtensionContext } from "@selesai/code"; import { createEventBus } from "../../../../core/event-bus.ts"; import { RpcReleaseBridge } from "../../../../modes/rpc/release-bridge.ts"; import { RELEASE_READINESS_EVENT, STOP_ALL_BACKGROUND_EVENT, findUndeliveredResults, registerReleaseReadiness, schedulerReleaseBlockers, stopSchedules, stopSubagentWork, subagentReleaseBlockers, type SubagentReleaseDeps, type SubagentStopDeps, } from "../../src/extension/release-readiness.ts"; import { reconcileAsyncRun } from "../../src/runs/background/stale-run-reconciler.ts"; import { writeAsyncResultFile } from "../../src/runs/background/result-files.ts"; import { createResultDeliveryOwnership } from "../../src/runs/background/result-delivery-ownership.ts"; import { ownerRecordPath } from "../../src/shared/owner-record.ts"; import { createScheduledRunManager } from "../../src/runs/background/scheduled-runs.ts"; import type { AsyncJobState, SubagentState } from "../../src/shared/types.ts"; const roots: string[] = []; afterEach(() => { for (const root of roots.splice(0)) fs.rmSync(root, { recursive: true, force: true }); }); const tempRoot = (): string => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "pi-release-test-")); roots.push(root); return root; }; function makeState(jobs: Array & { asyncId: string }> = []): SubagentState { return { currentSessionId: "/managed/session", supervisorOwnerSessionId: "host", asyncJobs: new Map(jobs.map((job) => [job.asyncId, { asyncDir: `/nowhere/${job.asyncId}`, status: "running", startedAt: 1000, updatedAt: 1000, ...job } as AsyncJobState])), foregroundControls: new Map(), } as unknown as SubagentState; } function deps(state: SubagentState, overrides: Partial = {}): SubagentReleaseDeps { return { state, hasPendingDelivery: () => false, listUndeliveredResults: () => [], ...overrides }; } const dead = () => { throw Object.assign(new Error("no such process"), { code: "ESRCH" }); }; const alive = () => true; describe("subagent release blockers", () => { it("is empty for an idle session", () => { assert.deepEqual(subagentReleaseBlockers(deps(makeState())), []); }); it("reports a live running job as async_run with its id and start time", () => { const state = makeState([{ asyncId: "run-live", status: "running", startedAt: 1234, mode: "chain", agents: ["scout", "worker"] }]); const blockers = subagentReleaseBlockers(deps(state, { reconcile: () => ({ status: { runId: "run-live", state: "running" } as never, repaired: false }), })); assert.deepEqual(blockers, [{ source: "subagents", kind: "async_run", id: "run-live", detail: "Async chain running: scout, worker", since: 1234 }]); }); it("treats a job whose status is not yet written (within grace) as live", () => { const state = makeState([{ asyncId: "run-new", status: "queued" }]); const blockers = subagentReleaseBlockers(deps(state, { reconcile: () => ({ status: null, repaired: false }) })); assert.deepEqual(blockers.map((b) => [b.kind, b.id]), [["async_run", "run-new"]]); }); it("reconciles a dead-pid running job first and reports its repaired result as undelivered_result, not async_run", () => { const state = makeState([{ asyncId: "run-dead", status: "running", pid: 999_999 }]); const seen: string[] = []; const blockers = subagentReleaseBlockers(deps(state, { reconcile: (asyncDir) => { seen.push(asyncDir); return { status: { runId: "run-dead", state: "failed", endedAt: 2000 } as never, repaired: true }; }, })); assert.deepEqual(seen, ["/nowhere/run-dead"]); assert.equal(blockers.some((b) => b.kind === "async_run"), false); assert.deepEqual(blockers.map((b) => [b.source, b.kind, b.id, b.since]), [["subagents", "undelivered_result", "run-dead", 2000]]); assert.match(blockers[0]!.detail, /without a live runner/); }); it("runs the real reconciler against a dead runner pid and surfaces the persisted failure result once", () => { const root = tempRoot(); const asyncDir = path.join(root, "async", "run-dead"); const resultsDir = path.join(root, "results"); fs.mkdirSync(asyncDir, { recursive: true }); fs.mkdirSync(resultsDir, { recursive: true }); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify({ runId: "run-dead", sessionId: "/managed/session", completionOwnerId: "owner-1", mode: "single", state: "running", pid: 999_999, startedAt: 1000, lastUpdate: 1000, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: 1000 }], })); const state = makeState([{ asyncId: "run-dead", asyncDir, status: "running", pid: 999_999, sessionId: "/managed/session", completionOwnerId: "owner-1" }]); const real = (dir: string, options: Parameters[1]) => reconcileAsyncRun(dir, { ...options, kill: dead }); const owns = (sessionId: string, owner: unknown) => sessionId === "/managed/session" && owner === "owner-1"; const blockers = subagentReleaseBlockers(deps(state, { resultsDir, reconcile: real, listUndeliveredResults: () => findUndeliveredResults({ resultsDir, sessionIds: ["/managed/session"], owns }), })); assert.deepEqual(blockers.map((b) => [b.kind, b.id]), [["undelivered_result", "run-dead"]]); // The repair is durable: a second call still reports it (until delivered), and the status is terminal on disk. assert.equal(JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")).state, "failed"); assert.deepEqual(findUndeliveredResults({ resultsDir, sessionIds: ["/managed/session"], owns }), ["run-dead"]); }); it("keeps a run with a live pid as async_run through the real reconciler", () => { const root = tempRoot(); const asyncDir = path.join(root, "async", "run-alive"); fs.mkdirSync(asyncDir, { recursive: true }); const now = Date.now(); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify({ runId: "run-alive", sessionId: "/managed/session", mode: "single", state: "running", pid: process.pid, startedAt: now, lastUpdate: now, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: now }], })); const state = makeState([{ asyncId: "run-alive", asyncDir, status: "running", pid: process.pid, startedAt: now }]); const blockers = subagentReleaseBlockers(deps(state, { resultsDir: path.join(root, "results"), reconcile: (dir, options) => reconcileAsyncRun(dir, { ...options, kill: alive }), })); assert.deepEqual(blockers.map((b) => [b.kind, b.id, b.since]), [["async_run", "run-alive", now]]); }); it("fails closed when reconciliation throws: contributor_error plus the job stays an async_run", () => { const state = makeState([{ asyncId: "run-bad" }]); const blockers = subagentReleaseBlockers(deps(state, { reconcile: () => { throw new Error("status.json is corrupt"); } })); assert.deepEqual(blockers.map((b) => [b.kind, b.id]), [["contributor_error", "run-bad"], ["async_run", "run-bad"]]); assert.match(blockers[0]!.detail, /status\.json is corrupt/); }); it("does not reconcile jobs that are already terminal in memory", () => { const state = makeState([{ asyncId: "run-done", status: "complete" }]); const blockers = subagentReleaseBlockers(deps(state, { reconcile: () => { throw new Error("must not be called"); } })); assert.deepEqual(blockers, []); }); it("reports foreground runs, retained nested routes, nested descendants of finished jobs and pending notifications", () => { const state = makeState([{ asyncId: "run-nested", status: "complete", nestedChildren: [{ id: "child", state: "running" }] as never }]); state.foregroundControls.set("fg-1", { runId: "fg-1", mode: "single", startedAt: 77, updatedAt: 77, currentAgent: "scout" } as never); state.retainedForegroundNestedRoutes = new Map([["root-1", {} as never]]); const blockers = subagentReleaseBlockers(deps(state, { hasPendingDelivery: () => true })); assert.deepEqual(blockers.map((b) => [b.kind, b.id]), [ ["async_run", "run-nested"], ["foreground_run", "fg-1"], ["async_run", "root-1"], ["undelivered_result", undefined], ]); assert.equal(blockers.find((b) => b.kind === "foreground_run")!.since, 77); }); it("reports persisted undelivered results once and does not double-count a reconciled run", () => { const state = makeState([{ asyncId: "run-dead" }]); const blockers = subagentReleaseBlockers(deps(state, { reconcile: () => ({ status: { state: "failed" } as never, repaired: true }), listUndeliveredResults: () => ["run-dead", "run-old"], })); assert.deepEqual(blockers.map((b) => [b.kind, b.id]), [["undelivered_result", "run-dead"], ["undelivered_result", "run-old"]]); }); it("uses delivery demand only as a catch-all", () => { assert.deepEqual(subagentReleaseBlockers(deps(makeState(), { hasDeliveryDemand: () => false })), []); const blockers = subagentReleaseBlockers(deps(makeState(), { hasDeliveryDemand: () => true })); assert.deepEqual(blockers.map((b) => b.kind), ["undelivered_result"]); }); it("reports a probe failure instead of dropping it", () => { const blockers = subagentReleaseBlockers(deps(makeState(), { hasPendingDelivery: () => { throw new Error("notifier gone"); } })); assert.deepEqual(blockers.map((b) => b.kind), ["contributor_error"]); assert.match(blockers[0]!.detail, /notifier gone/); }); }); describe("undelivered result discovery", () => { const write = (resultsDir: string, runId: string, extra: Record = {}) => writeAsyncResultFile(path.join(resultsDir, `${runId}.json`), { id: runId, runId, sessionId: "/managed/session", completionOwnerId: "owner-1", success: false, state: "failed", ...extra, }); const owns = (sessionId: string, owner: unknown) => sessionId === "/managed/session" && owner === "owner-1"; it("lists owned undelivered results and skips delivered or foreign-owned ones", () => { const resultsDir = path.join(tempRoot(), "results"); fs.mkdirSync(resultsDir, { recursive: true }); write(resultsDir, "mine"); write(resultsDir, "delivered", { notificationDeliveredAt: 5 }); write(resultsDir, "foreign", { completionOwnerId: "someone-else" }); assert.deepEqual(findUndeliveredResults({ resultsDir, sessionIds: ["/managed/session", null, undefined], owns }), ["mine"]); assert.deepEqual(findUndeliveredResults({ resultsDir, sessionIds: [], owns }), []); }); it("reports a foreign result this process will take over, but not one with a live owner, and never claims", () => { const root = tempRoot(); const resultsDir = path.join(root, "results"); const ownersDir = path.join(root, "owners"); fs.mkdirSync(resultsDir, { recursive: true }); fs.mkdirSync(ownersDir, { recursive: true }); const record = (id: string, pid: number) => fs.writeFileSync(ownerRecordPath(ownersDir, id)!, JSON.stringify({ version: 1, completionOwnerId: id, pid, hostname: "host-a", createdAt: 1 })); record("dead-owner", 100); record("live-owner", 200); write(resultsDir, "mine", { completionOwnerId: "me" }); write(resultsDir, "orphan", { completionOwnerId: "dead-owner" }); write(resultsDir, "orphan-delivered", { completionOwnerId: "dead-owner", notificationDeliveredAt: 5 }); write(resultsDir, "busy", { completionOwnerId: "live-owner" }); write(resultsDir, "unknown-owner", { completionOwnerId: "no-record" }); write(resultsDir, "other-session", { completionOwnerId: "dead-owner", sessionId: "/other/session" }); const ownership = createResultDeliveryOwnership({ currentSessionId: "/managed/session", completionOwnerId: "me" }, { ownersDir, hostname: () => "host-a", processLiveness: (pid) => (pid === 100 ? "dead" : "alive"), processStartKey: () => undefined, }); const options = { resultsDir, sessionIds: ["/managed/session", "/other/session"], owns: (sessionId: string, owner: unknown) => ownership.owns(sessionId, owner), canTakeOver: (sessionId: string, owner: unknown, runKey: string) => ownership.canTakeOver(sessionId, owner, runKey), }; // "other-session" is only scanned because the caller listed it; ownership still refuses to take it over. assert.deepEqual(findUndeliveredResults(options).sort(), ["mine", "orphan"]); assert.equal(fs.existsSync(path.join(ownersDir, "claims")), false, "the readiness probe is read-only"); // Once the watcher has actually claimed it, the result is still reported until it is delivered. assert.equal(ownership.owns("/managed/session", "dead-owner", "orphan"), true); assert.deepEqual(findUndeliveredResults(options).sort(), ["mine", "orphan"]); // Without the takeover hook the old behaviour holds: foreign results are skipped. assert.deepEqual(findUndeliveredResults({ ...options, canTakeOver: undefined }).sort(), ["mine"]); }); }); describe("scheduler release blockers", () => { it("blocks on session-only and unreadable armed schedules and on pending scheduled runs, but not on project schedules", () => { const next = Date.parse("2030-01-01T00:10:00Z"); const blockers = schedulerReleaseBlockers({ armedSchedules: () => [ { id: "nightly", nextRunAt: next, sessionOnly: true }, { id: "project", nextRunAt: next, sessionOnly: false }, { id: "unreadable" }, ], observedCompletionRunIds: () => ["run-9"], }); assert.deepEqual(blockers, [ { source: "scheduler", kind: "armed_schedule", id: "nightly", detail: "Armed session-only schedule, next run 2030-01-01T00:10:00.000Z", since: next }, { source: "scheduler", kind: "armed_schedule", id: "unreadable", detail: "Armed unreadable schedule" }, { source: "scheduler", kind: "scheduled_run_pending", id: "run-9", detail: "Scheduled run run-9 is awaiting completion" }, ]); }); it("does not block on a project schedule alone, yet still blocks on its pending run (live work)", () => { assert.deepEqual(schedulerReleaseBlockers({ armedSchedules: () => [{ id: "p", sessionOnly: false }], observedCompletionRunIds: () => [] }), []); assert.deepEqual( schedulerReleaseBlockers({ armedSchedules: () => [{ id: "p", sessionOnly: false }], observedCompletionRunIds: () => ["run-1"] }).map((b) => b.kind), ["scheduled_run_pending"], ); }); it("fails closed when a real manager cannot read an armed schedule's record", async () => { const root = tempRoot(); const project = path.join(root, "project"); fs.mkdirSync(project, { recursive: true }); const ctx = { cwd: project, sessionManager: { getSessionId: () => "session-a", getSessionFile: () => path.join(project, "session-a.jsonl") } } as unknown as ExtensionContext; const manager = createScheduledRunManager({ config: { scheduledRuns: { enabled: true } }, storeRoot: path.join(root, "stores"), now: () => Date.parse("2030-01-01T00:00:00Z"), timers: { setTimeout: (() => 1) as never, clearTimeout: (() => undefined) as never }, launch: () => new Promise(() => undefined) as never, }); manager.bindSession(ctx); const script = "return runs.run('main', { agent: 'reviewer', task: 'Review the diff' })"; await manager.handleToolCall({ action: "schedule.create", id: "project-one", at: "+10m", workflowScript: script }, ctx); await manager.handleToolCall({ action: "schedule.create", id: "session-one", at: "+20m", sessionOnly: true, workflowScript: script }, ctx); const blockers = (armed = manager.armedSchedules()) => schedulerReleaseBlockers({ armedSchedules: () => armed, observedCompletionRunIds: () => [] }); assert.deepEqual(blockers().map((b) => b.id), ["session-one"]); // Break the stores so `find` throws: both armed timers now have an unknown kind and must block. const stores = (manager as unknown as { stores: Map }).stores; for (const store of stores.values()) store.find = () => { throw new Error("corrupt record"); }; assert.deepEqual(blockers().map((b) => b.id).sort(), ["project-one", "session-one"]); manager.stop(); }); it("reads armed schedules from a real ScheduledRunManager and clears them when paused or stopped", async () => { const root = tempRoot(); const project = path.join(root, "project"); fs.mkdirSync(project, { recursive: true }); const timers = new Map void>(); let nextTimer = 0; const clock = { now: Date.parse("2030-01-01T00:00:00Z") }; const ctx = { cwd: project, sessionManager: { getSessionId: () => "session-a", getSessionFile: () => path.join(project, "session-a.jsonl") } } as unknown as ExtensionContext; const manager = createScheduledRunManager({ config: { scheduledRuns: { enabled: true } }, storeRoot: path.join(root, "stores"), now: () => clock.now, timers: { setTimeout: ((callback: () => void) => { const id = ++nextTimer; timers.set(id, callback); return id; }) as never, clearTimeout: ((id: number) => void timers.delete(id)) as never, }, launch: () => new Promise(() => undefined) as never, }); manager.bindSession(ctx); assert.deepEqual(manager.armedSchedules(), []); const script = "return runs.run('main', { agent: 'reviewer', task: 'Review the diff' })"; await manager.handleToolCall({ action: "schedule.create", id: "project-one", at: "+10m", workflowScript: script }, ctx); await manager.handleToolCall({ action: "schedule.create", id: "session-one", at: "+20m", sessionOnly: true, workflowScript: script }, ctx); const armed = manager.armedSchedules().sort((a, b) => a.id.localeCompare(b.id)); assert.deepEqual(armed, [ { id: "project-one", nextRunAt: clock.now + 600_000, sessionOnly: false }, { id: "session-one", nextRunAt: clock.now + 1_200_000, sessionOnly: true }, ]); await manager.handleToolCall({ action: "schedule.pause", id: "project-one" }, ctx); assert.deepEqual(manager.armedSchedules().map((s) => s.id), ["session-one"]); manager.stop(); assert.deepEqual(manager.armedSchedules(), []); }); }); describe("registerReleaseReadiness", () => { function wire(state: SubagentState, hostSessionId: () => string | null, overrides: Partial = {}, armed: Array<{ id: string }> = []) { const bus = createEventBus(); const registration = registerReleaseReadiness({ events: bus, getHostSessionId: hostSessionId, subagents: deps(state, overrides), scheduler: { armedSchedules: () => armed, observedCompletionRunIds: () => [] }, }); return { bus, registration }; } const scope = (sessionId: string) => ({ version: 1 as const, sessionId, epoch: "e1" }); const command = { type: "release_readiness" as const, version: 1 as const }; const read = (bridge: RpcReleaseBridge) => { const reply = bridge.request(command); if (!reply.success || reply.command !== "release_readiness") throw new Error(JSON.stringify(reply)); return reply.data; }; it("uses the same event name as the host package", () => { assert.equal(RELEASE_READINESS_EVENT, "release:readiness:v1"); }); it("answers as subagents and scheduler for the host session (supervisorOwnerSessionId, not currentSessionId)", () => { const state = makeState([{ asyncId: "run-live" }]); const { bus } = wire(state, () => state.supervisorOwnerSessionId ?? null, { reconcile: () => ({ status: null, repaired: false }) }, [{ id: "armed-1" }]); const result = read(new RpcReleaseBridge(bus, scope("host"))); assert.deepEqual(result.contributors, ["subagents", "scheduler"]); assert.equal(result.safe, false); assert.deepEqual(result.blockers.map((b) => `${b.source}:${b.kind}:${b.id}`), ["subagents:async_run:run-live", "scheduler:armed_schedule:armed-1"]); // The runtime session id (a different value) must NOT be used as the identity. assert.deepEqual(read(new RpcReleaseBridge(bus, scope("/managed/session"))).contributors, []); }); it("is safe but still lists contributors when idle", () => { const state = makeState(); const { bus } = wire(state, () => "host"); assert.deepEqual(read(new RpcReleaseBridge(bus, scope("host"))), { ...scope("host"), safe: true, blockers: [], contributors: ["subagents", "scheduler"] }); }); it("ignores requests for other sessions or before a host session is bound", () => { const state = makeState([{ asyncId: "run-live" }]); let host: string | null = null; const { bus } = wire(state, () => host); assert.deepEqual(read(new RpcReleaseBridge(bus, scope("host"))).contributors, []); host = "someone-else"; assert.deepEqual(read(new RpcReleaseBridge(bus, scope("host"))).contributors, []); }); it("turns a throwing contributor into contributor_error through the bridge and stops answering after dispose", () => { const state = makeState(); (state as { asyncJobs: unknown }).asyncJobs = { values() { throw new Error("state exploded"); } }; const { bus, registration } = wire(state, () => "host"); const failed = read(new RpcReleaseBridge(bus, scope("host"))); assert.equal(failed.safe, false); assert.equal(failed.blockers.some((b) => b.kind === "contributor_error" && b.source === "subagents" && /state exploded/.test(b.detail)), true); registration.dispose(); assert.deepEqual(read(new RpcReleaseBridge(bus, scope("host"))).contributors, []); }); }); describe("stop_all_background: subagents handler", () => { const SESSION = "/managed/session"; const farDeadline = () => Date.now() + 5_000; /** A run directory with a live-looking status (this process's pid is alive), like a detached runner would leave. */ function writeRun(root: string, runId: string, overrides: Record = {}): string { const asyncDir = path.join(root, "async", runId); fs.mkdirSync(asyncDir, { recursive: true }); const now = Date.now(); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify({ runId, sessionId: SESSION, mode: "single", state: "running", pid: process.pid, startedAt: now, lastUpdate: now, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: now }], ...overrides, })); return asyncDir; } /** What the detached runner does when it consumes the stop request: ends with status stopped. */ function runnerStops(asyncDir: string): void { const file = path.join(asyncDir, "status.json"); const status = JSON.parse(fs.readFileSync(file, "utf-8")); fs.writeFileSync(file, JSON.stringify({ ...status, state: "stopped", endedAt: Date.now(), lastUpdate: Date.now(), steps: status.steps.map((step: object) => ({ ...step, status: "stopped" })) })); } function stopDeps(root: string, state: SubagentState, overrides: Partial = {}): SubagentStopDeps & { stopped: string[] } { const stopped: string[] = []; return { state, resultsDir: path.join(root, "results"), asyncDirRoot: path.join(root, "async"), observedCompletionRunIds: () => [], stopRun: ({ dir }) => { stopped.push(path.basename(dir)); runnerStops(dir); }, sleep: () => new Promise((resolve) => setTimeout(resolve, 1)), pollIntervalMs: 1, ...overrides, stopped, }; } const job = (root: string, runId: string, overrides: Partial = {}) => ({ asyncId: runId, asyncDir: path.join(root, "async", runId), sessionId: SESSION, pid: process.pid, ...overrides }); it("uses the same event name as the host package", () => { assert.equal(STOP_ALL_BACKGROUND_EVENT, "release:stop-all:v1"); }); it("is empty (and calls nothing) for an idle session", async () => { const root = tempRoot(); const d = stopDeps(root, makeState()); assert.deepEqual(await stopSubagentWork(d, { deadline: farDeadline() }), []); assert.deepEqual(d.stopped, []); }); it("stops every live async run of the session through stopRun and reports each once its runner is terminal", async () => { const root = tempRoot(); const one = writeRun(root, "run-1"); const two = writeRun(root, "run-2", { mode: "chain" }); const state = makeState([job(root, "run-1", { agents: ["worker"] }), job(root, "run-2", { mode: "chain" })]); const d = stopDeps(root, state); const items = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(d.stopped.sort(), ["run-1", "run-2"]); assert.deepEqual(items.map((i) => [i.source, i.kind, i.id, i.ok]), [["subagents", "async_run", "run-1", true], ["subagents", "async_run", "run-2", true]]); assert.match(items[0]!.detail, /^Stopped Async run \(worker\) run-1/); assert.equal(JSON.parse(fs.readFileSync(path.join(one, "status.json"), "utf-8")).state, "stopped"); assert.equal(JSON.parse(fs.readFileSync(path.join(two, "status.json"), "utf-8")).state, "stopped"); }); it("delivers every stop request synchronously, before the first await (the process may exit right after)", () => { const root = tempRoot(); writeRun(root, "run-1"); const state = makeState([job(root, "run-1")]); const d = stopDeps(root, state, { stopRun: ({ dir }) => { d.stopped.push(path.basename(dir)); } }); const pending = stopSubagentWork(d, { deadline: Date.now() + 30 }); assert.deepEqual(d.stopped, ["run-1"], "the request was written during the synchronous call"); return pending.then((items) => assert.deepEqual(items.map((i) => [i.id, i.ok, i.error]), [["run-1", false, "timeout"]])); }); it("leaves a run of ANOTHER session untouched and says so", async () => { const root = tempRoot(); writeRun(root, "mine"); writeRun(root, "theirs", { sessionId: "/other/session" }); const state = makeState([job(root, "mine"), job(root, "theirs", { sessionId: "/other/session" })]); const d = stopDeps(root, state); const items = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(d.stopped, ["mine"]); assert.equal(JSON.parse(fs.readFileSync(path.join(root, "async", "theirs", "status.json"), "utf-8")).state, "running"); assert.deepEqual(items.map((i) => [i.id, i.ok, i.error]), [["theirs", false, "foreign_session"], ["mine", true, undefined]]); }); it("does not stop a run whose stopRun refuses (foreign session on disk) and reports the refusal as ok:false", async () => { const root = tempRoot(); writeRun(root, "run-1"); const state = makeState([job(root, "run-1")]); const d = stopDeps(root, state, { stopRun: () => { throw new Error("Async run 'run-1' was not found in the active session."); } }); const items = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(items.map((i) => [i.kind, i.id, i.ok, i.error]), [["async_run", "run-1", false, "Async run 'run-1' was not found in the active session."]]); }); it("skips runs that already ended, including a dead runner (repaired to a failed result) and is idempotent", async () => { const root = tempRoot(); writeRun(root, "run-dead", { pid: 999_999 }); writeRun(root, "run-done", { state: "complete" }); writeRun(root, "run-live"); const state = makeState([job(root, "run-dead", { pid: 999_999 }), job(root, "run-done"), job(root, "run-live")]); (state.asyncJobs.get("run-done") as AsyncJobState).status = "running"; // in-memory state lags the disk const d = stopDeps(root, state, { reconcile: (dir, options) => reconcileAsyncRun(dir, { ...options, kill: (pid) => { if (pid === 999_999) dead(); return true; } }) }); const first = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(d.stopped, ["run-live"]); assert.deepEqual(first.map((i) => [i.id, i.ok]), [["run-live", true]]); assert.equal(JSON.parse(fs.readFileSync(path.join(root, "async", "run-dead", "status.json"), "utf-8")).state, "failed"); const second = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(second, []); assert.deepEqual(d.stopped, ["run-live"]); }); it("stops runs a scheduled run is awaiting (scheduled_run_pending) even when they are not tracked in memory, and ignores finished ones", async () => { const root = tempRoot(); writeRun(root, "sched-live"); writeRun(root, "sched-done", { state: "complete" }); const d = stopDeps(root, makeState(), { observedCompletionRunIds: () => ["sched-live", "sched-done", "sched-gone"] }); const items = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(d.stopped, ["sched-live"]); assert.deepEqual(items.map((i) => [i.kind, i.id, i.ok]), [["scheduled_run", "sched-live", true]]); }); it("stops live nested descendants that have their own run directory", async () => { const root = tempRoot(); writeRun(root, "nested-1"); const state = makeState([{ ...job(root, "run-parent"), status: "complete", nestedChildren: [{ id: "nested-1", asyncDir: path.join(root, "async", "nested-1"), state: "running", children: [] }, { id: "nested-done", asyncDir: path.join(root, "async", "nested-done"), state: "complete" }] as never }]); const d = stopDeps(root, state); const items = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(d.stopped, ["nested-1"]); assert.deepEqual(items.map((i) => [i.kind, i.id, i.ok]), [["nested_run", "nested-1", true]]); }); it("reports a runner that never acknowledges as ok:false error timeout, bounded by the deadline", async () => { const root = tempRoot(); writeRun(root, "run-wedged"); const state = makeState([job(root, "run-wedged")]); const d = stopDeps(root, state, { stopRun: ({ dir }) => { d.stopped.push(path.basename(dir)); } }); const startedAt = Date.now(); const items = await stopSubagentWork(d, { deadline: Date.now() + 80 }); assert.ok(Date.now() - startedAt < 2_000); assert.deepEqual(items.map((i) => [i.id, i.ok, i.error]), [["run-wedged", false, "timeout"]]); assert.match(items[0]!.detail, /runner has not acknowledged/); }); it("one failing run does not stop the others", async () => { const root = tempRoot(); writeRun(root, "bad"); writeRun(root, "good"); const state = makeState([job(root, "bad"), job(root, "good")]); const d = stopDeps(root, state, { stopRun: ({ dir }) => { if (path.basename(dir) === "bad") throw new Error("disk full"); d.stopped.push(path.basename(dir)); runnerStops(dir); } }); const items = await stopSubagentWork(d, { deadline: farDeadline() }); assert.deepEqual(items.map((i) => [i.id, i.ok]), [["bad", false], ["good", true]]); }); it("waits for foreground runs to end with the aborted turn, and reports one that does not", async () => { const root = tempRoot(); const state = makeState(); state.foregroundControls.set("fg-ends", { runId: "fg-ends", mode: "single", startedAt: 1, updatedAt: 1, currentAgent: "scout" } as never); state.foregroundControls.set("fg-stuck", { runId: "fg-stuck", mode: "chain", startedAt: 1, updatedAt: 1 } as never); setTimeout(() => state.foregroundControls.delete("fg-ends"), 10); const items = await stopSubagentWork(stopDeps(root, state), { deadline: Date.now() + 120 }); assert.deepEqual(items.map((i) => [i.kind, i.id, i.ok, i.error]), [["foreground_run", "fg-ends", true, undefined], ["foreground_run", "fg-stuck", false, "timeout"]]); }); }); describe("stop_all_background: session-only schedules", () => { const script = "return runs.run('main', { agent: 'reviewer', task: 'Review the diff' })"; function setup() { const root = tempRoot(); const project = path.join(root, "project"); fs.mkdirSync(project, { recursive: true }); const timers = new Map void>(); let nextTimer = 0; const clock = { now: Date.parse("2030-01-01T00:00:00Z") }; const ctx = { cwd: project, sessionManager: { getSessionId: () => "session-a", getSessionFile: () => path.join(project, "session-a.jsonl") } } as unknown as ExtensionContext; const make = () => createScheduledRunManager({ config: { scheduledRuns: { enabled: true } }, storeRoot: path.join(root, "stores"), now: () => clock.now, timers: { setTimeout: ((callback: () => void) => { const id = ++nextTimer; timers.set(id, callback); return id; }) as never, clearTimeout: ((id: number) => void timers.delete(id)) as never, }, launch: () => new Promise(() => undefined) as never, }); const manager = make(); manager.bindSession(ctx); return { root, ctx, manager, make, timers, clock }; } it("pauses session-only schedules (persisted) and clears their timers; project schedules are untouched", async () => { const { ctx, manager, timers } = setup(); await manager.handleToolCall({ action: "schedule.create", id: "project-one", at: "+10m", workflowScript: script }, ctx); await manager.handleToolCall({ action: "schedule.create", id: "session-one", at: "+20m", sessionOnly: true, workflowScript: script }, ctx); await manager.handleToolCall({ action: "schedule.create", id: "session-two", every: "1h", sessionOnly: true, workflowScript: script }, ctx); assert.equal(timers.size, 3); const results = manager.disarmSessionOnlySchedules().sort((a, b) => a.id.localeCompare(b.id)); assert.deepEqual(results.map((r) => [r.id, r.ok]), [["session-one", true], ["session-two", true]]); assert.deepEqual(manager.armedSchedules().map((s) => s.id), ["project-one"]); assert.equal(timers.size, 1); // Persisted: the records are paused, the project schedule is not. const show = async (id: string) => (await manager.handleToolCall({ action: "schedule.show", id }, ctx)).content[0]!.text; assert.match(await show("session-one"), /State: paused/); assert.match(await show("session-two"), /State: paused/); assert.match(await show("project-one"), /State: scheduled/); // Readiness has no armed_schedule blocker for them (and none for the project schedule either, per D11). assert.deepEqual(schedulerReleaseBlockers({ armedSchedules: () => manager.armedSchedules(), observedCompletionRunIds: () => [] }), []); // Idempotent. assert.deepEqual(manager.disarmSessionOnlySchedules(), []); manager.stop(); }); it("does not fire or re-arm on the next launch, even when the planned time has passed (catchUp latest)", async () => { const { ctx, manager, make, timers, clock } = setup(); await manager.handleToolCall({ action: "schedule.create", id: "session-one", at: "+5m", sessionOnly: true, workflowScript: script }, ctx); manager.disarmSessionOnlySchedules(); manager.stop(); clock.now += 3_600_000; // the process was gone for an hour: the schedule is overdue const next = make(); next.bindSession(ctx); assert.deepEqual(next.armedSchedules(), []); assert.equal(timers.size, 0); const due = await next.handleToolCall({ action: "schedule.run-due" }, ctx); assert.match(due.content[0]!.text, /No schedules are due/); // Reversible: resuming re-arms it. await next.handleToolCall({ action: "schedule.resume", id: "session-one" }, ctx); assert.deepEqual(next.armedSchedules().map((s) => s.id), ["session-one"]); next.stop(); }); it("leaves an armed timer alone when its record cannot be read (it may be a project schedule) and reports it", async () => { const { ctx, manager, timers } = setup(); await manager.handleToolCall({ action: "schedule.create", id: "who-knows", at: "+10m", sessionOnly: true, workflowScript: script }, ctx); const stores = (manager as unknown as { stores: Map }).stores; for (const store of stores.values()) store.find = () => { throw new Error("corrupt record"); }; const results = manager.disarmSessionOnlySchedules(); assert.deepEqual(results.map((r) => [r.id, r.ok, r.error]), [["who-knows", false, "corrupt record"]]); assert.equal(timers.size, 1); manager.stop(); }); it("reports ok:false when the pause cannot be persisted, but still clears the in-memory timer", async () => { const { ctx, manager, timers } = setup(); await manager.handleToolCall({ action: "schedule.create", id: "session-one", at: "+10m", sessionOnly: true, workflowScript: script }, ctx); const stores = (manager as unknown as { stores: Map }).stores; for (const store of stores.values()) store.write = () => { throw new Error("read-only file system"); }; const results = manager.disarmSessionOnlySchedules(); assert.deepEqual(results.map((r) => [r.id, r.ok, r.error]), [["session-one", false, "read-only file system"]]); assert.match(results[0]!.detail, /may fire on the next launch/); assert.equal(timers.size, 0); manager.stop(); }); it("maps results to scheduler stop items", () => { assert.deepEqual(stopSchedules({ disarmSessionSchedules: () => [{ id: "a", ok: true, detail: "Paused" }, { id: "b", ok: false, detail: "nope", error: "boom" }] }), [ { source: "scheduler", kind: "schedule", id: "a", detail: "Paused", ok: true }, { source: "scheduler", kind: "schedule", id: "b", detail: "nope", ok: false, error: "boom" }, ]); }); }); describe("stop_all_background: wired through the RPC bridge", () => { const script = "return runs.run('main', { agent: 'reviewer', task: 'Review the diff' })"; const stopCommand = { type: "stop_all_background" as const, version: 1 as const }; const scope = { version: 1 as const, sessionId: "host", epoch: "e1" }; it("stops this session's runs and session-only schedules in one command, leaving readiness without async_run or armed_schedule blockers", async () => { const root = tempRoot(); const project = path.join(root, "project"); fs.mkdirSync(project, { recursive: true }); const ctx = { cwd: project, sessionManager: { getSessionId: () => "session-a", getSessionFile: () => path.join(project, "session-a.jsonl") } } as unknown as ExtensionContext; const manager = createScheduledRunManager({ config: { scheduledRuns: { enabled: true } }, storeRoot: path.join(root, "stores"), now: () => Date.parse("2030-01-01T00:00:00Z"), timers: { setTimeout: (() => 1) as never, clearTimeout: (() => undefined) as never }, launch: () => new Promise(() => undefined) as never, }); manager.bindSession(ctx); await manager.handleToolCall({ action: "schedule.create", id: "project-one", at: "+10m", workflowScript: script }, ctx); await manager.handleToolCall({ action: "schedule.create", id: "session-one", at: "+20m", sessionOnly: true, workflowScript: script }, ctx); const asyncRoot = path.join(root, "async"); const run = (id: string, sessionId: string) => { const asyncDir = path.join(asyncRoot, id); fs.mkdirSync(asyncDir, { recursive: true }); const now = Date.now(); fs.writeFileSync(path.join(asyncDir, "status.json"), JSON.stringify({ runId: id, sessionId, mode: "single", state: "running", pid: process.pid, startedAt: now, lastUpdate: now, currentStep: 0, steps: [{ agent: "worker", status: "running", startedAt: now }] })); return asyncDir; }; const mine = run("mine", "/managed/session"); run("theirs", "/other/session"); // on disk, not tracked by this session const state = makeState([{ asyncId: "mine", asyncDir: mine, sessionId: "/managed/session", pid: process.pid }]); const bus = createEventBus(); const reconcile: SubagentReleaseDeps["reconcile"] = (dir, options) => reconcileAsyncRun(dir, options); const stoppedDirs: string[] = []; registerReleaseReadiness({ events: bus, getHostSessionId: () => "host", subagents: deps(state, { resultsDir: path.join(root, "results"), reconcile }), scheduler: { armedSchedules: () => manager.armedSchedules(), observedCompletionRunIds: () => manager.observedCompletionRunIds() }, stopAll: { subagents: { state, resultsDir: path.join(root, "results"), asyncDirRoot: asyncRoot, observedCompletionRunIds: () => manager.observedCompletionRunIds(), stopRun: ({ dir }) => { stoppedDirs.push(path.basename(dir)); const file = path.join(dir, "status.json"); setTimeout(() => fs.writeFileSync(file, JSON.stringify({ ...JSON.parse(fs.readFileSync(file, "utf-8")), state: "stopped", endedAt: Date.now() })), 5); }, pollIntervalMs: 5, }, scheduler: { disarmSessionSchedules: () => manager.disarmSessionOnlySchedules() }, }, }); const bridge = new RpcReleaseBridge(bus, scope); const before = bridge.request({ type: "release_readiness", version: 1 }); assert.ok(before.success && before.command === "release_readiness"); assert.deepEqual(before.data.blockers.map((b) => `${b.kind}:${b.id}`), ["async_run:mine", "armed_schedule:session-one"]); const reply = await bridge.stopAll(stopCommand, () => [], () => []); assert.ok(reply.success && reply.command === "stop_all_background"); assert.deepEqual(reply.data.stopped.map((i) => `${i.source}:${i.kind}:${i.id}:${i.ok}`), ["scheduler:schedule:session-one:true", "subagents:async_run:mine:true"]); assert.deepEqual(stoppedDirs, ["mine"], "the run of another session is never touched"); assert.equal(JSON.parse(fs.readFileSync(path.join(asyncRoot, "theirs", "status.json"), "utf-8")).state, "running"); assert.deepEqual(reply.data.remaining.blockers.filter((b) => b.kind === "async_run" || b.kind === "armed_schedule"), []); assert.deepEqual(manager.armedSchedules().map((s) => s.id), ["project-one"], "the project schedule stays armed"); const again = await bridge.stopAll(stopCommand, () => [], () => []); assert.ok(again.success && again.command === "stop_all_background"); assert.deepEqual(again.data.stopped, []); manager.stop(); }); it("ignores stop requests for another session and after dispose", async () => { const bus = createEventBus(); let calls = 0; const registration = registerReleaseReadiness({ events: bus, getHostSessionId: () => "host", subagents: deps(makeState()), scheduler: { armedSchedules: () => [], observedCompletionRunIds: () => [] }, stopAll: { subagents: { state: makeState(), observedCompletionRunIds: () => { calls++; return []; }, stopRun: () => undefined }, scheduler: { disarmSessionSchedules: () => { calls++; return []; } }, }, }); await new RpcReleaseBridge(bus, { ...scope, sessionId: "someone-else" }).stopAll(stopCommand, () => [], () => []); assert.equal(calls, 0); registration.dispose(); await new RpcReleaseBridge(bus, scope).stopAll(stopCommand, () => [], () => []); assert.equal(calls, 0); }); });