diff --git a/docs/configuration.md b/docs/configuration.md index 3502a98b..c44da6d3 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -153,6 +153,10 @@ lists existing subagent sessions for the current workspace, scoped by the workspace environment injected into shell commands. The `subagent-delegation` skill teaches the model to use only the minimal `devspace agents ls`, `devspace agents run`, and `devspace agents show` workflow. +`devspace agents run` submits work to the running DevSpace server so subagent +execution can reuse long-lived provider resources. Start `devspace serve` with +`DEVSPACE_SUBAGENTS=1` before using the agent commands directly from another +terminal. Starter profile templates are available under `examples/agents/`. Copy or adapt them into one of the active profile directories before use. diff --git a/package.json b/package.json index 12983912..ebebd3be 100644 --- a/package.json +++ b/package.json @@ -28,7 +28,7 @@ "dev": "node scripts/dev-server.mjs", "postinstall": "node scripts/fix-node-pty-permissions.mjs", "start": "node dist/cli.js serve", - "test": "tsx src/config.test.ts && tsx src/request-meta.test.ts && tsx src/incoming-artifacts.test.ts && tsx src/artifact-download.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/workspace-conversation.test.ts && tsx src/review-checkpoints.test.ts && tsx src/server.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts", + "test": "tsx src/config.test.ts && tsx src/request-meta.test.ts && tsx src/incoming-artifacts.test.ts && tsx src/artifact-download.test.ts && tsx src/ui/card-types.test.ts && tsx src/ui/patch-display.test.ts && tsx src/ui/tool-display.test.ts && tsx src/apply-patch.test.ts && tsx src/process-platform.test.ts && tsx src/process-sessions.test.ts && tsx src/mcp-sessions.test.ts && tsx src/server-shutdown.test.ts && tsx src/local-agent-runtime.test.ts && tsx src/local-agent-adapters.test.ts && tsx src/local-agent-availability.test.ts && tsx src/local-agent-profiles.test.ts && tsx src/local-agent-targets.test.ts && tsx src/local-agent-store.test.ts && tsx src/local-agent-manager.test.ts && tsx src/local-agent-control.test.ts && tsx src/roots.test.ts && tsx src/skills.test.ts && tsx src/workspaces.test.ts && tsx src/workspace-conversation.test.ts && tsx src/review-checkpoints.test.ts && tsx src/server.test.ts && tsx src/oauth-store.test.ts && tsx src/cli.test.ts", "typecheck": "tsc -p tsconfig.json --noEmit" }, "keywords": [], diff --git a/src/cli.ts b/src/cli.ts index 7a1ac63f..15af1042 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -1,33 +1,19 @@ #!/usr/bin/env node import { createRequire } from "node:module"; import { stdin as input, stdout as output } from "node:process"; -import { spawn } from "node:child_process"; -import { mkdtempSync, writeFileSync } from "node:fs"; -import { readFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import { join, resolve } from "node:path"; -import { fileURLToPath } from "node:url"; +import { resolve } from "node:path"; import * as prompts from "@clack/prompts"; import { getShellConfig } from "@earendil-works/pi-coding-agent"; import { satisfies } from "semver"; import { loadConfig } from "./config.js"; -import { runLocalAgentProvider } from "./local-agent-adapters.js"; import { - isLocalAgentProvider, - loadLocalAgentProfiles, - type LocalAgentProfile, -} from "./local-agent-profiles.js"; -import { - assertLocalAgentProviderAvailable, formatLocalAgentProviderAvailabilitySummary, } from "./local-agent-availability.js"; import { - formatAvailableLocalAgentTargets, parseLocalAgentRunArgs, - resolveLocalAgentTarget, } from "./local-agent-targets.js"; import { createLocalAgentStore, type LocalAgentRecord } from "./local-agent-store.js"; -import type { LocalAgentRunResult } from "./local-agent-runtime.js"; +import { requestLocalAgentRun } from "./local-agent-control.js"; import { ensureDevspaceDefaultSkills, generateOwnerToken, @@ -213,7 +199,8 @@ async function serve(): Promise { const { createServer } = await import("./server.js"); const config = loadConfig(); - const { app, close, localAgentProviders } = createServer(config); + const { app, close, localAgentProviders, startAgentControl } = createServer(config); + await startAgentControl(); const httpServer = app.listen(config.port, config.host, () => { console.log(`devspace listening on http://${config.host}:${config.port}/mcp`); console.log(`public base url: ${config.publicBaseUrl}`); @@ -333,9 +320,6 @@ async function runAgentsCommand(args: string[]): Promise { case "show": await runAgentsShow(rest); return; - case "__worker": - await runAgentsWorker(rest); - return; case undefined: case "help": case "--help": @@ -364,56 +348,16 @@ async function runAgentsList(): Promise { async function runAgentsRun(args: string[]): Promise { const parsed = parseLocalAgentRunArgs(args); - const config = loadConfig(); - const workspaceRoot = resolveCurrentWorkspaceRoot(); - const store = createLocalAgentStore(config); - const existing = store.get(parsed.target); - - if (existing) { - if (!isLocalAgentProvider(existing.provider)) { - throw new Error(`Unknown subagent provider for existing session: ${existing.provider}`); - } - assertLocalAgentProviderAvailable(existing.provider); - const promptFile = writeAgentPromptFile(parsed.prompt); - store.update(existing.id, { - status: "starting", - model: parsed.model ?? existing.model, - thinking: parsed.thinking ?? existing.thinking, - latestResponse: undefined, - error: undefined, - }); - spawnAgentWorker(existing.id, promptFile); - console.log(formatAgentLine({ - ...existing, - status: "running", - model: parsed.model ?? existing.model, - thinking: parsed.thinking ?? existing.thinking, - })); - return; - } - - const profiles = await loadLocalAgentProfiles(config, workspaceRoot); - const target = resolveLocalAgentTarget(parsed.target, profiles, parsed.model, parsed.thinking); - if (!target) { - throw new Error( - `Unknown subagent profile, provider, or id: ${parsed.target}. Available ${formatAvailableLocalAgentTargets(profiles)}`, - ); - } - assertLocalAgentProviderAvailable(target.provider); - - const promptFile = writeAgentPromptFile(parsed.prompt); - const record = store.create({ + const record = await requestLocalAgentRun(config, { workspaceId: process.env.DEVSPACE_WORKSPACE_ID, - workspaceRoot, - profileName: target.name, - provider: target.provider, - model: target.model, - thinking: target.thinking, + workspaceRoot: resolveCurrentWorkspaceRoot(), + target: parsed.target, + prompt: parsed.prompt, + model: parsed.model, + thinking: parsed.thinking, }); - - spawnAgentWorker(record.id, promptFile); - console.log(formatAgentLine({ ...record, status: "running" })); + console.log(formatAgentLine(record)); } async function runAgentsShow(args: string[]): Promise { @@ -445,98 +389,6 @@ async function runAgentsShow(args: string[]): Promise { } } -async function runAgentsWorker(args: string[]): Promise { - const [id, promptFileFlag, promptFile] = args; - if (!id || promptFileFlag !== "--prompt-file" || !promptFile) { - throw new Error("Usage: devspace agents __worker --prompt-file "); - } - - const config = loadConfig(); - const store = createLocalAgentStore(config); - const record = store.get(id); - if (!record) throw new Error(`Unknown subagent id: ${id}`); - - store.update(record.id, { status: "running", error: undefined }); - try { - const profiles = await loadLocalAgentProfiles(config, record.workspaceRoot); - const profile = profiles.find((candidate) => candidate.name === record.profileName); - const prompt = await readFile(promptFile, "utf8"); - const result = profile - ? await runLocalAgentProfile(profile, record, prompt) - : await runRawLocalAgentProvider(record, prompt); - store.update(record.id, { - providerSessionId: result.providerSessionId ?? undefined, - status: "idle", - latestResponse: result.finalResponse, - error: undefined, - }); - } catch (error) { - store.update(record.id, { - status: "error", - error: error instanceof Error ? error.message : String(error), - }); - } -} - -async function runLocalAgentProfile( - profile: LocalAgentProfile, - record: LocalAgentRecord, - prompt: string, -): Promise { - const body = profile.body.trim(); - const fullPrompt = body ? `${body}\n\nTask:\n${prompt}` : prompt; - return runLocalAgentProvider(profile.provider, { - prompt: fullPrompt, - workspace: record.workspaceRoot, - providerSessionId: record.providerSessionId, - writeMode: "allowed", - model: record.model ?? profile.model, - thinking: record.thinking ?? profile.thinking, - }); -} - -async function runRawLocalAgentProvider( - record: LocalAgentRecord, - prompt: string, -): Promise { - if (record.profileName !== record.provider || !isLocalAgentProvider(record.provider)) { - throw new Error(`Subagent profile not found: ${record.profileName}`); - } - - return runLocalAgentProvider(record.provider, { - prompt, - workspace: record.workspaceRoot, - providerSessionId: record.providerSessionId, - writeMode: "allowed", - model: record.model, - thinking: record.thinking, - }); -} - -function spawnAgentWorker(agentId: string, promptFile: string): void { - const child = spawn(process.execPath, [ - ...process.execArgv, - fileURLToPath(import.meta.url), - "agents", - "__worker", - agentId, - "--prompt-file", - promptFile, - ], { - detached: true, - stdio: "ignore", - env: process.env, - }); - child.unref(); -} - -function writeAgentPromptFile(prompt: string): string { - const directory = mkdtempSync(join(tmpdir(), "devspace-agent-prompt-")); - const filePath = join(directory, "prompt.txt"); - writeFileSync(filePath, prompt, { mode: 0o600 }); - return filePath; -} - function resolveCurrentWorkspaceRoot(): string { return resolve(process.env.DEVSPACE_WORKSPACE_ROOT || process.cwd()); } diff --git a/src/local-agent-control.test.ts b/src/local-agent-control.test.ts new file mode 100644 index 00000000..26d28e76 --- /dev/null +++ b/src/local-agent-control.test.ts @@ -0,0 +1,82 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { createConnection } from "node:net"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { loadConfig } from "./config.js"; +import { + LocalAgentControlServer, + localAgentControlAddress, + requestLocalAgentRun, + type LocalAgentCommandHandler, +} from "./local-agent-control.js"; +import type { LocalAgentRunCommand } from "./local-agent-manager.js"; +import type { LocalAgentRecord } from "./local-agent-store.js"; + +const root = mkdtempSync(join(tmpdir(), "devspace-local-agent-control-test-")); +const config = loadConfig({ + ...process.env, + DEVSPACE_CONFIG_DIR: join(root, "config"), + DEVSPACE_STATE_DIR: join(root, "state"), + DEVSPACE_WORKTREE_ROOT: join(root, "worktrees"), + DEVSPACE_ALLOWED_ROOTS: root, + DEVSPACE_OAUTH_OWNER_TOKEN: "test-owner-token-long-enough", +}); +let received: LocalAgentRunCommand | undefined; +const expected: LocalAgentRecord = { + id: "agt_control", + workspaceId: "ws_control", + workspaceRoot: root, + profileName: "codex", + provider: "codex", + status: "running", + model: undefined, + thinking: undefined, + providerSessionId: undefined, + latestResponse: undefined, + error: undefined, + createdAt: "2026-01-01T00:00:00.000Z", + updatedAt: "2026-01-01T00:00:00.000Z", +}; +const handler: LocalAgentCommandHandler = { + async enqueue(command) { + received = command; + return expected; + }, +}; +const server = new LocalAgentControlServer(config, handler); + +try { + await server.start(); + const record = await requestLocalAgentRun(config, { + workspaceId: "ws_control", + workspaceRoot: root, + target: "codex", + prompt: "inspect", + model: "gpt-test", + thinking: "high", + }); + + assert.deepEqual(received, { + workspaceId: "ws_control", + workspaceRoot: root, + target: "codex", + prompt: "inspect", + model: "gpt-test", + thinking: "high", + }); + assert.deepEqual(record, expected); + + const idleSocket = createConnection(localAgentControlAddress(config.stateDir)); + await new Promise((resolve, reject) => { + idleSocket.once("connect", resolve); + idleSocket.once("error", reject); + }); + const socketClosed = new Promise((resolve) => idleSocket.once("close", () => resolve())); + await server.close(); + await socketClosed; + assert.equal(idleSocket.destroyed, true); +} finally { + await server.close(); + rmSync(root, { recursive: true, force: true }); +} diff --git a/src/local-agent-control.ts b/src/local-agent-control.ts new file mode 100644 index 00000000..468f4821 --- /dev/null +++ b/src/local-agent-control.ts @@ -0,0 +1,250 @@ +import { createHash } from "node:crypto"; +import { createConnection, createServer, type Server, type Socket } from "node:net"; +import { mkdir, rm } from "node:fs/promises"; +import { join } from "node:path"; +import type { ServerConfig } from "./config.js"; +import { + type LocalAgentRunCommand, +} from "./local-agent-manager.js"; +import type { LocalAgentRecord } from "./local-agent-store.js"; + +const MAX_CONTROL_REQUEST_BYTES = 512 * 1024; +const CONTROL_SOCKET_IDLE_TIMEOUT_MS = 30_000; +const CONTROL_REQUEST_TIMEOUT_MS = 30_000; + +type ControlRequest = { type: "run"; command: LocalAgentRunCommand }; +type ControlResponse = + | { ok: true; record: LocalAgentRecord } + | { ok: false; error: string }; + +export interface LocalAgentCommandHandler { + enqueue(command: LocalAgentRunCommand): Promise; +} + +export class LocalAgentControlServer { + private server: Server | undefined; + private readonly sockets = new Set(); + + constructor( + private readonly config: ServerConfig, + private readonly manager: LocalAgentCommandHandler, + ) {} + + async start(): Promise { + if (this.server) return; + const address = localAgentControlAddress(this.config.stateDir); + if (process.platform !== "win32") { + await mkdir(this.config.stateDir, { recursive: true }); + await removeStaleSocket(address); + } + + const server = createServer((socket) => this.handleConnection(socket)); + await new Promise((resolve, reject) => { + const onError = (error: Error) => { + server.off("listening", onListening); + reject(error); + }; + const onListening = () => { + server.off("error", onError); + resolve(); + }; + server.once("error", onError); + server.once("listening", onListening); + server.listen(address); + }); + this.server = server; + } + + async close(): Promise { + const server = this.server; + this.server = undefined; + for (const socket of this.sockets) socket.destroy(); + this.sockets.clear(); + if (server) { + await new Promise((resolve) => server.close(() => resolve())); + } + if (process.platform !== "win32") { + await rm(localAgentControlAddress(this.config.stateDir), { force: true }); + } + } + + private handleConnection(socket: Socket): void { + this.sockets.add(socket); + socket.once("close", () => this.sockets.delete(socket)); + socket.on("error", () => undefined); + socket.setTimeout(CONTROL_SOCKET_IDLE_TIMEOUT_MS, () => socket.destroy()); + socket.setEncoding("utf8"); + let input = ""; + let handled = false; + socket.on("data", (chunk: string) => { + if (handled) return; + input += chunk; + if (Buffer.byteLength(input, "utf8") > MAX_CONTROL_REQUEST_BYTES) { + handled = true; + socket.end(`${JSON.stringify({ ok: false, error: "DevSpace subagent control request is too large." })}\n`); + return; + } + const newline = input.indexOf("\n"); + if (newline === -1) return; + const line = input.slice(0, newline); + handled = true; + socket.pause(); + void this.handleLine(line).then((response) => { + if (!socket.destroyed) socket.end(`${JSON.stringify(response)}\n`); + }).catch(() => socket.destroy()); + }); + } + + private async handleLine(line: string): Promise { + try { + const request = parseControlRequest(JSON.parse(line) as unknown); + const record = await this.manager.enqueue(request.command); + return { ok: true, record }; + } catch (error) { + return { ok: false, error: error instanceof Error ? error.message : String(error) }; + } + } +} + +export async function requestLocalAgentRun( + config: ServerConfig, + command: LocalAgentRunCommand, +): Promise { + const response = await requestControl(localAgentControlAddress(config.stateDir), { + type: "run", + command, + }); + if (!response.ok) throw new Error(response.error); + return response.record; +} + +export function localAgentControlAddress(stateDir: string): string { + if (process.platform !== "win32") return join(stateDir, "subagents.sock"); + const digest = createHash("sha256").update(stateDir).digest("hex").slice(0, 16); + return `\\\\.\\pipe\\devspace-subagents-${digest}`; +} + +async function requestControl(address: string, request: ControlRequest): Promise { + return new Promise((resolve, reject) => { + const socket = createConnection(address); + socket.setEncoding("utf8"); + let response = ""; + let settled = false; + const fail = (error: Error) => { + if (settled) return; + settled = true; + socket.destroy(); + reject(error); + }; + socket.setTimeout(CONTROL_REQUEST_TIMEOUT_MS, () => { + fail(new Error("DevSpace subagent runtime did not respond in time.")); + }); + socket.once("error", (error) => { + fail(new Error( + `DevSpace subagent runtime is unavailable. Start \`devspace serve\` and try again. (${error.message})`, + )); + }); + socket.on("data", (chunk: string) => { + response += chunk; + }); + socket.on("end", () => { + if (settled) return; + settled = true; + try { + resolve(parseControlResponse(JSON.parse(response.trim()) as unknown)); + } catch (error) { + reject(error); + } + }); + socket.once("connect", () => { + socket.write(`${JSON.stringify(request)}\n`); + }); + }); +} + +async function removeStaleSocket(address: string): Promise { + const active = await canConnect(address); + if (active) { + throw new Error(`DevSpace subagent control socket is already active at ${address}.`); + } + await rm(address, { force: true }); +} + +async function canConnect(address: string): Promise { + return new Promise((resolve) => { + const socket = createConnection(address); + socket.once("connect", () => { + socket.destroy(); + resolve(true); + }); + socket.once("error", () => resolve(false)); + }); +} + +function parseControlRequest(value: unknown): ControlRequest { + if (!isRecord(value) || value.type !== "run" || !isRecord(value.command)) { + throw new Error("Invalid DevSpace subagent control request."); + } + return { type: "run", command: parseRunCommand(value.command) }; +} + +function parseRunCommand(value: Record): LocalAgentRunCommand { + const workspaceRoot = requiredString(value.workspaceRoot, "workspaceRoot"); + const target = requiredString(value.target, "target"); + const prompt = requiredString(value.prompt, "prompt"); + return { + workspaceRoot, + target, + prompt, + workspaceId: optionalString(value.workspaceId), + model: optionalString(value.model), + thinking: optionalString(value.thinking), + }; +} + +function parseControlResponse(value: unknown): ControlResponse { + if (!isRecord(value) || typeof value.ok !== "boolean") { + throw new Error("Invalid response from DevSpace subagent runtime."); + } + if (value.ok === false) { + return { ok: false, error: requiredString(value.error, "error") }; + } + if (!isRecord(value.record)) throw new Error("Invalid subagent record in control response."); + return { ok: true, record: parseAgentRecord(value.record) }; +} + +function parseAgentRecord(value: Record): LocalAgentRecord { + const status = value.status; + if (status !== "starting" && status !== "running" && status !== "idle" && status !== "error" && status !== "stopped") { + throw new Error("Invalid subagent status in control response."); + } + return { + id: requiredString(value.id, "id"), + workspaceRoot: requiredString(value.workspaceRoot, "workspaceRoot"), + profileName: requiredString(value.profileName, "profileName"), + provider: requiredString(value.provider, "provider"), + status, + createdAt: requiredString(value.createdAt, "createdAt"), + updatedAt: requiredString(value.updatedAt, "updatedAt"), + workspaceId: optionalString(value.workspaceId), + model: optionalString(value.model), + thinking: optionalString(value.thinking), + providerSessionId: optionalString(value.providerSessionId), + latestResponse: optionalString(value.latestResponse), + error: optionalString(value.error), + }; +} + +function requiredString(value: unknown, name: string): string { + const result = optionalString(value); + if (!result) throw new Error(`Invalid ${name} in DevSpace subagent control message.`); + return result; +} + +function optionalString(value: unknown): string | undefined { + return typeof value === "string" && value.length > 0 ? value : undefined; +} + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} diff --git a/src/local-agent-manager.test.ts b/src/local-agent-manager.test.ts new file mode 100644 index 00000000..16cd972c --- /dev/null +++ b/src/local-agent-manager.test.ts @@ -0,0 +1,126 @@ +import assert from "node:assert/strict"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { loadConfig } from "./config.js"; +import { LocalAgentManager } from "./local-agent-manager.js"; +import { LocalAgentStore } from "./local-agent-store.js"; +import type { LocalAgentRunInput, LocalAgentRunResult } from "./local-agent-runtime.js"; + +const root = mkdtempSync(join(tmpdir(), "devspace-local-agent-manager-test-")); +const config = loadConfig({ + ...process.env, + DEVSPACE_CONFIG_DIR: join(root, "config"), + DEVSPACE_STATE_DIR: join(root, "state"), + DEVSPACE_WORKTREE_ROOT: join(root, "worktrees"), + DEVSPACE_ALLOWED_ROOTS: root, + DEVSPACE_OAUTH_OWNER_TOKEN: "test-owner-token-long-enough", + DEVSPACE_SUBAGENTS: "1", +}); +const store = new LocalAgentStore(config.stateDir); +const firstRun = deferred(); +const secondRun = deferred(); +const calls: Array<{ provider: string; input: LocalAgentRunInput }> = []; + +const manager = new LocalAgentManager(config, { + store, + assertProviderAvailable: () => undefined, + runProvider: async (provider, input) => { + calls.push({ provider, input }); + return calls.length === 1 ? firstRun.promise : secondRun.promise; + }, +}); + +try { + const first = await manager.enqueue({ + workspaceId: "ws_test", + workspaceRoot: root, + target: "codex", + prompt: "first", + model: "gpt-test", + }); + await waitFor(() => calls.length === 1); + + await assert.rejects( + manager.enqueue({ + workspaceId: "ws_test", + workspaceRoot: join(root, "other-workspace"), + target: first.id, + prompt: "wrong workspace", + }), + /belongs to a different workspace/, + ); + + await assert.rejects( + manager.enqueue({ + workspaceRoot: join(root, "..", "outside-root"), + target: "codex", + prompt: "outside", + }), + /outside allowed roots/, + ); + + const queued = await manager.enqueue({ + workspaceId: "ws_test", + workspaceRoot: root, + target: first.id, + prompt: "second", + thinking: "high", + }); + assert.equal(queued.id, first.id); + await immediate(); + assert.equal(calls.length, 1, "a second turn for the same agent must wait for the first"); + + firstRun.resolve({ + provider: "codex", + providerSessionId: "thread_first", + finalResponse: "first result", + items: [], + }); + await waitFor(() => calls.length === 2); + + assert.equal(calls[0]?.provider, "codex"); + assert.equal(calls[0]?.input.prompt, "first"); + assert.equal(calls[0]?.input.model, "gpt-test"); + assert.equal(calls[1]?.input.prompt, "second"); + assert.equal(calls[1]?.input.providerSessionId, "thread_first"); + assert.equal(calls[1]?.input.thinking, "high"); + + secondRun.resolve({ + provider: "codex", + providerSessionId: "thread_second", + finalResponse: "second result", + items: [], + }); + await waitFor(() => store.get(first.id)?.status === "idle"); + + const completed = store.get(first.id); + assert.equal(completed?.providerSessionId, "thread_second"); + assert.equal(completed?.latestResponse, "second result"); +} finally { + await manager.shutdown(); + rmSync(root, { recursive: true, force: true }); +} + +function deferred(): { + promise: Promise; + resolve(value: T): void; +} { + let resolve!: (value: T) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + +async function waitFor(predicate: () => boolean): Promise { + for (let attempt = 0; attempt < 100; attempt += 1) { + if (predicate()) return; + await immediate(); + } + throw new Error("Timed out waiting for local agent manager state."); +} + +function immediate(): Promise { + return new Promise((resolve) => setImmediate(resolve)); +} diff --git a/src/local-agent-manager.ts b/src/local-agent-manager.ts new file mode 100644 index 00000000..fe1fbb4d --- /dev/null +++ b/src/local-agent-manager.ts @@ -0,0 +1,216 @@ +import type { ServerConfig } from "./config.js"; +import { runLocalAgentProvider } from "./local-agent-adapters.js"; +import { assertLocalAgentProviderAvailable } from "./local-agent-availability.js"; +import { + isLocalAgentProvider, + loadLocalAgentProfiles, + type LocalAgentProfile, +} from "./local-agent-profiles.js"; +import { + formatAvailableLocalAgentTargets, + resolveLocalAgentTarget, +} from "./local-agent-targets.js"; +import { + createLocalAgentStore, + type LocalAgentRecord, + type LocalAgentStore, +} from "./local-agent-store.js"; +import type { LocalAgentRunResult } from "./local-agent-runtime.js"; +import { assertAllowedPath } from "./roots.js"; + +export interface LocalAgentRunCommand { + workspaceId?: string; + workspaceRoot: string; + target: string; + prompt: string; + model?: string; + thinking?: string; +} + +type RunProvider = ( + provider: LocalAgentProfile["provider"], + input: Parameters[1], +) => Promise; + +interface LocalAgentManagerOptions { + store?: LocalAgentStore; + runProvider?: RunProvider; + assertProviderAvailable?: typeof assertLocalAgentProviderAvailable; +} + +interface QueuedTurn { + prompt: string; + model?: string; + thinking?: string; +} + +interface AgentQueue { + tail: Promise; + pending: number; +} + +/** + * Owns logical subagent turn serialization and durable state updates. + * + * Provider runtimes are intentionally not durable state: the store owns the + * provider session id, while live provider resources may be recreated later. + */ +export class LocalAgentManager { + private readonly store: LocalAgentStore; + private readonly runProvider: RunProvider; + private readonly assertProviderAvailable: typeof assertLocalAgentProviderAvailable; + private readonly queues = new Map(); + private closing = false; + + constructor( + private readonly config: ServerConfig, + options: LocalAgentManagerOptions = {}, + ) { + this.store = options.store ?? createLocalAgentStore(config); + this.runProvider = options.runProvider ?? runLocalAgentProvider; + this.assertProviderAvailable = options.assertProviderAvailable ?? assertLocalAgentProviderAvailable; + } + + async enqueue(command: LocalAgentRunCommand): Promise { + if (this.closing) throw new Error("DevSpace subagent manager is shutting down."); + + const normalizedCommand = { + ...command, + workspaceRoot: assertAllowedPath(command.workspaceRoot, this.config.allowedRoots), + }; + const existing = this.store.getByAgentId(normalizedCommand.target); + const record = existing + ? this.prepareExistingAgent(existing, normalizedCommand) + : await this.createAgent(normalizedCommand); + + this.schedule(record.id, { + prompt: normalizedCommand.prompt, + model: normalizedCommand.model ?? record.model, + thinking: normalizedCommand.thinking ?? record.thinking, + }); + + return this.store.update(record.id, { status: "running" }); + } + + async shutdown(): Promise { + if (this.closing) return; + this.closing = true; + await Promise.allSettled(Array.from(this.queues.values(), (queue) => queue.tail)); + this.store.close(); + } + + private prepareExistingAgent( + existing: LocalAgentRecord, + command: LocalAgentRunCommand, + ): LocalAgentRecord { + const workspaceRoot = assertAllowedPath(existing.workspaceRoot, this.config.allowedRoots); + if (workspaceRoot !== command.workspaceRoot) { + throw new Error(`Subagent ${existing.id} belongs to a different workspace.`); + } + if (command.workspaceId && existing.workspaceId && command.workspaceId !== existing.workspaceId) { + throw new Error(`Subagent ${existing.id} belongs to a different workspace.`); + } + if (!isLocalAgentProvider(existing.provider)) { + throw new Error(`Unknown subagent provider for existing session: ${existing.provider}`); + } + this.assertProviderAvailable(existing.provider); + return this.store.update(existing.id, { + status: "starting", + model: command.model ?? existing.model, + thinking: command.thinking ?? existing.thinking, + latestResponse: undefined, + error: undefined, + }); + } + + private async createAgent(command: LocalAgentRunCommand): Promise { + const profiles = await loadLocalAgentProfiles(this.config, command.workspaceRoot); + const target = resolveLocalAgentTarget(command.target, profiles, command.model, command.thinking); + if (!target) { + throw new Error( + `Unknown subagent profile, provider, or id: ${command.target}. Available ${formatAvailableLocalAgentTargets(profiles)}`, + ); + } + this.assertProviderAvailable(target.provider); + return this.store.create({ + workspaceId: command.workspaceId, + workspaceRoot: command.workspaceRoot, + profileName: target.name, + provider: target.provider, + model: target.model, + thinking: target.thinking, + }); + } + + private schedule(agentId: string, turn: QueuedTurn): void { + const current = this.queues.get(agentId); + const previous = current?.tail ?? Promise.resolve(); + const queue: AgentQueue = current ?? { tail: Promise.resolve(), pending: 0 }; + queue.pending += 1; + + const next = previous + .catch(() => undefined) + .then(async () => { + this.store.update(agentId, { status: "running", error: undefined }); + await this.executeTurn(agentId, turn); + }) + .catch((error) => { + try { + this.store.update(agentId, { + status: "error", + error: error instanceof Error ? error.message : String(error), + }); + } catch { + // The store may already be unavailable during shutdown. + } + }) + .finally(() => { + queue.pending -= 1; + if (queue.pending === 0 && this.queues.get(agentId) === queue) { + this.queues.delete(agentId); + } + }); + + queue.tail = next; + this.queues.set(agentId, queue); + } + + private async executeTurn(agentId: string, turn: QueuedTurn): Promise { + const record = this.store.getByAgentId(agentId); + if (!record) throw new Error(`Unknown subagent id: ${agentId}`); + if (!isLocalAgentProvider(record.provider)) { + throw new Error(`Unknown subagent provider for existing session: ${record.provider}`); + } + + const workspaceRoot = assertAllowedPath(record.workspaceRoot, this.config.allowedRoots); + const profiles = await loadLocalAgentProfiles(this.config, workspaceRoot); + const profile = profiles.find((candidate) => candidate.name === record.profileName); + const prompt = profile ? profilePrompt(profile, turn.prompt) : rawProviderPrompt(record, turn.prompt); + const result = await this.runProvider(record.provider, { + prompt, + workspace: workspaceRoot, + providerSessionId: record.providerSessionId, + writeMode: "allowed", + model: turn.model, + thinking: turn.thinking, + }); + this.store.update(record.id, { + providerSessionId: result.providerSessionId ?? undefined, + status: "idle", + latestResponse: result.finalResponse, + error: undefined, + }); + } +} + +function profilePrompt(profile: LocalAgentProfile, prompt: string): string { + const body = profile.body.trim(); + return body ? `${body}\n\nTask:\n${prompt}` : prompt; +} + +function rawProviderPrompt(record: LocalAgentRecord, prompt: string): string { + if (record.profileName !== record.provider || !isLocalAgentProvider(record.provider)) { + throw new Error(`Subagent profile not found: ${record.profileName}`); + } + return prompt; +} diff --git a/src/local-agent-store.test.ts b/src/local-agent-store.test.ts index cf7265a9..f78d7173 100644 --- a/src/local-agent-store.test.ts +++ b/src/local-agent-store.test.ts @@ -24,6 +24,7 @@ try { assert.equal(store.get(created.id)?.thinking, "high"); assert.equal(store.get(created.id)?.profileName, "reviewer"); assert.equal(store.get(created.id.slice(0, 7))?.id, created.id); + assert.equal(store.getByAgentId(created.id.slice(0, 7))?.id, created.id); const updated = store.update(created.id, { status: "idle", @@ -35,6 +36,7 @@ try { assert.equal(updated.status, "idle"); assert.equal(updated.thinking, "medium"); assert.equal(store.get("thread_123")?.id, created.id); + assert.equal(store.getByAgentId("thread_123"), undefined); assert.equal(store.get(created.id)?.thinking, "medium"); assert.equal(store.update(created.id, { latestResponse: undefined }).latestResponse, undefined); assert.deepEqual( @@ -54,6 +56,11 @@ try { provider: "claude", }); + assert.throws( + () => store.getByAgentId("agt_"), + /Ambiguous subagent id prefix: agt_/, + ); + assert.deepEqual( store.list({ workspaceId: "ws_1" }).map((agent) => agent.id).sort(), [created.id, createdFromOtherStore.id].sort(), diff --git a/src/local-agent-store.ts b/src/local-agent-store.ts index a850ca9f..55d28073 100644 --- a/src/local-agent-store.ts +++ b/src/local-agent-store.ts @@ -152,6 +152,26 @@ export class LocalAgentStore { return matches.length === 1 ? rowToLocalAgentRecord(matches[0]!) : undefined; } + getByAgentId(idOrPrefix: string): LocalAgentRecord | undefined { + const exact = this.database.sqlite + .prepare("select * from local_agent_sessions where id = ?") + .get(idOrPrefix) as LocalAgentRow | undefined; + if (exact) return rowToLocalAgentRecord(exact); + + const matches = this.database.sqlite + .prepare( + `select * from local_agent_sessions + where id like ? escape '\\' + order by updated_at desc`, + ) + .all(`${escapeLike(idOrPrefix)}%`) as LocalAgentRow[]; + + if (matches.length > 1) { + throw new Error(`Ambiguous subagent id prefix: ${idOrPrefix}`); + } + return matches.length === 1 ? rowToLocalAgentRecord(matches[0]!) : undefined; + } + update(id: string, patch: Partial>): LocalAgentRecord { const current = this.getById(id); if (!current) throw new Error(`Unknown subagent id: ${id}`); diff --git a/src/server.ts b/src/server.ts index 840594ab..844a299a 100644 --- a/src/server.ts +++ b/src/server.ts @@ -61,6 +61,8 @@ import { getLocalAgentProviderAvailabilitySnapshot, type LocalAgentProviderAvailability, } from "./local-agent-availability.js"; +import { LocalAgentManager } from "./local-agent-manager.js"; +import { LocalAgentControlServer } from "./local-agent-control.js"; type Transport = StreamableHTTPServerTransport; // MCP clients can reconnect without closing the previous transport. Bound stale @@ -92,6 +94,7 @@ interface RunningServer { app: ReturnType; config: ServerConfig; localAgentProviders: LocalAgentProviderAvailability[]; + startAgentControl(): Promise; close(): Promise; } @@ -1694,6 +1697,10 @@ export function createServer( const localAgentProviders = config.subagents ? getLocalAgentProviderAvailabilitySnapshot() : []; + const localAgentManager = config.subagents ? new LocalAgentManager(config) : undefined; + const localAgentControl = localAgentManager + ? new LocalAgentControlServer(config, localAgentManager) + : undefined; const logSessionCloseResults = ( reason: "idle_timeout" | "server_shutdown", @@ -1875,21 +1882,33 @@ export function createServer( }); let closePromise: Promise | undefined; + const close = () => { + closePromise ??= (async () => { + clearInterval(sessionCleanupTimer); + const results = await transports.closeAll(); + logSessionCloseResults("server_shutdown", results); + processSessions.shutdown(); + await localAgentControl?.close(); + await localAgentManager?.shutdown(); + oauthProvider.close(); + workspaceStore.close?.(); + })(); + return closePromise; + }; + return { app, config, localAgentProviders, - close: () => { - closePromise ??= (async () => { - clearInterval(sessionCleanupTimer); - const results = await transports.closeAll(); - logSessionCloseResults("server_shutdown", results); - processSessions.shutdown(); - oauthProvider.close(); - workspaceStore.close?.(); - })(); - return closePromise; + startAgentControl: async () => { + try { + await localAgentControl?.start(); + } catch (error) { + await close().catch(() => undefined); + throw error; + } }, + close, }; } @@ -1902,7 +1921,8 @@ async function isMainModule(): Promise { } if (await isMainModule()) { - const { app, config, close, localAgentProviders } = createServer(); + const { app, config, close, localAgentProviders, startAgentControl } = createServer(); + await startAgentControl(); const httpServer = app.listen(config.port, config.host, () => { console.log( `devspace listening on http://${config.host}:${config.port}/mcp`,