diff --git a/docs/agent-profile-schema.md b/docs/agent-profile-schema.md index 0dc3db95..7817f6e5 100644 --- a/docs/agent-profile-schema.md +++ b/docs/agent-profile-schema.md @@ -2,7 +2,8 @@ DevSpace agent profiles are user-owned markdown files with YAML frontmatter. They describe roles such as reviewer, explorer, or implementer. -DevSpace owns provider invocation. +The internal on-demand `devspace-agentd` process owns provider invocation. The +CLI and MCP server use it as clients when they need agent execution. Profiles are discovered from: @@ -133,7 +134,8 @@ The Subagent skill teaches only: ```bash devspace agents ls -devspace agents run "" +devspace agents run "" +devspace agents continue "" devspace agents show ``` @@ -152,6 +154,10 @@ devspace agents show `devspace agents ls` lists existing subagent sessions for the current workspace; it does not list profile definitions. +Use `devspace agents continue ` for a later turn. The logical agent ID is +the `agt_...` value returned by `run` or `ls`; provider session IDs are not +accepted as substitutes. + The full profile body stays out of the model context until DevSpace launches the profile. @@ -161,5 +167,6 @@ profile. - Inferring changed files, tests, or diffs from worker output. - Exposing raw provider transcripts by default. - Teaching the model provider-specific CLIs. -- First-class MCP agent tools. Future tools should wrap the same provider - adapter registry used by `devspace agents`. +- First-class MCP agent tools. Future tools should call the same local agent + daemon used by `devspace agents` rather than executing providers in the MCP + server process. diff --git a/docs/chatgpt-coding-workflow.md b/docs/chatgpt-coding-workflow.md index ac2bdc7c..20f263aa 100644 --- a/docs/chatgpt-coding-workflow.md +++ b/docs/chatgpt-coding-workflow.md @@ -143,7 +143,8 @@ Skill paths may be outside the workspace. DevSpace only permits reading: Set `DEVSPACE_SKILLS=0` to hide skills from workspace output. Set `DEVSPACE_SUBAGENTS=1` to expose the experimental subagent catalog and `subagent-delegation` skill. That skill teaches the minimal -`devspace agents ls`, `devspace agents run`, and `devspace agents show` +`devspace agents ls`, `devspace agents run`, `devspace agents continue`, and +`devspace agents show` workflow. The catalog comes from `open_workspace`; `devspace agents ls` lists existing subagent sessions for that workspace. diff --git a/docs/configuration.md b/docs/configuration.md index 3502a98b..2a22ea03 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -152,7 +152,8 @@ agent without reading provider-specific launch details. `devspace agents ls` 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`, `devspace agents continue`, and `devspace agents show` +workflow. Starter profile templates are available under `examples/agents/`. Copy or adapt them into one of the active profile directories before use. diff --git a/docs/gotchas.md b/docs/gotchas.md index 639f54a9..5769d3e0 100644 --- a/docs/gotchas.md +++ b/docs/gotchas.md @@ -224,7 +224,10 @@ When `DEVSPACE_SUBAGENTS=1`, DevSpace loads agent profiles from `~/.devspace/agents/*.md` and project `.devspace/agents/*.md`, then exposes a compact profile catalog through `open_workspace`. The bundled `subagent-delegation` skill keeps the model-facing workflow to -`devspace agents ls`, `devspace agents run`, and `devspace agents show`. +`devspace agents ls`, `devspace agents run`, `devspace agents continue`, and +`devspace agents show`. +Those commands automatically manage the internal local agent daemon; `devspace +serve` is not a prerequisite. `devspace agents ls` lists existing subagent sessions, not profile definitions. diff --git a/docs/local-agent-daemon.md b/docs/local-agent-daemon.md new file mode 100644 index 00000000..b8195e28 --- /dev/null +++ b/docs/local-agent-daemon.md @@ -0,0 +1,54 @@ +# Local agent daemon + +Local agent execution is owned by an on-demand `devspace-agentd` process, not +by the MCP server and not by an individual CLI invocation. The daemon is an +internal implementation detail: the normal workflow remains: + +```text +devspace agents run/continue/show/ls + │ + ▼ + devspace-agentd + │ + ├── LocalAgentManager + ├── LocalAgentStore + ├── LocalAgentRuntimePool + └── provider runtimes +``` + +The CLI starts the daemon automatically when an agent command needs it. The +MCP server can use the same local client when an MCP operation needs agent +functionality, but `devspace serve` is not required for local-agent execution. +The daemon is scoped to one DevSpace `stateDir`, so one SQLite store and one +runtime owner serve all clients using that configuration. + +Communication uses a private Unix domain socket on Linux/macOS or a named pipe +on Windows. The endpoint is not exposed through the public MCP HTTP port. +Provider session identifiers and logical agent records are durable; live +provider runtimes are disposable and may be recreated after a daemon restart. + +The daemon state directory contains the socket or pipe identity, an atomic +lock, a PID marker, and diagnostic logs. A second client cannot start another +daemon for the same state directory. Stale lock and socket files are recovered +only after the recorded PID is no longer alive. + +The daemon is started on demand and may exit after its active turns, clients, +and warm runtime idle periods have ended. Users do not need to manage it during +normal operation. Diagnostic commands are available for startup, process, and +cleanup problems: + +```bash +devspace agents daemon status +devspace agents daemon stop +devspace agents daemon logs +``` + +Agent identity is explicit at the client boundary. `agents run` starts a new +logical agent from a profile or provider; `agents continue ` continues an +existing logical agent. Provider session IDs are never accepted as logical +agent IDs, and the daemon does not resolve ambiguous prefixes. + +Shutdown gives active turns a bounded graceful window. If that window expires, +the process exits with active records left durable; the next daemon startup +reconciles stale `starting` and `running` records to `error` without discarding +their `providerSessionId` or `latestResponse`. diff --git a/package-lock.json b/package-lock.json index 79993030..9b5de5f5 100644 --- a/package-lock.json +++ b/package-lock.json @@ -31,7 +31,8 @@ "zod": "^4.4.3" }, "bin": { - "devspace": "dist/cli.js" + "devspace": "dist/cli.js", + "devspace-agentd": "dist/local-agent-daemon-main.js" }, "devDependencies": { "@types/better-sqlite3": "^7.6.13", diff --git a/package.json b/package.json index 12983912..3ac3f08a 100644 --- a/package.json +++ b/package.json @@ -8,7 +8,8 @@ "node": ">=22.19 <27" }, "bin": { - "devspace": "dist/cli.js" + "devspace": "dist/cli.js", + "devspace-agentd": "dist/local-agent-daemon-main.js" }, "files": [ "dist", @@ -28,7 +29,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-daemon-lifecycle.test.ts && tsx src/local-agent-daemon-protocol.test.ts && tsx src/local-agent-daemon.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/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/skills/subagent-delegation/SKILL.md b/skills/subagent-delegation/SKILL.md index fb269df5..d143d272 100644 --- a/skills/subagent-delegation/SKILL.md +++ b/skills/subagent-delegation/SKILL.md @@ -18,12 +18,15 @@ Use only these commands for normal delegation: ```bash devspace agents ls -devspace agents run "" +devspace agents run "" +devspace agents continue "" devspace agents show ``` `ls` shows existing subagent sessions for the current workspace. DevSpace scopes it automatically from the shell environment injected by the workspace tool. +Use the returned logical `agt_...` ID with `continue`; provider session IDs and +prefixes are not interchangeable with logical agent IDs. `run ""` starts a new configured profile and prints a DevSpace agent id. @@ -37,6 +40,11 @@ profile is needed. Built-in providers are listed by `open_workspace`. running, `show` waits briefly. If there is still no final response, call `show` again later. +The commands automatically start the internal `devspace-agentd` process when +needed. `devspace serve` is not required for local-agent execution. The daemon +owns shared agent sessions and provider runtimes for the configured DevSpace +state directory. + Do not run provider CLIs such as `codex`, `claude`, `opencode`, `pi`, `cursor-agent`, or `copilot` directly unless you are explicitly debugging DevSpace agent integration. diff --git a/src/cli.test.ts b/src/cli.test.ts index 97b7084a..560a4eb8 100644 --- a/src/cli.test.ts +++ b/src/cli.test.ts @@ -1,11 +1,17 @@ import assert from "node:assert/strict"; -import { execFileSync } from "node:child_process"; +import { execFile, execFileSync } from "node:child_process"; import { mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; +import { createServer as createNetServer } from "node:net"; import { tmpdir } from "node:os"; import { join } from "node:path"; +import { promisify } from "node:util"; import { loadConfig } from "./config.js"; +import { localAgentDaemonPaths } from "./local-agent-daemon-lifecycle.js"; +import { encodeLocalAgentDaemonResponse } from "./local-agent-daemon-protocol.js"; import { LocalAgentStore } from "./local-agent-store.js"; +const execFileAsync = promisify(execFile); + const packageJson = JSON.parse(readFileSync(new URL("../package.json", import.meta.url), "utf8")) as { version: string; }; @@ -65,24 +71,66 @@ try { ); store.close(); - const output = execFileSync("node", ["--import", "tsx", "src/cli.ts", "agents", "ls"], { - cwd: process.cwd(), - encoding: "utf8", - env: { - ...process.env, - DEVSPACE_CONFIG_DIR: configDir, - DEVSPACE_ALLOWED_ROOTS: projectRoot, - DEVSPACE_STATE_DIR: stateDir, - DEVSPACE_WORKSPACE_ID: "ws_current", - DEVSPACE_WORKSPACE_ROOT: projectRoot, - DEVSPACE_SUBAGENTS: "1", - DEVSPACE_OAUTH_OWNER_TOKEN: "test-owner-token-that-is-long-enough", - }, + const daemonSocket = localAgentDaemonPaths(stateDir).endpoint; + const daemon = createNetServer((socket) => { + let buffer = ""; + socket.setEncoding("utf8"); + socket.on("data", (chunk: string | Buffer) => { + buffer += chunk.toString(); + const newline = buffer.indexOf("\n"); + if (newline === -1) return; + const request = JSON.parse(buffer.slice(0, newline)) as { requestId: string; method: string }; + const result = request.method === "agent.list" + ? [current] + : request.method === "hello" + ? { + state: "ready", + protocolVersion: 1, + pid: process.pid, + endpoint: daemonSocket, + startedAt: "now", + activeTurns: 0, + runtimeCount: 0, + clientConnections: 1, + } + : null; + socket.end(encodeLocalAgentDaemonResponse({ + requestId: request.requestId, + protocolVersion: 1, + ok: true, + result, + })); + }); }); + await new Promise((resolveListen, rejectListen) => { + daemon.once("error", rejectListen); + daemon.listen(daemonSocket, resolveListen); + }); + + try { + const { stdout: output } = await execFileAsync("node", ["--import", "tsx", "src/cli.ts", "agents", "ls"], { + cwd: process.cwd(), + encoding: "utf8", + env: { + ...process.env, + DEVSPACE_CONFIG_DIR: configDir, + DEVSPACE_ALLOWED_ROOTS: projectRoot, + DEVSPACE_STATE_DIR: stateDir, + DEVSPACE_WORKSPACE_ID: "ws_current", + DEVSPACE_WORKSPACE_ROOT: projectRoot, + DEVSPACE_SUBAGENTS: "1", + DEVSPACE_OAUTH_OWNER_TOKEN: "test-owner-token-that-is-long-enough", + }, + }); - assert.match(output, new RegExp(`${current.id} idle reviewer codex gpt-5\\.4 thinking=high`)); - assert.doesNotMatch(output, /profile reviewer/); - assert.doesNotMatch(output, new RegExp(other.id)); + assert.match(output, new RegExp(`${current.id} idle reviewer codex gpt-5\\.4 thinking=high`)); + assert.doesNotMatch(output, /profile reviewer/); + assert.doesNotMatch(output, new RegExp(other.id)); + } finally { + await new Promise((resolveClose, rejectClose) => { + daemon.close((error) => error ? rejectClose(error) : resolveClose()); + }); + } assert.equal(loadConfig({ DEVSPACE_CONFIG_DIR: configDir, diff --git a/src/cli.ts b/src/cli.ts index 7a1ac63f..222f1611 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -1,33 +1,20 @@ #!/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, + parseLocalAgentContinueArgs, 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 { createLocalAgentClient } from "./local-agent-client.js"; +import type { LocalAgentRecord } from "./local-agent-store.js"; import { ensureDevspaceDefaultSkills, generateOwnerToken, @@ -310,8 +297,10 @@ function printHelp(): void { " devspace config get Print persisted config", " devspace config set publicBaseUrl ", " devspace agents ls List subagent sessions", - " devspace agents run [--model ] ", + " devspace agents run [--model ] ", + " devspace agents continue [--model ] ", " devspace agents show ", + " devspace agents daemon ", " devspace -v, --version Print the installed version", "", "For temporary tunnels:", @@ -330,11 +319,14 @@ async function runAgentsCommand(args: string[]): Promise { case "run": await runAgentsRun(rest); return; + case "continue": + await runAgentsContinue(rest); + return; case "show": await runAgentsShow(rest); return; - case "__worker": - await runAgentsWorker(rest); + case "daemon": + await runAgentsDaemon(rest); return; case undefined: case "help": @@ -349,8 +341,8 @@ async function runAgentsCommand(args: string[]): Promise { async function runAgentsList(): Promise { const config = loadConfig(); - const store = createLocalAgentStore(config); - const agents = store.list(resolveCurrentWorkspaceScope()); + const client = createLocalAgentClient(config); + const agents = await client.list(resolveCurrentWorkspaceScope()); if (agents.length === 0) { console.log("No subagent sessions found for this workspace."); @@ -364,56 +356,30 @@ 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({ - workspaceId: process.env.DEVSPACE_WORKSPACE_ID, + const client = createLocalAgentClient(config); + const record = await client.start({ + target: parsed.target, + prompt: parsed.prompt, workspaceRoot, - profileName: target.name, - provider: target.provider, - model: target.model, - thinking: target.thinking, + workspaceId: process.env.DEVSPACE_WORKSPACE_ID, + model: parsed.model, + thinking: parsed.thinking, }); + console.log(formatAgentLine(record)); +} - spawnAgentWorker(record.id, promptFile); - console.log(formatAgentLine({ ...record, status: "running" })); +async function runAgentsContinue(args: string[]): Promise { + const parsed = parseLocalAgentContinueArgs(args); + const config = loadConfig(); + const client = createLocalAgentClient(config); + const scope = resolveCurrentWorkspaceScope(); + const record = await client.continue(parsed.agentId, parsed.prompt, { + model: parsed.model, + thinking: parsed.thinking, + }, scope); + console.log(formatAgentLine(record)); } async function runAgentsShow(args: string[]): Promise { @@ -421,14 +387,15 @@ async function runAgentsShow(args: string[]): Promise { if (!id) throw new Error("Usage: devspace agents show "); const config = loadConfig(); - const store = createLocalAgentStore(config); - let record = store.get(id); + const client = createLocalAgentClient(config); + const scope = resolveCurrentWorkspaceScope(); + let record = await client.get(id, scope); if (!record) throw new Error(`Unknown subagent id: ${id}`); const deadline = Date.now() + 15_000; while ((record.status === "starting" || record.status === "running") && Date.now() < deadline) { await sleep(500); - record = store.get(id) ?? record; + record = await client.get(id, scope) ?? record; } console.log(formatAgentLine(record)); @@ -445,96 +412,24 @@ 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 "); - } - +async function runAgentsDaemon(args: string[]): Promise { + const [subcommand] = args; 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}`); + const client = createLocalAgentClient(config); + switch (subcommand) { + case "status": + console.log(JSON.stringify(await client.status(), null, 2)); + return; + case "stop": + await client.stop(); + console.log("Local agent daemon stopped."); + return; + case "logs": + console.log(await client.logs()); + return; + default: + throw new Error("Usage: devspace agents daemon "); } - - 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 { @@ -568,8 +463,10 @@ function printAgentsHelp(): void { "", "Usage:", " devspace agents ls", - " devspace agents run [--model ] [--thinking ] ", + " devspace agents run [--model ] [--thinking ] ", + " devspace agents continue [--model ] [--thinking ] ", " devspace agents show ", + " devspace agents daemon ", ].join("\n"), ); } diff --git a/src/local-agent-adapters.ts b/src/local-agent-adapters.ts index 457b8e08..2bb88d13 100644 --- a/src/local-agent-adapters.ts +++ b/src/local-agent-adapters.ts @@ -1,14 +1,22 @@ import { spawn, spawnSync, type ChildProcessWithoutNullStreams } from "node:child_process"; import { resolve } from "node:path"; import { Readable, Writable } from "node:stream"; +import type { + ModelReasoningEffort, + SandboxMode, + ThreadOptions, +} from "@openai/codex-sdk"; import type { EffortLevel } from "@anthropic-ai/claude-agent-sdk"; import type { LocalAgentProvider } from "./local-agent-profiles.js"; import { removeDevspaceNodeModulesBinFromPath } from "./local-agent-path.js"; import { - createCodexSdkLocalAgentRuntime, type LocalAgentRunInput, type LocalAgentRunResult, + type LocalAgentRuntime, + type LocalAgentRuntimeContext, + type LocalAgentDriver, } from "./local-agent-runtime.js"; +import { LOCAL_AGENT_PROVIDERS } from "./local-agent-profiles.js"; export interface LocalAgentAdapter { readonly provider: LocalAgentProvider; @@ -44,15 +52,70 @@ export function createLocalAgentAdapter(provider: LocalAgentProvider): LocalAgen } } +export function createLocalAgentDrivers(): LocalAgentDriver[] { + return LOCAL_AGENT_PROVIDERS.map((provider) => new LegacyLocalAgentDriver(provider)); +} + +class LegacyLocalAgentDriver implements LocalAgentDriver { + readonly idleTimeoutMs = 60_000; + + constructor(readonly provider: LocalAgentProvider) {} + + runtimeKey(context: LocalAgentRuntimeContext): string { + return `legacy:${this.provider}:${context.agentId}`; + } + + async createRuntime(_context: LocalAgentRuntimeContext): Promise { + const adapter = createLocalAgentAdapter(this.provider); + return { + provider: this.provider, + run: (input) => adapter.run(input), + releaseSession: async () => undefined, + close: async () => undefined, + isAlive: () => true, + }; + } +} + class CodexLocalAgentAdapter implements LocalAgentAdapter { readonly provider = "codex" as const; async run(input: LocalAgentRunInput): Promise { - const runtime = await createCodexSdkLocalAgentRuntime(); - return runtime.run(input); + const module = await import("@openai/codex-sdk"); + const codex = new module.Codex(); + const options = threadOptionsFor(input); + const thread = input.providerSessionId + ? codex.resumeThread(input.providerSessionId, options) + : codex.startThread(options); + const turn = await thread.run(input.prompt); + return { + provider: this.provider, + providerSessionId: thread.id, + finalResponse: turn.finalResponse, + items: turn.items, + }; } } +function sandboxModeFor(writeMode: LocalAgentRunInput["writeMode"]): SandboxMode { + switch (writeMode) { + case "allowed": return "workspace-write"; + case "full_access": return "danger-full-access"; + case "read_only": + case undefined: return "read-only"; + } +} + +function threadOptionsFor(input: LocalAgentRunInput): ThreadOptions { + return { + workingDirectory: input.workspace, + sandboxMode: sandboxModeFor(input.writeMode), + approvalPolicy: "never", + model: input.model, + modelReasoningEffort: input.thinking as ModelReasoningEffort | undefined, + }; +} + class ClaudeLocalAgentAdapter implements LocalAgentAdapter { readonly provider = "claude" as const; diff --git a/src/local-agent-availability.ts b/src/local-agent-availability.ts index 747f304f..ca2372a1 100644 --- a/src/local-agent-availability.ts +++ b/src/local-agent-availability.ts @@ -1,4 +1,4 @@ -import { spawnSync } from "node:child_process"; +import { accessSync, constants } from "node:fs"; import { delimiter, resolve } from "node:path"; import { removeDevspaceNodeModulesBinFromPath } from "./local-agent-path.js"; import { @@ -101,10 +101,10 @@ function commandAvailability( function resolveCommand(command: string, env: NodeJS.ProcessEnv = process.env): string | undefined { const commandHasPath = command.includes("/") || command.includes("\\"); - if (commandHasPath) return executableExists(command, env) ? command : undefined; + if (commandHasPath) return executableExists(command) ? command : undefined; for (const candidate of candidateCommandPaths(command, env)) { - if (executableExists(candidate, env)) return candidate; + if (executableExists(candidate)) return candidate; } return undefined; } @@ -127,17 +127,14 @@ function candidateCommandPaths(command: string, env: NodeJS.ProcessEnv): string[ return candidates; } -function executableExists(command: string, env: NodeJS.ProcessEnv): boolean { - const result = spawnSync(command, ["--version"], { - encoding: "utf8", - env, - windowsHide: true, - timeout: 5_000, - }); - const code = typeof result.error === "object" && result.error && "code" in result.error - ? result.error.code - : undefined; - return code !== "ENOENT"; +function executableExists(command: string): boolean { + const mode = process.platform === "win32" ? constants.F_OK : constants.X_OK; + try { + accessSync(command, mode); + return true; + } catch { + return false; + } } function piAvailabilityEnvironment(env: NodeJS.ProcessEnv): NodeJS.ProcessEnv { diff --git a/src/local-agent-client.ts b/src/local-agent-client.ts new file mode 100644 index 00000000..efdae7e9 --- /dev/null +++ b/src/local-agent-client.ts @@ -0,0 +1,252 @@ +import { existsSync } from "node:fs"; +import { randomUUID } from "node:crypto"; +import { spawn } from "node:child_process"; +import { createConnection } from "node:net"; +import { fileURLToPath } from "node:url"; +import type { ServerConfig } from "./config.js"; +import { + decodeAgentRecord, + decodeAgentRecordList, + decodeDaemonLogs, + decodeDaemonStatus, + decodeLocalAgentDaemonResponse, + encodeLocalAgentDaemonRequest, + LocalAgentDaemonProtocolError, + type LocalAgentDaemonRequest, + type LocalAgentDaemonResponse, + type LocalAgentDaemonStatus, +} from "./local-agent-daemon-protocol.js"; +import { + LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + ensureLocalAgentDaemonSecret, + localAgentDaemonPaths, + type LocalAgentDaemonPaths, +} from "./local-agent-daemon-lifecycle.js"; +import type { RunOverrides, StartLocalAgentInput } from "./local-agent-manager.js"; +import type { LocalAgentListScope, LocalAgentRecord, LocalAgentWorkspaceScope } from "./local-agent-store.js"; + +const DEFAULT_STARTUP_TIMEOUT_MS = 8_000; +const DEFAULT_REQUEST_TIMEOUT_MS = 30_000; +const RETRY_DELAY_MS = 40; + +export interface LocalAgentClientOptions { + stateDir: string; + startupTimeoutMs?: number; + requestTimeoutMs?: number; + spawnDaemon?: () => void; + endpoint?: string; +} + +export class LocalAgentDaemonClientError extends Error { + constructor(readonly code: string, message: string) { + super(message); + this.name = "LocalAgentDaemonClientError"; + } +} + +export class LocalAgentClient { + private readonly stateDir: string; + private readonly paths: LocalAgentDaemonPaths; + private readonly endpoint: string; + private readonly startupTimeoutMs: number; + private readonly requestTimeoutMs: number; + private readonly spawnDaemon: () => void; + private startupPromise?: Promise; + + constructor(options: LocalAgentClientOptions) { + this.stateDir = options.stateDir; + this.paths = localAgentDaemonPaths(options.stateDir); + this.endpoint = options.endpoint ?? this.paths.endpoint; + this.startupTimeoutMs = options.startupTimeoutMs ?? DEFAULT_STARTUP_TIMEOUT_MS; + this.requestTimeoutMs = options.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS; + this.spawnDaemon = options.spawnDaemon ?? (() => spawnLocalAgentDaemon(options.stateDir)); + } + + async run(input: StartLocalAgentInput): Promise { + return this.start(input); + } + + async start(input: StartLocalAgentInput): Promise { + const result = await this.request("agent.start", input); + return decodeAgentRecord(result); + } + + async continue(agentId: string, prompt: string, overrides: RunOverrides = {}, scope: LocalAgentWorkspaceScope): Promise { + const result = await this.request("agent.continue", { + id: agentId, + prompt, + scope, + ...(Object.keys(overrides).length > 0 ? { overrides } : {}), + }); + return decodeAgentRecord(result); + } + + async get(agentId: string, scope: LocalAgentWorkspaceScope): Promise { + const result = await this.request("agent.get", { id: agentId, scope }); + return result === null ? undefined : decodeAgentRecord(result); + } + + async list(scope: LocalAgentListScope): Promise { + return decodeAgentRecordList(await this.request("agent.list", scope)); + } + + async status(): Promise { + return decodeDaemonStatus(await this.request("daemon.status", {})); + } + + async stop(): Promise { + return decodeDaemonStatus(await this.request("daemon.stop", {})); + } + + async logs(lines = 200): Promise { + return decodeDaemonLogs(await this.request("daemon.logs", { lines })); + } + + async ensureReady(): Promise { + if (this.startupPromise) return this.startupPromise; + this.startupPromise = this.ensureReadyInternal().finally(() => { + this.startupPromise = undefined; + }); + return this.startupPromise; + } + + private async ensureReadyInternal(): Promise { + const existing = await this.tryHello(); + if (existing) return existing; + + this.spawnDaemon(); + const deadline = Date.now() + this.startupTimeoutMs; + let lastError: unknown; + while (Date.now() < deadline) { + await delay(RETRY_DELAY_MS); + try { + const ready = await this.tryHello(); + if (ready) return ready; + } catch (error) { + lastError = error; + if (error instanceof LocalAgentDaemonClientError && error.code === "PROTOCOL_MISMATCH") throw error; + } + } + const suffix = lastError instanceof Error ? `: ${lastError.message}` : ""; + throw new LocalAgentDaemonClientError( + "DAEMON_START_FAILED", + `Unable to start the local agent daemon in ${this.stateDir}${suffix}`, + ); + } + + private async tryHello(): Promise { + try { + const response = await sendRequest(this.endpoint, { + requestId: randomUUID(), + protocolVersion: LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + authToken: ensureLocalAgentDaemonSecret(this.paths), + method: "hello", + params: {}, + }, this.requestTimeoutMs); + if (!response.ok) { + if (response.error.code === "PROTOCOL_MISMATCH") { + throw new LocalAgentDaemonClientError(response.error.code, response.error.message); + } + return undefined; + } + const status = decodeDaemonStatus(response.result); + return status.state === "ready" ? status : undefined; + } catch (error) { + if (error instanceof LocalAgentDaemonClientError && error.code === "PROTOCOL_MISMATCH") throw error; + return undefined; + } + } + + private async request( + method: M, + params: Extract['params'], + ): Promise { + await this.ensureReady(); + const response = await sendRequest(this.endpoint, { + requestId: randomUUID(), + protocolVersion: LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + authToken: ensureLocalAgentDaemonSecret(this.paths), + method, + params, + } as LocalAgentDaemonRequest, this.requestTimeoutMs); + if (!response.ok) { + throw new LocalAgentDaemonClientError(response.error.code, response.error.message); + } + return response.result; + } +} + +export function createLocalAgentClient(config: Pick): LocalAgentClient { + return new LocalAgentClient({ stateDir: config.stateDir }); +} + +export function spawnLocalAgentDaemon(stateDir: string, env: NodeJS.ProcessEnv = process.env): void { + const entrypoint = resolveDaemonEntrypoint(); + const child = spawn(process.execPath, [...process.execArgv, entrypoint], { + detached: true, + stdio: "ignore", + windowsHide: true, + env: { ...env, DEVSPACE_STATE_DIR: stateDir }, + }); + child.unref(); +} + +export function resolveDaemonEntrypoint(): string { + const compiled = fileURLToPath(new URL("./local-agent-daemon-main.js", import.meta.url)); + if (existsSync(compiled)) return compiled; + return fileURLToPath(new URL("./local-agent-daemon-main.ts", import.meta.url)); +} + +async function sendRequest( + endpoint: string, + request: LocalAgentDaemonRequest, + timeoutMs: number, +): Promise { + return new Promise((resolve, reject) => { + const socket = createConnection(endpoint); + let buffer = ""; + let settled = false; + const timer = setTimeout(() => { + finish(new LocalAgentDaemonClientError("REQUEST_TIMEOUT", "Timed out waiting for the local agent daemon."), true); + }, timeoutMs); + + const finish = (error?: unknown, destroy = false) => { + if (settled) return; + settled = true; + clearTimeout(timer); + if (destroy) socket.destroy(); + if (error) reject(error); + }; + + socket.setEncoding("utf8"); + socket.on("data", (chunk: string | Buffer) => { + buffer += chunk.toString(); + const newline = buffer.indexOf("\n"); + if (newline === -1) return; + try { + const response = decodeLocalAgentDaemonResponse(JSON.parse(buffer.slice(0, newline)) as unknown); + if (response.requestId !== request.requestId) { + throw new LocalAgentDaemonProtocolError("INVALID_RESPONSE", "Daemon response request id did not match."); + } + settled = true; + clearTimeout(timer); + resolve(response); + socket.end(); + } catch (error) { + finish(error, true); + } + }); + socket.once("error", (error) => finish(new LocalAgentDaemonClientError( + (error as NodeJS.ErrnoException).code ?? "DAEMON_UNAVAILABLE", + error.message, + ))); + socket.once("close", () => { + if (!settled) finish(new LocalAgentDaemonClientError("DAEMON_UNAVAILABLE", "Local agent daemon closed the connection.")); + }); + socket.once("connect", () => socket.write(encodeLocalAgentDaemonRequest(request))); + }); +} + +function delay(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} diff --git a/src/local-agent-daemon-lifecycle.test.ts b/src/local-agent-daemon-lifecycle.test.ts new file mode 100644 index 00000000..b4438cdc --- /dev/null +++ b/src/local-agent-daemon-lifecycle.test.ts @@ -0,0 +1,64 @@ +import assert from "node:assert/strict"; +import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { + LocalAgentDaemonAlreadyRunningError, + LocalAgentDaemonLock, + ensureLocalAgentDaemonStateDir, + isProcessAlive, + localAgentDaemonPaths, + removeLocalAgentDaemonFiles, + ensureLocalAgentDaemonSecret, + writeLocalAgentDaemonPid, +} from "./local-agent-daemon-lifecycle.js"; + +const root = await mkdtemp(join(tmpdir(), "devspace-agentd-lifecycle-test-")); +try { + const paths = localAgentDaemonPaths(join(root, "state")); + ensureLocalAgentDaemonStateDir(paths.stateDir); + const lock = new LocalAgentDaemonLock(paths); + lock.acquire(); + assert.equal(await readFile(paths.lockPath, "utf8"), `${process.pid}\n`); + assert.throws( + () => new LocalAgentDaemonLock(paths).acquire(), + (error: unknown) => error instanceof LocalAgentDaemonAlreadyRunningError, + ); + await writeFile(paths.pidPath, "999999\n", { mode: 0o600 }); + assert.throws( + () => new LocalAgentDaemonLock(paths).acquire(), + (error: unknown) => error instanceof LocalAgentDaemonAlreadyRunningError, + "a stale diagnostic PID must not override the live lock owner", + ); + assert.equal(ensureLocalAgentDaemonSecret(paths).length, 64); + lock.release(); + + await writeFile(paths.lockPath, "999999999\n", { mode: 0o600 }); + await writeFile(paths.pidPath, "999999999\n", { mode: 0o600 }); + const recovered = new LocalAgentDaemonLock(paths); + recovered.acquire(); + assert.equal(await readFile(paths.lockPath, "utf8"), `${process.pid}\n`); + writeLocalAgentDaemonPid(paths); + assert.equal(await readFile(paths.pidPath, "utf8"), `${process.pid}\n`); + assert.equal(isProcessAlive(process.pid), true); + recovered.release(); + + await writeFile(paths.lockPath, "not-a-pid\n", { mode: 0o600 }); + assert.throws( + () => new LocalAgentDaemonLock(paths).acquire(), + (error: unknown) => error instanceof LocalAgentDaemonAlreadyRunningError, + "an undecodable lock must fail closed instead of being deleted by age", + ); + assert.equal(await readFile(paths.lockPath, "utf8"), "not-a-pid\n"); + await rm(paths.lockPath, { force: true }); + + await writeFile(paths.secretPath, "not-a-hex-secret\n", { mode: 0o600 }); + assert.throws( + () => ensureLocalAgentDaemonSecret(paths), + /secret is invalid/, + "daemon secrets must be exactly 64 hexadecimal characters", + ); + removeLocalAgentDaemonFiles(paths); +} finally { + await rm(root, { recursive: true, force: true }); +} diff --git a/src/local-agent-daemon-lifecycle.ts b/src/local-agent-daemon-lifecycle.ts new file mode 100644 index 00000000..29129c54 --- /dev/null +++ b/src/local-agent-daemon-lifecycle.ts @@ -0,0 +1,210 @@ +import { createHash, randomBytes } from "node:crypto"; +import { + chmodSync, + closeSync, + linkSync, + mkdirSync, + openSync, + readFileSync, + renameSync, + rmSync, + writeSync, +} from "node:fs"; +import { join, resolve } from "node:path"; + +export const LOCAL_AGENT_DAEMON_PROTOCOL_VERSION = 1; +export const LOCAL_AGENT_DAEMON_SOCKET_NAME = "agentd.sock"; +export const LOCAL_AGENT_DAEMON_PID_NAME = "agentd.pid"; +export const LOCAL_AGENT_DAEMON_LOCK_NAME = "agentd.lock"; +export const LOCAL_AGENT_DAEMON_SECRET_NAME = "agentd.secret"; +export const LOCAL_AGENT_DAEMON_LOG_NAME = "agentd.log"; + +export interface LocalAgentDaemonPaths { + stateDir: string; + socketPath: string; + pidPath: string; + lockPath: string; + secretPath: string; + logPath: string; + endpoint: string; +} + +export function localAgentDaemonPaths( + stateDir: string, + platform: NodeJS.Platform = process.platform, +): LocalAgentDaemonPaths { + const resolvedStateDir = resolve(stateDir); + const socketPath = join(resolvedStateDir, LOCAL_AGENT_DAEMON_SOCKET_NAME); + return { + stateDir: resolvedStateDir, + socketPath, + pidPath: join(resolvedStateDir, LOCAL_AGENT_DAEMON_PID_NAME), + lockPath: join(resolvedStateDir, LOCAL_AGENT_DAEMON_LOCK_NAME), + secretPath: join(resolvedStateDir, LOCAL_AGENT_DAEMON_SECRET_NAME), + logPath: join(resolvedStateDir, LOCAL_AGENT_DAEMON_LOG_NAME), + endpoint: platform === "win32" + ? `\\\\.\\pipe\\devspace-agentd-${hashStateDir(resolvedStateDir)}` + : socketPath, + }; +} + +export function ensureLocalAgentDaemonStateDir(stateDir: string): void { + mkdirSync(stateDir, { recursive: true, mode: 0o700 }); + chmodSync(stateDir, 0o700); +} + +export class LocalAgentDaemonAlreadyRunningError extends Error { + readonly code = "DAEMON_ALREADY_RUNNING" as const; + + constructor(readonly pid?: number) { + super(pid ? `Local agent daemon is already running (pid ${pid}).` : "Local agent daemon is already running."); + this.name = "LocalAgentDaemonAlreadyRunningError"; + } +} + +export class LocalAgentDaemonLock { + private acquired = false; + + constructor(readonly paths: LocalAgentDaemonPaths) {} + + acquire(): void { + ensureLocalAgentDaemonStateDir(this.paths.stateDir); + for (let attempt = 0; attempt < 2; attempt += 1) { + const temporaryPath = `${this.paths.lockPath}.${process.pid}.${randomBytes(8).toString("hex")}.tmp`; + let published = false; + try { + writeFileSecure(temporaryPath, `${process.pid}\n`); + // Publish the owner record atomically. An empty lock must never be + // visible to stale-lock recovery between create and write. + linkSync(temporaryPath, this.paths.lockPath); + published = true; + rmSync(temporaryPath, { force: true }); + chmodSync(this.paths.lockPath, 0o600); + writeFileSecure(this.paths.pidPath, `${process.pid}\n`); + this.acquired = true; + return; + } catch (error) { + rmSync(temporaryPath, { force: true }); + if (published && readDaemonPid(this.paths.lockPath) === process.pid) { + rmSync(this.paths.lockPath, { force: true }); + } + if (!isFileExistsError(error)) throw error; + const pid = readDaemonPid(this.paths.lockPath); + if (pid !== undefined && isProcessAlive(pid)) { + throw new LocalAgentDaemonAlreadyRunningError(pid); + } + if (pid === undefined) { + // An undecodable lock may belong to a process that has not finished + // publishing its owner record. Refuse to delete it automatically. + throw new LocalAgentDaemonAlreadyRunningError(); + } + if (!removeStaleLock(this.paths.lockPath)) continue; + } + } + throw new LocalAgentDaemonAlreadyRunningError(readDaemonPid(this.paths.lockPath)); + } + + release(): void { + if (!this.acquired) return; + this.acquired = false; + if (readDaemonPid(this.paths.pidPath) === process.pid) { + rmSync(this.paths.pidPath, { force: true }); + } + if (readDaemonPid(this.paths.lockPath) === process.pid) { + rmSync(this.paths.lockPath, { force: true }); + } + } +} + +export function writeLocalAgentDaemonPid(paths: LocalAgentDaemonPaths): void { + writeFileSecure(paths.pidPath, `${process.pid}\n`); +} + +export function ensureLocalAgentDaemonSecret(paths: LocalAgentDaemonPaths): string { + ensureLocalAgentDaemonStateDir(paths.stateDir); + try { + const secret = readFileSync(paths.secretPath, "utf8").trim(); + if (isDaemonSecret(secret)) return secret; + } catch { + // Create the secret below. + } + + const secret = randomBytes(32).toString("hex"); + try { + const fileDescriptor = openSync(paths.secretPath, "wx", 0o600); + try { + writeSync(fileDescriptor, `${secret}\n`); + chmodSync(paths.secretPath, 0o600); + return secret; + } finally { + closeSync(fileDescriptor); + } + } catch (error) { + if (!isFileExistsError(error)) throw error; + const existing = readFileSync(paths.secretPath, "utf8").trim(); + if (!isDaemonSecret(existing)) throw new Error("Local agent daemon secret is invalid."); + return existing; + } +} + +export function removeLocalAgentDaemonFiles(paths: LocalAgentDaemonPaths): void { + rmSync(paths.pidPath, { force: true }); + if (process.platform !== "win32") rmSync(paths.socketPath, { force: true }); +} + +export function readDaemonPid(pidPath: string): number | undefined { + try { + const value = readFileSync(pidPath, "utf8").trim(); + if (!/^\d+$/.test(value)) return undefined; + const pid = Number(value); + return Number.isSafeInteger(pid) && pid > 0 ? pid : undefined; + } catch { + return undefined; + } +} + +export function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (error) { + return (error as NodeJS.ErrnoException).code === "EPERM"; + } +} + +function writeFileSecure(path: string, content: string): void { + const fileDescriptor = openSync(path, "w", 0o600); + try { + writeSync(fileDescriptor, content); + chmodSync(path, 0o600); + } finally { + closeSync(fileDescriptor); + } +} + +function isFileExistsError(error: unknown): boolean { + return (error as NodeJS.ErrnoException).code === "EEXIST"; +} + +function removeStaleLock(path: string): boolean { + const stalePath = `${path}.stale-${process.pid}-${randomBytes(8).toString("hex")}`; + try { + // Rename moves the exact lock we inspected out of the ownership path. If + // another contender publishes a new lock after this point, it is never + // removed with the stale one. + renameSync(path, stalePath); + rmSync(stalePath, { force: true }); + return true; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return false; + throw error; + } +} + +function isDaemonSecret(secret: string): boolean { + return /^[0-9a-f]{64}$/i.test(secret); +} + +function hashStateDir(stateDir: string): string { + return createHash("sha256").update(stateDir).digest("hex").slice(0, 24); +} diff --git a/src/local-agent-daemon-main.ts b/src/local-agent-daemon-main.ts new file mode 100644 index 00000000..59a685f9 --- /dev/null +++ b/src/local-agent-daemon-main.ts @@ -0,0 +1,88 @@ +#!/usr/bin/env node +import { loadConfig } from "./config.js"; +import { createLocalAgentDrivers } from "./local-agent-adapters.js"; +import { loadLocalAgentProfiles } from "./local-agent-profiles.js"; +import { LocalAgentDaemon, writeLocalAgentDaemonLog } from "./local-agent-daemon.js"; +import { + LocalAgentDaemonAlreadyRunningError, + localAgentDaemonPaths, +} from "./local-agent-daemon-lifecycle.js"; +import { LocalAgentManager } from "./local-agent-manager.js"; +import { LocalAgentRuntimePool } from "./local-agent-runtime-pool.js"; +import { LocalAgentStore } from "./local-agent-store.js"; + +const config = loadConfig(); +const DEFAULT_DAEMON_SHUTDOWN_TIMEOUT_MS = 10_000; +const paths = localAgentDaemonPaths(config.stateDir); +const log = ( + level: "info" | "warn" | "error", + event: string, + fields: Record, +) => writeLocalAgentDaemonLog(paths, level, event, fields); +const store = new LocalAgentStore(paths.stateDir); +const manager = new LocalAgentManager({ + store, + drivers: createLocalAgentDrivers(), + pool: new LocalAgentRuntimePool({ logger: log }), + loadProfiles: (workspaceRoot) => loadLocalAgentProfiles(config, workspaceRoot), + agentDir: config.agentDir, + allowedRoots: config.allowedRoots, + logger: log, +}); +const daemon = new LocalAgentDaemon({ + stateDir: paths.stateDir, + manager, + onLockAcquired: () => { manager.reconcileActiveRuns(); }, + onClosed: () => { if (!shuttingDown) process.exit(0); }, + idleShutdownMs: parseIdleShutdownMs(process.env.DEVSPACE_AGENTD_IDLE_TIMEOUT_MS), +}); + +let shuttingDown = false; +const shutdown = () => { + if (shuttingDown) return; + shuttingDown = true; + const forceTimer = setTimeout(() => { + log("error", "daemon_forced_shutdown", { + activeTurns: manager.activeTurnCount, + runtimeCount: manager.runtimeCount, + }); + // Active records intentionally remain durable. The next daemon startup + // reconciles them to error while preserving provider continuation data. + process.exit(1); + }, parseShutdownTimeoutMs(process.env.DEVSPACE_AGENTD_SHUTDOWN_TIMEOUT_MS)); + forceTimer.unref(); + void daemon.close().finally(() => process.exit(0)); +}; +process.once("SIGINT", shutdown); +process.once("SIGTERM", shutdown); + +try { + await daemon.start(); +} catch (error) { + if (error instanceof LocalAgentDaemonAlreadyRunningError) { + await manager.close(); + process.exit(0); + } + log("error", "daemon_start_failed", { error: error instanceof Error ? error.message : String(error) }); + await manager.close(); + console.error(error instanceof Error ? error.message : String(error)); + process.exit(1); +} + +function parseIdleShutdownMs(value: string | undefined): number { + if (value === undefined || value.trim() === "") return 30_000; + const parsed = Number(value); + if (!Number.isFinite(parsed) || parsed < 0) { + throw new Error("DEVSPACE_AGENTD_IDLE_TIMEOUT_MS must be a non-negative duration."); + } + return parsed; +} + +function parseShutdownTimeoutMs(value: string | undefined): number { + if (value === undefined || value.trim() === "") return DEFAULT_DAEMON_SHUTDOWN_TIMEOUT_MS; + const parsed = Number(value); + if (!Number.isFinite(parsed) || parsed < 0) { + throw new Error("DEVSPACE_AGENTD_SHUTDOWN_TIMEOUT_MS must be a non-negative duration."); + } + return parsed; +} diff --git a/src/local-agent-daemon-protocol.test.ts b/src/local-agent-daemon-protocol.test.ts new file mode 100644 index 00000000..480a8b33 --- /dev/null +++ b/src/local-agent-daemon-protocol.test.ts @@ -0,0 +1,54 @@ +import assert from "node:assert/strict"; +import { + decodeAgentRecord, + decodeLocalAgentDaemonRequest, + decodeLocalAgentDaemonResponse, + encodeLocalAgentDaemonRequest, + LocalAgentDaemonProtocolError, +} from "./local-agent-daemon-protocol.js"; + +const request = decodeLocalAgentDaemonRequest({ + requestId: "req_1", + protocolVersion: 1, + authToken: "test-secret", + method: "agent.start", + params: { + target: "reviewer", + prompt: "Review this", + workspaceRoot: "/tmp/project", + writeMode: "read_only", + }, +}); +assert.equal(request.method, "agent.start"); +assert.equal(request.params.writeMode, "read_only"); +assert.match(encodeLocalAgentDaemonRequest(request), /"method":"agent.start"/); + +assert.throws( + () => decodeLocalAgentDaemonRequest({ + requestId: "req_2", + protocolVersion: 1, + authToken: "test-secret", + method: "agent.start", + params: { target: "reviewer", prompt: "" }, + }), + (error: unknown) => error instanceof LocalAgentDaemonProtocolError && error.code === "INVALID_PARAMS", +); + +const record = decodeAgentRecord({ + id: "agt_1234", + workspaceRoot: "/tmp/project", + profileName: "reviewer", + provider: "codex", + status: "idle", + createdAt: "now", + updatedAt: "now", +}); +assert.equal(record.id, "agt_1234"); + +const response = decodeLocalAgentDaemonResponse({ + requestId: "req_1", + protocolVersion: 1, + ok: true, + result: record, +}); +assert.equal(response.ok, true); diff --git a/src/local-agent-daemon-protocol.ts b/src/local-agent-daemon-protocol.ts new file mode 100644 index 00000000..dabb0549 --- /dev/null +++ b/src/local-agent-daemon-protocol.ts @@ -0,0 +1,326 @@ +import type { + LocalAgentListScope, + LocalAgentRecord, + LocalAgentStatus, + LocalAgentWorkspaceScope, +} from "./local-agent-store.js"; +import type { + RunOverrides, + StartLocalAgentInput, +} from "./local-agent-manager.js"; +import type { LocalAgentWriteMode } from "./local-agent-runtime.js"; +import { LOCAL_AGENT_DAEMON_PROTOCOL_VERSION } from "./local-agent-daemon-lifecycle.js"; + +export type LocalAgentDaemonMethod = + | "hello" + | "agent.start" + | "agent.continue" + | "agent.get" + | "agent.list" + | "daemon.status" + | "daemon.stop" + | "daemon.logs"; + +export type LocalAgentDaemonRequest = + | AgentDaemonRequestBase<"hello", Record> + | AgentDaemonRequestBase<"agent.start", StartLocalAgentInput> + | AgentDaemonRequestBase<"agent.continue", { id: string; prompt: string; scope: LocalAgentWorkspaceScope; overrides?: RunOverrides }> + | AgentDaemonRequestBase<"agent.get", { id: string; scope: LocalAgentWorkspaceScope }> + | AgentDaemonRequestBase<"agent.list", LocalAgentListScope> + | AgentDaemonRequestBase<"daemon.status", Record> + | AgentDaemonRequestBase<"daemon.stop", Record> + | AgentDaemonRequestBase<"daemon.logs", { lines?: number }>; + +interface AgentDaemonRequestBase< + M extends LocalAgentDaemonMethod, + P, +> { + requestId: string; + protocolVersion: number; + authToken: string; + method: M; + params: P; +} + +export interface LocalAgentDaemonStatus { + state: "ready" | "stopping"; + protocolVersion: number; + pid: number; + endpoint: string; + startedAt: string; + activeTurns: number; + runtimeCount: number; + clientConnections: number; +} + +export interface LocalAgentDaemonErrorPayload { + code: string; + message: string; +} + +export type LocalAgentDaemonResponse = + | { + requestId: string; + protocolVersion: number; + ok: true; + result: unknown; + } + | { + requestId: string; + protocolVersion: number; + ok: false; + error: LocalAgentDaemonErrorPayload; + }; + +export function encodeLocalAgentDaemonRequest(request: LocalAgentDaemonRequest): string { + return `${JSON.stringify(request)}\n`; +} + +export function encodeLocalAgentDaemonResponse(response: LocalAgentDaemonResponse): string { + return `${JSON.stringify(response)}\n`; +} + +export function decodeLocalAgentDaemonRequest(value: unknown): LocalAgentDaemonRequest { + const record = asRecord(value); + const requestId = requiredString(record?.requestId, "requestId"); + const protocolVersion = requiredInteger(record?.protocolVersion, "protocolVersion"); + const authToken = requiredString(record?.authToken, "authToken"); + const method = requiredString(record?.method, "method") as LocalAgentDaemonMethod; + const params = record?.params; + + switch (method) { + case "hello": + case "daemon.status": + case "daemon.stop": + return { requestId, protocolVersion, authToken, method, params: decodeEmptyParams(params) } as LocalAgentDaemonRequest; + case "agent.start": + return { + requestId, + protocolVersion, + authToken, + method, + params: decodeStartInput(params), + } as LocalAgentDaemonRequest; + case "agent.continue": + return { + requestId, + protocolVersion, + authToken, + method, + params: decodeContinueInput(params), + } as LocalAgentDaemonRequest; + case "agent.get": + return { + requestId, + protocolVersion, + method, + authToken, + params: { + id: requiredString(asRecord(params)?.id, "id"), + scope: decodeWorkspaceScope(asRecord(params)?.scope), + }, + } as LocalAgentDaemonRequest; + case "agent.list": + return { + requestId, + protocolVersion, + authToken, + method, + params: decodeListScope(params), + } as LocalAgentDaemonRequest; + case "daemon.logs": + return { + requestId, + protocolVersion, + authToken, + method, + params: decodeLogsParams(params), + } as LocalAgentDaemonRequest; + default: + throw new LocalAgentDaemonProtocolError("UNKNOWN_METHOD", `Unknown daemon method: ${method}`); + } +} + +export function decodeLocalAgentDaemonResponse(value: unknown): LocalAgentDaemonResponse { + const record = asRecord(value); + const requestId = requiredString(record?.requestId, "requestId"); + const protocolVersion = requiredInteger(record?.protocolVersion, "protocolVersion"); + if (record?.ok === true) { + return { requestId, protocolVersion, ok: true, result: record.result }; + } + if (record?.ok === false) { + const error = asRecord(record.error); + return { + requestId, + protocolVersion, + ok: false, + error: { + code: requiredString(error?.code, "error.code"), + message: requiredString(error?.message, "error.message"), + }, + }; + } + throw new LocalAgentDaemonProtocolError("INVALID_RESPONSE", "Daemon returned an invalid response."); +} + +export function decodeAgentRecord(value: unknown): LocalAgentRecord { + const record = asRecord(value); + const status = requiredString(record?.status, "status"); + if (!isLocalAgentStatus(status)) throw new LocalAgentDaemonProtocolError("INVALID_RECORD", "Invalid agent status."); + return { + id: requiredString(record?.id, "id"), + workspaceId: optionalString(record?.workspaceId), + workspaceRoot: requiredString(record?.workspaceRoot, "workspaceRoot"), + profileName: requiredString(record?.profileName, "profileName"), + provider: requiredString(record?.provider, "provider"), + model: optionalString(record?.model), + thinking: optionalString(record?.thinking), + providerSessionId: optionalString(record?.providerSessionId), + status, + latestResponse: optionalString(record?.latestResponse), + error: optionalString(record?.error), + createdAt: requiredString(record?.createdAt, "createdAt"), + updatedAt: requiredString(record?.updatedAt, "updatedAt"), + }; +} + +export function decodeAgentRecordList(value: unknown): LocalAgentRecord[] { + if (!Array.isArray(value)) throw new LocalAgentDaemonProtocolError("INVALID_RESULT", "Daemon returned an invalid agent list."); + return value.map(decodeAgentRecord); +} + +export function decodeDaemonStatus(value: unknown): LocalAgentDaemonStatus { + const record = asRecord(value); + const state = requiredString(record?.state, "state"); + if (state !== "ready" && state !== "stopping") { + throw new LocalAgentDaemonProtocolError("INVALID_RESULT", "Daemon returned an invalid status."); + } + return { + state, + protocolVersion: requiredInteger(record?.protocolVersion, "protocolVersion"), + pid: requiredInteger(record?.pid, "pid"), + endpoint: requiredString(record?.endpoint, "endpoint"), + startedAt: requiredString(record?.startedAt, "startedAt"), + activeTurns: requiredInteger(record?.activeTurns, "activeTurns"), + runtimeCount: requiredInteger(record?.runtimeCount, "runtimeCount"), + clientConnections: requiredInteger(record?.clientConnections, "clientConnections"), + }; +} + +export function decodeDaemonLogs(value: unknown): string { + if (typeof value !== "string") throw new LocalAgentDaemonProtocolError("INVALID_RESULT", "Daemon returned invalid logs."); + return value; +} + +export class LocalAgentDaemonProtocolError extends Error { + constructor(readonly code: string, message: string) { + super(message); + this.name = "LocalAgentDaemonProtocolError"; + } +} + +function decodeEmptyParams(value: unknown): Record { + if (value === undefined) return {}; + const record = asRecord(value); + if (!record || Object.keys(record).length > 0) { + throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", "This daemon method does not accept parameters."); + } + return {}; +} + +function decodeStartInput(value: unknown): StartLocalAgentInput { + const record = asRecord(value); + return { + target: requiredString(record?.target, "target"), + prompt: requiredString(record?.prompt, "prompt"), + workspaceRoot: requiredString(record?.workspaceRoot, "workspaceRoot"), + workspaceId: optionalString(record?.workspaceId), + model: optionalString(record?.model), + thinking: optionalString(record?.thinking), + writeMode: decodeWriteMode(record?.writeMode), + }; +} + +function decodeContinueInput(value: unknown): { id: string; prompt: string; scope: LocalAgentWorkspaceScope; overrides?: RunOverrides } { + const record = asRecord(value); + const overrides = asRecord(record?.overrides); + return { + id: requiredString(record?.id, "id"), + prompt: requiredString(record?.prompt, "prompt"), + scope: decodeWorkspaceScope(record?.scope), + ...(overrides ? { overrides: { + model: optionalString(overrides.model), + thinking: optionalString(overrides.thinking), + writeMode: decodeWriteMode(overrides.writeMode), + } } : {}), + }; +} + +function decodeWorkspaceScope(value: unknown): LocalAgentWorkspaceScope { + const record = asRecord(value); + if (!record) throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", "Workspace scope is required."); + return { + workspaceId: optionalString(record.workspaceId), + workspaceRoot: requiredString(record.workspaceRoot, "scope.workspaceRoot"), + }; +} + +function decodeListScope(value: unknown): LocalAgentListScope { + if (value === undefined) return {}; + const record = asRecord(value); + if (!record) throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", "List scope must be an object."); + return { + workspaceId: optionalString(record.workspaceId), + workspaceRoot: optionalString(record.workspaceRoot), + }; +} + +function decodeLogsParams(value: unknown): { lines?: number } { + if (value === undefined) return {}; + const record = asRecord(value); + if (!record) throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", "Log options must be an object."); + const lines = record.lines; + if (lines === undefined) return {}; + if (typeof lines !== "number" || !Number.isInteger(lines) || lines < 1 || lines > 10_000) { + throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", "Log lines must be an integer between 1 and 10000."); + } + return { lines }; +} + +function decodeWriteMode(value: unknown): LocalAgentWriteMode | undefined { + if (value === undefined) return undefined; + if (value === "read_only" || value === "allowed" || value === "full_access") return value; + throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", "Invalid write mode."); +} + +function isLocalAgentStatus(value: string): value is LocalAgentStatus { + return value === "starting" || value === "running" || value === "idle" || value === "error" || value === "stopped"; +} + +function requiredString(value: unknown, field: string): string { + const result = optionalString(value); + if (!result) throw new LocalAgentDaemonProtocolError("INVALID_PARAMS", `Missing ${field}.`); + return result; +} + +function requiredInteger(value: unknown, field: string): number { + if (typeof value !== "number" || !Number.isSafeInteger(value)) { + throw new LocalAgentDaemonProtocolError("INVALID_PROTOCOL", `Invalid ${field}.`); + } + return value; +} + +function optionalString(value: unknown): string | undefined { + if (typeof value !== "string") return undefined; + const trimmed = value.trim(); + return trimmed || undefined; +} + +function asRecord(value: unknown): Record | undefined { + if (!value || typeof value !== "object" || Array.isArray(value)) return undefined; + return value as Record; +} + +export function supportedDaemonProtocolVersion(): number { + return LOCAL_AGENT_DAEMON_PROTOCOL_VERSION; +} diff --git a/src/local-agent-daemon.test.ts b/src/local-agent-daemon.test.ts new file mode 100644 index 00000000..04099bcd --- /dev/null +++ b/src/local-agent-daemon.test.ts @@ -0,0 +1,197 @@ +import assert from "node:assert/strict"; +import { existsSync, readFileSync } from "node:fs"; +import { mkdtemp, rm } from "node:fs/promises"; +import { createConnection } from "node:net"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { LocalAgentClient } from "./local-agent-client.js"; +import { LocalAgentDaemon, type LocalAgentDaemonManager } from "./local-agent-daemon.js"; +import type { RunOverrides, StartLocalAgentInput } from "./local-agent-manager.js"; +import type { LocalAgentListScope, LocalAgentRecord } from "./local-agent-store.js"; + +const root = await mkdtemp(join(tmpdir(), "devspace-agentd-test-")); +const record: LocalAgentRecord = { + id: "agt_test", + workspaceRoot: join(root, "project"), + profileName: "reviewer", + provider: "codex", + status: "running", + createdAt: "now", + updatedAt: "now", +}; + +class FakeManager implements LocalAgentDaemonManager { + activeTurnCount = 1; + runtimeCount = 0; + closed = false; + lastInput?: StartLocalAgentInput; + + async start(input: StartLocalAgentInput): Promise { + this.lastInput = input; + return record; + } + + async continue(_agentId: string, _prompt: string, _overrides?: RunOverrides): Promise { + return { ...record, status: "running" }; + } + + get(id: string): LocalAgentRecord | undefined { + return id === record.id ? record : undefined; + } + + list(_scope?: LocalAgentListScope): LocalAgentRecord[] { + return [record]; + } + + async evictIdle(): Promise {} + + async close(): Promise { + this.closed = true; + this.activeTurnCount = 0; + } +} + +const manager = new FakeManager(); +const daemon = new LocalAgentDaemon({ + stateDir: join(root, "state"), + manager, + idleShutdownMs: 60_000, +}); +const client = new LocalAgentClient({ + stateDir: join(root, "state"), + startupTimeoutMs: 2_000, + requestTimeoutMs: 2_000, + spawnDaemon: () => { void daemon.start(); }, +}); + +let shutdownSocket: ReturnType | undefined; +try { + const started = await client.run({ + target: "reviewer", + prompt: "Review this", + workspaceRoot: join(root, "project"), + }); + assert.equal(started.id, record.id); + assert.equal(manager.lastInput?.prompt, "Review this"); + assert.equal((await client.get(record.id, { workspaceRoot: record.workspaceRoot }))?.id, record.id); + assert.equal((await client.list({ workspaceRoot: record.workspaceRoot }))[0]?.id, record.id); + assert.equal((await client.status()).state, "ready"); + + await client.stop(); + await waitFor(() => manager.closed && !existsSync(daemon.paths.socketPath)); +} finally { + await daemon.close(); +} + +const idleStateDir = join(root, "idle-state"); +const idleManager = new FakeManager(); +idleManager.activeTurnCount = 0; +const idleDaemon = new LocalAgentDaemon({ + stateDir: idleStateDir, + manager: idleManager, + idleShutdownMs: 200, + idleCheckIntervalMs: 10, +}); +const idleClient = new LocalAgentClient({ + stateDir: idleStateDir, + startupTimeoutMs: 2_000, + requestTimeoutMs: 2_000, + spawnDaemon: () => { void idleDaemon.start(); }, +}); + +try { + await idleClient.ensureReady(); + await waitFor(() => idleManager.closed && !existsSync(idleDaemon.paths.socketPath)); +} finally { + await idleDaemon.close(); + await rm(root, { recursive: true, force: true }); +} + +const ownershipStateDir = join(root, "ownership-state"); +const ownerManager = new FakeManager(); +const competingManager = new FakeManager(); +const ownerDaemon = new LocalAgentDaemon({ + stateDir: ownershipStateDir, + manager: ownerManager, + idleShutdownMs: 60_000, +}); +const competingDaemon = new LocalAgentDaemon({ + stateDir: ownershipStateDir, + manager: competingManager, + idleShutdownMs: 60_000, +}); + +try { + const startupResults = await Promise.allSettled([ + ownerDaemon.start(), + competingDaemon.start(), + ]); + assert.equal( + startupResults.filter((result) => result.status === "fulfilled").length, + 1, + "only one competing daemon may acquire the state-directory lock", + ); + assert.equal( + startupResults.filter((result) => result.status === "rejected").length, + 1, + ); + const lockBefore = readFileSync(ownerDaemon.paths.lockPath, "utf8"); + const pidBefore = readFileSync(ownerDaemon.paths.pidPath, "utf8"); + assert.notEqual(ownerDaemon.paths.endpoint, ""); + assert.equal(readFileSync(ownerDaemon.paths.lockPath, "utf8"), lockBefore); + assert.equal(readFileSync(ownerDaemon.paths.pidPath, "utf8"), pidBefore); + const ownerClient = new LocalAgentClient({ + stateDir: ownershipStateDir, + spawnDaemon: () => { throw new Error("the winning daemon should already be reachable"); }, + }); + assert.equal((await ownerClient.status()).pid, process.pid); +} finally { + await competingDaemon.close(); + await ownerDaemon.close(); +} + +const socketStateDir = join(root, "socket-state"); +const socketManager = new FakeManager(); +socketManager.activeTurnCount = 0; +const socketDaemon = new LocalAgentDaemon({ + stateDir: socketStateDir, + manager: socketManager, + requestReadTimeoutMs: 30, + shutdownTimeoutMs: 100, + idleShutdownMs: 60_000, +}); + +try { + await socketDaemon.start(); + const idleSocket = createConnection(socketDaemon.paths.endpoint); + await onceSocket(idleSocket, "connect"); + await waitFor(() => socketDaemon.status().clientConnections === 0); + idleSocket.destroy(); + + shutdownSocket = createConnection(socketDaemon.paths.endpoint); + await onceSocket(shutdownSocket, "connect"); + const shutdownSocketClosed = onceSocket(shutdownSocket, "close"); + const startedAt = Date.now(); + await socketDaemon.close(); + await shutdownSocketClosed; + assert.ok(Date.now() - startedAt < 500, "shutdown should destroy idle client sockets before closing the server"); +} finally { + shutdownSocket?.destroy(); + await socketDaemon.close(); + await rm(root, { recursive: true, force: true }); +} + +async function waitFor(check: () => boolean): Promise { + const deadline = Date.now() + 2_000; + while (!check() && Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + assert.equal(check(), true, "condition did not become true before timeout"); +} + +function onceSocket(socket: ReturnType, event: "connect" | "close"): Promise { + return new Promise((resolve, reject) => { + socket.once(event, () => resolve()); + socket.once("error", reject); + }); +} diff --git a/src/local-agent-daemon.ts b/src/local-agent-daemon.ts new file mode 100644 index 00000000..b80b3749 --- /dev/null +++ b/src/local-agent-daemon.ts @@ -0,0 +1,391 @@ +import { timingSafeEqual } from "node:crypto"; +import { appendFileSync, chmodSync, readFileSync, rmSync } from "node:fs"; +import { createServer, type Server as NetServer, type Socket } from "node:net"; +import { + LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + LocalAgentDaemonAlreadyRunningError, + LocalAgentDaemonLock, + ensureLocalAgentDaemonStateDir, + ensureLocalAgentDaemonSecret, + localAgentDaemonPaths, + removeLocalAgentDaemonFiles, + type LocalAgentDaemonPaths, +} from "./local-agent-daemon-lifecycle.js"; +import { + decodeLocalAgentDaemonRequest, + encodeLocalAgentDaemonResponse, + type LocalAgentDaemonRequest, + type LocalAgentDaemonResponse, + type LocalAgentDaemonStatus, + LocalAgentDaemonProtocolError, +} from "./local-agent-daemon-protocol.js"; +import { LocalAgentConflictError, type RunOverrides, type StartLocalAgentInput } from "./local-agent-manager.js"; +import type { LocalAgentListScope, LocalAgentRecord, LocalAgentWorkspaceScope } from "./local-agent-store.js"; + +const MAX_REQUEST_BYTES = 512 * 1024; +const DEFAULT_DAEMON_IDLE_SHUTDOWN_MS = 30_000; +const DEFAULT_IDLE_CHECK_INTERVAL_MS = 1_000; +const DEFAULT_REQUEST_READ_TIMEOUT_MS = 5_000; +const DEFAULT_DAEMON_SHUTDOWN_TIMEOUT_MS = 10_000; + +export interface LocalAgentDaemonManager { + start(input: StartLocalAgentInput): Promise; + continue(agentId: string, prompt: string, overrides?: RunOverrides, scope?: LocalAgentWorkspaceScope): Promise; + get(agentId: string, scope?: LocalAgentWorkspaceScope): LocalAgentRecord | undefined; + list(scope?: LocalAgentListScope): LocalAgentRecord[]; + evictIdle(now?: number): Promise; + close(): Promise; + readonly activeTurnCount: number; + readonly runtimeCount: number; +} + +export interface LocalAgentDaemonOptions { + stateDir: string; + manager: LocalAgentDaemonManager; + idleShutdownMs?: number; + idleCheckIntervalMs?: number; + requestReadTimeoutMs?: number; + shutdownTimeoutMs?: number; + now?: () => number; + paths?: LocalAgentDaemonPaths; + onLockAcquired?: () => void | Promise; + onClosed?: () => void; +} + +export class LocalAgentDaemon { + readonly paths: LocalAgentDaemonPaths; + private readonly manager: LocalAgentDaemonManager; + private readonly lock: LocalAgentDaemonLock; + private readonly idleShutdownMs: number; + private readonly idleCheckIntervalMs: number; + private readonly requestReadTimeoutMs: number; + private readonly shutdownTimeoutMs: number; + private readonly now: () => number; + private readonly onLockAcquired?: () => void | Promise; + private readonly onClosed?: () => void; + private readonly sockets = new Set(); + private server?: NetServer; + private idleTimer?: NodeJS.Timeout; + private idleSince?: number; + private closePromise?: Promise; + private startedAt?: string; + private accepting = false; + private stopping = false; + private authToken?: string; + private ownsLock = false; + + constructor(options: LocalAgentDaemonOptions) { + this.paths = options.paths ?? localAgentDaemonPaths(options.stateDir); + this.manager = options.manager; + this.lock = new LocalAgentDaemonLock(this.paths); + this.idleShutdownMs = options.idleShutdownMs ?? DEFAULT_DAEMON_IDLE_SHUTDOWN_MS; + this.idleCheckIntervalMs = options.idleCheckIntervalMs ?? DEFAULT_IDLE_CHECK_INTERVAL_MS; + this.requestReadTimeoutMs = options.requestReadTimeoutMs ?? DEFAULT_REQUEST_READ_TIMEOUT_MS; + this.shutdownTimeoutMs = options.shutdownTimeoutMs ?? DEFAULT_DAEMON_SHUTDOWN_TIMEOUT_MS; + this.now = options.now ?? Date.now; + this.onLockAcquired = options.onLockAcquired; + this.onClosed = options.onClosed; + if (!Number.isFinite(this.idleShutdownMs) || this.idleShutdownMs < 0) { + throw new Error("Agent daemon idle shutdown must be a non-negative finite duration."); + } + if (!Number.isFinite(this.requestReadTimeoutMs) || this.requestReadTimeoutMs <= 0) { + throw new Error("Agent daemon request read timeout must be a positive finite duration."); + } + if (!Number.isFinite(this.shutdownTimeoutMs) || this.shutdownTimeoutMs < 0) { + throw new Error("Agent daemon shutdown timeout must be a non-negative finite duration."); + } + } + + async start(): Promise { + if (this.server) return this.status(); + ensureLocalAgentDaemonStateDir(this.paths.stateDir); + let lockAcquired = false; + try { + this.lock.acquire(); + lockAcquired = true; + this.ownsLock = true; + this.authToken = ensureLocalAgentDaemonSecret(this.paths); + await this.onLockAcquired?.(); + if (process.platform !== "win32") rmSync(this.paths.socketPath, { force: true }); + const server = createServer((socket) => this.handleConnection(socket)); + this.server = server; + await listen(server, this.paths.endpoint); + if (process.platform !== "win32") chmodSync(this.paths.socketPath, 0o600); + this.startedAt = new Date(this.now()).toISOString(); + this.accepting = true; + this.stopping = false; + this.idleTimer = setInterval(() => { + void this.maintainIdle().catch((error) => { + writeLocalAgentDaemonLog(this.paths, "warn", "daemon_idle_check_failed", { + error: errorMessage(error), + }); + }); + }, this.idleCheckIntervalMs); + this.idleTimer.unref(); + writeLocalAgentDaemonLog(this.paths, "info", "daemon_started", { pid: process.pid }); + return this.status(); + } catch (error) { + this.server = undefined; + this.authToken = undefined; + if (lockAcquired) { + this.lock.release(); + this.ownsLock = false; + removeLocalAgentDaemonFiles(this.paths); + } + if (error instanceof LocalAgentDaemonAlreadyRunningError) throw error; + throw error; + } + } + + status(): LocalAgentDaemonStatus { + if (!this.startedAt) throw new Error("Local agent daemon is not started."); + return { + state: this.stopping ? "stopping" : "ready", + protocolVersion: LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + pid: process.pid, + endpoint: this.paths.endpoint, + startedAt: this.startedAt, + activeTurns: this.manager.activeTurnCount, + runtimeCount: this.manager.runtimeCount, + clientConnections: this.sockets.size, + }; + } + + async close(): Promise { + if (this.closePromise) return this.closePromise; + if (!this.ownsLock && !this.server) return; + this.accepting = false; + this.stopping = true; + if (this.idleTimer) clearInterval(this.idleTimer); + this.closePromise = (async () => { + writeLocalAgentDaemonLog(this.paths, "info", "daemon_stopping", { + activeTurns: this.manager.activeTurnCount, + runtimeCount: this.manager.runtimeCount, + }); + for (const socket of this.sockets) socket.destroy(); + this.sockets.clear(); + const [serverResult, managerResult] = await Promise.allSettled([ + withTimeout(closeServer(this.server), this.shutdownTimeoutMs, "daemon socket shutdown"), + withTimeout(this.manager.close(), this.shutdownTimeoutMs, "daemon manager shutdown"), + ]); + if (serverResult.status === "rejected") { + writeLocalAgentDaemonLog(this.paths, "warn", "daemon_socket_close_failed", { + error: errorMessage(serverResult.reason), + }); + } + if (managerResult.status === "rejected") { + writeLocalAgentDaemonLog(this.paths, "warn", "daemon_manager_close_failed", { + error: errorMessage(managerResult.reason), + }); + } + removeLocalAgentDaemonFiles(this.paths); + this.lock.release(); + writeLocalAgentDaemonLog(this.paths, "info", "daemon_stopped", {}); + this.server = undefined; + this.authToken = undefined; + this.onClosed?.(); + })(); + return this.closePromise; + } + + private handleConnection(socket: Socket): void { + this.sockets.add(socket); + socket.setEncoding("utf8"); + let buffer = ""; + let handled = false; + const requestTimer = setTimeout(() => { + if (handled) return; + handled = true; + this.writeError(socket, "", "REQUEST_TIMEOUT", "Timed out waiting for a complete daemon request."); + socket.destroy(); + }, this.requestReadTimeoutMs); + requestTimer.unref(); + socket.on("data", (chunk: string | Buffer) => { + if (handled) return; + buffer += chunk.toString(); + if (Buffer.byteLength(buffer, "utf8") > MAX_REQUEST_BYTES) { + handled = true; + this.writeError(socket, "", "REQUEST_TOO_LARGE", "Daemon request is too large."); + return; + } + const newline = buffer.indexOf("\n"); + if (newline === -1) return; + handled = true; + clearTimeout(requestTimer); + const line = buffer.slice(0, newline); + void this.handleLine(socket, line); + }); + socket.on("error", () => undefined); + socket.on("close", () => this.sockets.delete(socket)); + socket.on("error", () => clearTimeout(requestTimer)); + } + + private async handleLine(socket: Socket, line: string): Promise { + let requestId = ""; + try { + const parsed: unknown = JSON.parse(line); + requestId = readRequestId(parsed); + const request = decodeLocalAgentDaemonRequest(parsed); + const response = await this.dispatch(request); + socket.end(encodeLocalAgentDaemonResponse({ + requestId: request.requestId, + protocolVersion: LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + ok: true, + result: response, + })); + if (request.method === "daemon.stop") setImmediate(() => { void this.close(); }); + } catch (error) { + this.writeError(socket, requestId, errorCode(error), errorMessage(error)); + } + } + + private async dispatch(request: LocalAgentDaemonRequest): Promise { + if (request.protocolVersion !== LOCAL_AGENT_DAEMON_PROTOCOL_VERSION) { + throw new LocalAgentDaemonProtocolError( + "PROTOCOL_MISMATCH", + `Unsupported daemon protocol version ${request.protocolVersion}; expected ${LOCAL_AGENT_DAEMON_PROTOCOL_VERSION}.`, + ); + } + this.assertAuthenticated(request.authToken); + if (!this.accepting && request.method !== "hello" && request.method !== "daemon.status") { + throw new Error("Local agent daemon is stopping."); + } + + switch (request.method) { + case "hello": + return this.status(); + case "agent.start": + return this.manager.start(request.params); + case "agent.continue": + return this.manager.continue(request.params.id, request.params.prompt, request.params.overrides, request.params.scope); + case "agent.get": + return this.manager.get(request.params.id, request.params.scope) ?? null; + case "agent.list": + return this.manager.list(request.params); + case "daemon.status": + return this.status(); + case "daemon.stop": + this.stopping = true; + this.accepting = false; + return this.status(); + case "daemon.logs": + return readLocalAgentDaemonLogs(this.paths, request.params.lines); + } + } + + private writeError(socket: Socket, requestId: string, code: string, message: string): void { + socket.end(encodeLocalAgentDaemonResponse({ + requestId, + protocolVersion: LOCAL_AGENT_DAEMON_PROTOCOL_VERSION, + ok: false, + error: { code, message }, + }), () => socket.destroy()); + } + + private assertAuthenticated(authToken: string): void { + const expected = this.authToken; + if (!expected || !safeEqual(authToken, expected)) { + throw new LocalAgentDaemonProtocolError("UNAUTHORIZED", "Invalid local agent daemon credentials."); + } + } + + private async maintainIdle(): Promise { + await this.manager.evictIdle(this.now()); + if (this.stopping || this.manager.activeTurnCount > 0 || this.manager.runtimeCount > 0 || this.sockets.size > 0) { + this.idleSince = undefined; + return; + } + const now = this.now(); + this.idleSince ??= now; + if (now - this.idleSince >= this.idleShutdownMs) await this.close(); + } +} + +async function listen(server: NetServer, endpoint: string): Promise { + 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(endpoint); + }); +} + +async function closeServer(server: NetServer | undefined): Promise { + if (!server) return; + if (!server.listening) return; + await new Promise((resolve, reject) => { + server.close((error) => error ? reject(error) : resolve()); + }); +} + +async function withTimeout(promise: Promise, timeoutMs: number, operation: string): Promise { + if (timeoutMs === 0) { + throw new Error(`${operation} timed out.`); + } + let timer: NodeJS.Timeout | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(`${operation} timed out.`)), timeoutMs); + timer.unref(); + }), + ]); + } finally { + if (timer) clearTimeout(timer); + } +} + +function safeEqual(actual: string, expected: string): boolean { + const actualBuffer = Buffer.from(actual); + const expectedBuffer = Buffer.from(expected); + return actualBuffer.length === expectedBuffer.length && timingSafeEqual(actualBuffer, expectedBuffer); +} + +function readRequestId(value: unknown): string { + if (!value || typeof value !== "object" || Array.isArray(value)) return ""; + const requestId = (value as Record).requestId; + return typeof requestId === "string" ? requestId : ""; +} + +function errorCode(error: unknown): string { + if (error instanceof LocalAgentDaemonProtocolError) return error.code; + if (error instanceof LocalAgentConflictError) return "CONFLICT"; + if (errorMessage(error).includes("is stopping")) return "DAEMON_STOPPING"; + return "AGENT_ERROR"; +} + +export function writeLocalAgentDaemonLog( + paths: LocalAgentDaemonPaths, + level: "info" | "warn" | "error", + event: string, + fields: Record, +): void { + try { + ensureLocalAgentDaemonStateDir(paths.stateDir); + appendFileSync(paths.logPath, `${JSON.stringify({ at: new Date().toISOString(), level, event, ...fields })}\n`, { mode: 0o600 }); + chmodSync(paths.logPath, 0o600); + } catch { + // Diagnostics must never break agent execution or shutdown. + } +} + +export function readLocalAgentDaemonLogs(paths: LocalAgentDaemonPaths, lines = 200): string { + try { + const content = readFileSync(paths.logPath, "utf8"); + return content.split(/\r?\n/).filter(Boolean).slice(-Math.max(1, lines)).join("\n"); + } catch { + return ""; + } +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/src/local-agent-manager.test.ts b/src/local-agent-manager.test.ts new file mode 100644 index 00000000..e08db763 --- /dev/null +++ b/src/local-agent-manager.test.ts @@ -0,0 +1,183 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { LocalAgentConflictError, LocalAgentManager } from "./local-agent-manager.js"; +import type { LocalAgentProfile } from "./local-agent-profiles.js"; +import type { + LocalAgentDriver, + LocalAgentRunInput, + LocalAgentRunResult, + LocalAgentRuntime, + LocalAgentRuntimeContext, +} from "./local-agent-runtime.js"; +import { LocalAgentRuntimePool } from "./local-agent-runtime-pool.js"; +import { LocalAgentStore } from "./local-agent-store.js"; + +const root = await mkdtemp(join(tmpdir(), "devspace-agent-manager-test-")); +const stateDir = join(root, "state"); +const profile: LocalAgentProfile = { + name: "reviewer", + description: "Test reviewer", + provider: "codex", + filePath: join(root, "reviewer.md"), + body: "Review only.", + disabled: false, +}; + +class FakeRuntime implements LocalAgentRuntime { + readonly provider = "codex" as const; + readonly inputs: LocalAgentRunInput[] = []; + closed = false; + private releaseHold: (() => void) | undefined; + + async run(input: LocalAgentRunInput, callbacks?: { onSessionId?: (id: string) => void | Promise }): Promise { + this.inputs.push(input); + if (input.prompt.includes("early-fail")) { + await callbacks?.onSessionId?.("thread_early"); + throw new Error("provider failed after session creation"); + } + if (input.prompt.includes("fail")) throw new Error("provider failed"); + if (input.prompt.includes("hold")) { + await new Promise((resolve) => { this.releaseHold = resolve; }); + } + return { + provider: this.provider, + providerSessionId: "thread_test", + finalResponse: `response:${input.prompt}`, + items: [], + }; + } + + release(): void { + this.releaseHold?.(); + this.releaseHold = undefined; + } + + releaseSession(): Promise { + return Promise.resolve(); + } + + isAlive(): boolean { + return !this.closed; + } + + async close(): Promise { + this.closed = true; + this.release(); + } +} + +const runtimes = new Map(); +const driver: LocalAgentDriver = { + provider: "codex", + runtimeKey: (context: LocalAgentRuntimeContext) => context.agentId, + createRuntime: async (context) => { + const runtime = new FakeRuntime(); + runtimes.set(context.agentId, runtime); + return runtime; + }, +}; + +const store = new LocalAgentStore(stateDir); +const stale = store.create({ + workspaceRoot: root, + profileName: "reviewer", + provider: "codex", +}); +store.update(stale.id, { status: "running", latestResponse: "previous response" }); + +const manager = new LocalAgentManager({ + store, + drivers: [driver], + pool: new LocalAgentRuntimePool(), + loadProfiles: async () => [profile], + allowedRoots: [root], +}); + +await assert.rejects( + manager.start({ target: "reviewer", prompt: "outside", workspaceRoot: join(tmpdir(), "outside") }), + /outside allowed roots/, +); + +assert.equal(manager.get(stale.id)?.status, "running"); +manager.reconcileActiveRuns(); +assert.equal(manager.get(stale.id)?.status, "error"); +assert.equal(manager.get(stale.id)?.latestResponse, "previous response"); +assert.equal( + manager.get(stale.id)?.error, + "DevSpace restarted while this agent turn was running.", +); + +const first = await manager.start({ + target: "reviewer", + prompt: "hold", + workspaceRoot: root, +}); +assert.equal(first.status, "running"); +await waitFor(() => runtimes.get(first.id)?.inputs.length === 1); +await assert.rejects( + () => manager.continue(first.id, "another prompt"), + (error: unknown) => error instanceof LocalAgentConflictError && error.agentId === first.id, +); + +runtimes.get(first.id)!.release(); +await waitFor(() => manager.get(first.id)?.status === "idle"); +assert.equal(manager.get(first.id)?.providerSessionId, "thread_test"); +assert.match(manager.get(first.id)?.latestResponse ?? "", /Task:\nhold/); + +const continued = await manager.continue(first.id, "continue"); +assert.equal(continued.status, "running"); +await waitFor(() => manager.get(first.id)?.status === "idle"); + +const second = await manager.start({ + target: "reviewer", + prompt: "second agent", + workspaceRoot: root, +}); +await waitFor(() => manager.get(second.id)?.status === "idle"); +assert.notEqual(first.id, second.id); +assert.equal(runtimes.size, 2, "different agents receive independent logical runtimes"); + +const failed = await manager.start({ + target: "reviewer", + prompt: "fail", + workspaceRoot: root, +}); +await waitFor(() => manager.get(failed.id)?.status === "error"); +assert.equal(manager.get(failed.id)?.error, "provider failed"); + +const earlyFailure = await manager.start({ + target: "reviewer", + prompt: "early-fail", + workspaceRoot: root, +}); +await waitFor(() => manager.get(earlyFailure.id)?.status === "error"); +assert.equal(manager.get(earlyFailure.id)?.providerSessionId, "thread_early"); + +await assert.rejects( + () => manager.continue(first.id, "wrong workspace", {}, { workspaceRoot: join(root, "other") }), + /different workspace/, +); + +const shuttingDown = await manager.start({ + target: "reviewer", + prompt: "hold during shutdown", + workspaceRoot: root, +}); +await waitFor(() => runtimes.get(shuttingDown.id)?.inputs.length === 1); +const closing = manager.close(); +await new Promise((resolve) => setImmediate(resolve)); +assert.equal(runtimes.get(shuttingDown.id)?.closed, true); +await closing; + +await manager.close(); +await rm(root, { recursive: true, force: true }); + +async function waitFor(check: () => boolean): Promise { + const deadline = Date.now() + 2_000; + while (!check() && Date.now() < deadline) { + await new Promise((resolve) => setImmediate(resolve)); + } + assert.equal(check(), true, "condition did not become true before timeout"); +} diff --git a/src/local-agent-manager.ts b/src/local-agent-manager.ts new file mode 100644 index 00000000..c6ac0449 --- /dev/null +++ b/src/local-agent-manager.ts @@ -0,0 +1,335 @@ +import { + type LocalAgentProfile, + type LocalAgentProvider, + isLocalAgentProvider, +} from "./local-agent-profiles.js"; +import { + resolveLocalAgentTarget, +} from "./local-agent-targets.js"; +import { + type LocalAgentListScope, + type LocalAgentRecord, + type LocalAgentStore, + type LocalAgentWorkspaceScope, +} from "./local-agent-store.js"; +import { + type LocalAgentDriver, + type LocalAgentRunCallbacks, + type LocalAgentRunInput, + type LocalAgentRuntimeContext, + type LocalAgentWriteMode, +} from "./local-agent-runtime.js"; +import { LocalAgentRuntimePool } from "./local-agent-runtime-pool.js"; +import { assertAllowedPath } from "./roots.js"; + +export interface StartLocalAgentInput { + target: string; + prompt: string; + workspaceRoot: string; + workspaceId?: string; + model?: string; + thinking?: string; + writeMode?: LocalAgentWriteMode; +} + +export interface RunOverrides { + model?: string; + thinking?: string; + writeMode?: LocalAgentWriteMode; +} + +export interface LocalAgentManagerLogger { + (level: "info" | "warn" | "error", event: string, fields: Record): void; +} + +export interface LocalAgentManagerOptions { + store: LocalAgentStore; + drivers: readonly LocalAgentDriver[]; + pool: LocalAgentRuntimePool; + loadProfiles: (workspaceRoot: string) => Promise; + agentDir?: string; + allowedRoots?: readonly string[]; + logger?: LocalAgentManagerLogger; +} + +export class LocalAgentConflictError extends Error { + readonly code = "CONFLICT" as const; + + constructor(readonly agentId: string) { + super(`Agent ${agentId} already has a running turn.`); + this.name = "LocalAgentConflictError"; + } +} + +/** + * Owns one durable DevSpace agent's turn lifecycle. Provider runtimes remain + * below this seam; this class only translates records into provider inputs and + * persists the result. + */ +export class LocalAgentManager { + private readonly store: LocalAgentStore; + private readonly drivers = new Map(); + private readonly pool: LocalAgentRuntimePool; + private readonly loadProfiles: (workspaceRoot: string) => Promise; + private readonly agentDir?: string; + private readonly allowedRoots?: readonly string[]; + private readonly logger?: LocalAgentManagerLogger; + private readonly activeTurns = new Map>(); + private accepting = true; + private closePromise?: Promise; + + constructor(options: LocalAgentManagerOptions) { + this.store = options.store; + for (const driver of options.drivers) this.drivers.set(driver.provider, driver); + this.pool = options.pool; + this.loadProfiles = options.loadProfiles; + this.agentDir = options.agentDir; + this.allowedRoots = options.allowedRoots; + this.logger = options.logger; + } + + reconcileActiveRuns(message?: string): number { + return this.store.reconcileActiveRuns(message); + } + + async start(input: StartLocalAgentInput): Promise { + this.assertAccepting(); + const workspaceRoot = this.authorizeWorkspace(input.workspaceRoot); + const profiles = await this.loadProfiles(workspaceRoot); + const target = resolveLocalAgentTarget(input.target, profiles, input.model, input.thinking); + if (!target) { + throw new Error(`Unknown subagent profile or provider: ${input.target}`); + } + this.assertDriver(target.provider); + + const record = this.store.create({ + workspaceId: input.workspaceId, + workspaceRoot, + profileName: target.name, + provider: target.provider, + model: target.model, + thinking: target.thinking, + }); + return this.begin(record, input.prompt, { + model: target.model, + thinking: target.thinking, + writeMode: input.writeMode, + }); + } + + async continue( + agentId: string, + prompt: string, + overrides: RunOverrides = {}, + scope?: LocalAgentWorkspaceScope, + ): Promise { + this.assertAccepting(); + const record = this.store.getById(agentId); + if (!record) throw new Error(`Unknown subagent id: ${agentId}`); + if (scope) this.assertAgentWorkspace(record, scope); + this.assertDriver(record.provider); + return this.begin(record, prompt, overrides); + } + + get(agentId: string, scope?: LocalAgentWorkspaceScope): LocalAgentRecord | undefined { + const record = this.store.getById(agentId); + if (record && scope) this.assertAgentWorkspace(record, scope); + return record; + } + + list(scope: LocalAgentListScope = {}): LocalAgentRecord[] { + return this.store.list(scope.workspaceRoot + ? { ...scope, workspaceRoot: this.authorizeWorkspace(scope.workspaceRoot) } + : scope); + } + + async close(): Promise { + if (this.closePromise) return this.closePromise; + this.accepting = false; + const turns = Array.from(this.activeTurns.values()); + this.closePromise = (async () => { + // Closing pooled runtimes is what interrupts provider turns. Waiting for + // those turns first can strand a provider process indefinitely. + await this.pool.close(); + const turnResults = await Promise.allSettled(turns); + for (const result of turnResults) { + if (result.status === "rejected") { + this.log("warn", "local_agent_close_failed", { error: errorMessage(result.reason) }); + } + } + this.store.close(); + })(); + return this.closePromise; + } + + get activeTurnCount(): number { + return this.activeTurns.size; + } + + get runtimeCount(): number { + return this.pool.size; + } + + async evictIdle(now?: number): Promise { + await this.pool.evictIdle(now); + } + + private begin( + record: LocalAgentRecord, + prompt: string, + overrides: RunOverrides, + ): LocalAgentRecord { + if (this.activeTurns.has(record.id)) { + throw new LocalAgentConflictError(record.id); + } + + const updated = this.store.update(record.id, { + status: "running", + model: overrides.model ?? record.model, + thinking: overrides.thinking ?? record.thinking, + latestResponse: undefined, + error: undefined, + }); + // Defer invocation until after the tracking entry is visible. This keeps + // cleanup correct even if runTurn later gains a synchronous completion path. + const turn = Promise.resolve().then(() => this.runTurn(updated, prompt, overrides)); + this.activeTurns.set(record.id, turn); + void turn.catch(() => undefined); + return updated; + } + + private async runTurn( + record: LocalAgentRecord, + prompt: string, + overrides: RunOverrides, + ): Promise { + const startedAt = Date.now(); + this.log("info", "agent_run_started", { + provider: record.provider, + agentId: record.id, + providerSessionIdPrefix: record.providerSessionId?.slice(0, 8), + }); + try { + const profiles = await this.loadProfiles(record.workspaceRoot); + const profile = profiles.find((candidate) => candidate.name === record.profileName); + const input = this.buildRunInput(record, profile, prompt, overrides); + const driver = this.assertDriver(record.provider); + const context: LocalAgentRuntimeContext = { + agentId: record.id, + provider: driver.provider, + workspace: record.workspaceRoot, + providerSessionId: record.providerSessionId, + writeMode: input.writeMode, + model: input.model, + thinking: input.thinking, + agentDir: this.agentDir, + }; + const callbacks: LocalAgentRunCallbacks = { + onSessionId: (providerSessionId) => { + const current = this.store.getById(record.id); + if (!current || current.providerSessionId === providerSessionId) return; + this.store.update(record.id, { providerSessionId }); + }, + }; + const result = await this.pool.run(driver, context, input, callbacks); + const current = this.store.getById(record.id); + if (!current) return; + const updated = this.store.update(record.id, { + providerSessionId: result.providerSessionId ?? current.providerSessionId, + status: "idle", + latestResponse: result.finalResponse, + error: undefined, + }); + this.log("info", "agent_run_completed", { + provider: updated.provider, + agentId: updated.id, + providerSessionIdPrefix: updated.providerSessionId?.slice(0, 8), + durationMs: Math.max(0, Date.now() - startedAt), + }); + } catch (error) { + const current = this.store.getById(record.id); + if (current) { + this.store.update(record.id, { + status: "error", + error: errorMessage(error), + }); + } + this.log("error", "agent_run_failed", { + provider: record.provider, + agentId: record.id, + providerSessionIdPrefix: record.providerSessionId?.slice(0, 8), + durationMs: Math.max(0, Date.now() - startedAt), + error: errorMessage(error), + }); + throw error; + } finally { + this.activeTurns.delete(record.id); + } + } + + private buildRunInput( + record: LocalAgentRecord, + profile: LocalAgentProfile | undefined, + prompt: string, + overrides: RunOverrides, + ): LocalAgentRunInput { + const isRawProvider = record.profileName === record.provider; + if (!profile && !isRawProvider) { + throw new Error(`Subagent profile not found: ${record.profileName}`); + } + const body = profile?.body.trim(); + const fullPrompt = body ? `${body}\n\nTask:\n${prompt}` : prompt; + return { + prompt: fullPrompt, + workspace: record.workspaceRoot, + providerSessionId: record.providerSessionId, + writeMode: overrides.writeMode ?? "allowed", + model: record.model ?? profile?.model, + thinking: record.thinking ?? profile?.thinking, + }; + } + + private assertDriver(provider: string): LocalAgentDriver { + if (!isLocalAgentProvider(provider)) { + throw new Error(`No local agent driver is configured for provider: ${provider}`); + } + const driver = this.drivers.get(provider); + if (!driver) throw new Error(`No local agent driver is configured for provider: ${provider}`); + return driver; + } + + private assertAccepting(): void { + if (!this.accepting) throw new Error("Local agent manager is closed."); + } + + private authorizeWorkspace(workspaceRoot: string): string { + if (!this.allowedRoots) return workspaceRoot; + return assertAllowedPath(workspaceRoot, [...this.allowedRoots]); + } + + private assertAgentWorkspace(record: LocalAgentRecord, scope: LocalAgentWorkspaceScope): void { + const workspaceRoot = this.authorizeWorkspace(scope.workspaceRoot); + if (workspaceRoot !== record.workspaceRoot) { + throw new Error(`Subagent ${record.id} belongs to a different workspace.`); + } + if (record.workspaceId && record.workspaceId !== scope.workspaceId) { + throw new Error(`Subagent ${record.id} belongs to a different workspace.`); + } + } + + private log( + level: "info" | "warn" | "error", + event: string, + fields: Record, + ): void { + this.logger?.(level, event, fields); + } +} + +export function createLocalAgentManager(options: LocalAgentManagerOptions): LocalAgentManager { + return new LocalAgentManager(options); +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/src/local-agent-runtime-pool.ts b/src/local-agent-runtime-pool.ts new file mode 100644 index 00000000..fa6e7621 --- /dev/null +++ b/src/local-agent-runtime-pool.ts @@ -0,0 +1,428 @@ +import { createHash } from "node:crypto"; +import type { + LocalAgentDriver, + LocalAgentRunCallbacks, + LocalAgentRunInput, + LocalAgentRunResult, + LocalAgentRuntime, + LocalAgentRuntimeContext, +} from "./local-agent-runtime.js"; + +const DEFAULT_IDLE_TIMEOUT_MS = 5 * 60_000; +const DEFAULT_SESSION_IDLE_TIMEOUT_MS = 60_000; + +export interface LocalAgentRuntimePoolLogger { + (level: "info" | "warn" | "error", event: string, fields: Record): void; +} + +interface RuntimeEntry { + readonly key: string; + readonly driver: LocalAgentDriver; + readonly idleTimeoutMs: number; + readonly sessionIdleTimeoutMs: number; + readonly createPromise: Promise; + runtime?: LocalAgentRuntime; + activeRuns: number; + lastUsedAt: number; + closePromise?: Promise; + idleTimer?: NodeJS.Timeout; + closing: boolean; + readonly sessions: Map; + readonly activeRunWaiters: Set<() => void>; +} + +interface SessionEntry { + activeRuns: number; + lastUsedAt: number; + releasePromise?: Promise; +} + +export interface LocalAgentRuntimePoolOptions { + now?: () => number; + logger?: LocalAgentRuntimePoolLogger; + sessionIdleTimeoutMs?: number; +} + +/** + * Owns live provider resources, not logical agent identity. Acquisition is + * single-flight per runtime key and an entry is removed before its close + * begins, so a new caller can never race with a closing runtime. + */ +export class LocalAgentRuntimePool { + private readonly entries = new Map(); + private readonly now: () => number; + private readonly logger?: LocalAgentRuntimePoolLogger; + private readonly sessionIdleTimeoutMs: number; + private closing = false; + private closePromise?: Promise; + + constructor(options: LocalAgentRuntimePoolOptions = {}) { + this.now = options.now ?? Date.now; + this.logger = options.logger; + this.sessionIdleTimeoutMs = options.sessionIdleTimeoutMs ?? DEFAULT_SESSION_IDLE_TIMEOUT_MS; + if (!Number.isFinite(this.sessionIdleTimeoutMs) || this.sessionIdleTimeoutMs < 0) { + throw new Error("Local agent session idle timeout must be a non-negative finite duration."); + } + } + + async run( + driver: LocalAgentDriver, + context: LocalAgentRuntimeContext, + input: LocalAgentRunInput, + inputCallbacks?: LocalAgentRunCallbacks, + ): Promise { + if (this.closing) throw new Error("Local agent runtime pool is closed."); + + let entry = await this.acquire(driver, context); + let runtime = entry.runtime; + if (!runtime) throw new Error("Local agent runtime was created without a runtime."); + if (!runtime.isAlive()) { + await this.removeAndClose(entry, "runtime_not_alive"); + entry = await this.acquire(driver, context); + runtime = entry.runtime; + if (!runtime || !runtime.isAlive()) { + await this.removeAndClose(entry, "runtime_not_alive"); + throw new Error("Local agent runtime exited during startup."); + } + } + + this.clearIdleTimer(entry); + entry.activeRuns += 1; + const sessionIds = new Set(); + const reserveSession = async (providerSessionId: string): Promise => { + if (!providerSessionId || sessionIds.has(providerSessionId)) return; + while (true) { + const existing = entry.sessions.get(providerSessionId); + if (existing?.releasePromise) { + await existing.releasePromise; + continue; + } + if (entry.closing) throw new Error("Local agent runtime is closing."); + const session = existing ?? { activeRuns: 0, lastUsedAt: this.now() }; + sessionIds.add(providerSessionId); + session.activeRuns += 1; + session.lastUsedAt = this.now(); + entry.sessions.set(providerSessionId, session); + return; + } + }; + const callbacks: LocalAgentRunCallbacks = { + onSessionId: async (providerSessionId) => { + await reserveSession(providerSessionId); + await inputCallbacks?.onSessionId?.(providerSessionId); + }, + }; + const startedAt = this.now(); + try { + await reserveSession(input.providerSessionId ?? ""); + const result = await runtime.run(input, callbacks); + await reserveSession(result.providerSessionId ?? ""); + return result; + } catch (error) { + if (!runtime.isAlive()) { + try { + await this.removeAndClose(entry, "runtime_crashed"); + } catch (cleanupError) { + this.log("warn", "harness_runtime_close_failed", { + provider: driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + reason: "runtime_crashed", + error: errorMessage(cleanupError), + }); + } + this.log("warn", "harness_runtime_crashed", { + provider: driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + agentId: context.agentId, + providerSessionIdPrefix: input.providerSessionId?.slice(0, 8), + durationMs: Math.max(0, Math.round(this.now() - startedAt)), + error: errorMessage(error), + }); + } + throw error; + } finally { + for (const providerSessionId of sessionIds) { + const session = entry.sessions.get(providerSessionId); + if (!session) continue; + session.activeRuns = Math.max(0, session.activeRuns - 1); + session.lastUsedAt = this.now(); + } + entry.activeRuns -= 1; + if (entry.activeRuns === 0) { + for (const resolve of entry.activeRunWaiters) resolve(); + entry.activeRunWaiters.clear(); + } + entry.lastUsedAt = this.now(); + if (entry.activeRuns === 0 && !entry.closing) this.scheduleIdleClose(entry); + } + } + + /** Evict entries whose runtime has been idle beyond their driver's TTL. */ + async evictIdle(now = this.now()): Promise { + const evictions: Promise[] = []; + for (const entry of this.entries.values()) { + if (entry.closing || !entry.runtime) continue; + await this.releaseIdleSessions(entry, now); + if (entry.activeRuns === 0 && now - entry.lastUsedAt >= entry.idleTimeoutMs) { + evictions.push(this.removeAndClose(entry, "idle_timeout")); + } + } + await Promise.all(evictions); + } + + async close(): Promise { + if (this.closePromise) return this.closePromise; + this.closing = true; + const entries = Array.from(this.entries.values()); + this.entries.clear(); + this.closePromise = Promise.allSettled(entries.map((entry) => this.closeEntry(entry, "server_shutdown"))) + .then(() => undefined); + return this.closePromise; + } + + get size(): number { + return this.entries.size; + } + + private async acquire( + driver: LocalAgentDriver, + context: LocalAgentRuntimeContext, + ): Promise { + const key = driver.runtimeKey(context); + while (true) { + const existing = this.entries.get(key); + if (existing && !existing.closing) { + if (!existing.runtime || existing.runtime.isAlive()) { + this.clearIdleTimer(existing); + if (existing.runtime) { + this.log("info", "harness_runtime_reused", { + provider: driver.provider, + runtimeKeyHash: hashRuntimeKey(key), + agentId: context.agentId, + }); + } + await existing.createPromise; + if ( + !this.closing && + !existing.closing && + this.entries.get(key) === existing && + existing.runtime?.isAlive() + ) { + return existing; + } + } + await this.removeAndClose(existing, "runtime_not_alive"); + continue; + } + + if (this.closing) throw new Error("Local agent runtime pool is closed."); + + let entry!: RuntimeEntry; + const createPromise = Promise.resolve() + .then(() => driver.createRuntime(context)) + .then((runtime) => { + entry.runtime = runtime; + entry.lastUsedAt = this.now(); + this.log("info", "harness_runtime_started", { + provider: driver.provider, + runtimeKeyHash: hashRuntimeKey(key), + agentId: context.agentId, + }); + return runtime; + }) + .catch((error) => { + if (this.entries.get(key) === entry) this.entries.delete(key); + throw error; + }); + + entry = { + key, + driver, + idleTimeoutMs: driver.idleTimeoutMs ?? DEFAULT_IDLE_TIMEOUT_MS, + sessionIdleTimeoutMs: this.sessionIdleTimeoutMs, + createPromise, + activeRuns: 0, + lastUsedAt: this.now(), + closing: false, + sessions: new Map(), + activeRunWaiters: new Set(), + }; + this.entries.set(key, entry); + await createPromise; + if (this.closing || entry.closing || this.entries.get(key) !== entry) { + await this.closeEntry(entry, "pool_shutdown_during_creation"); + throw new Error("Local agent runtime pool is closed."); + } + return entry; + } + } + + private scheduleIdleClose(entry: RuntimeEntry): void { + this.clearIdleTimer(entry); + if (!Number.isFinite(entry.idleTimeoutMs) || entry.idleTimeoutMs <= 0) return; + entry.idleTimer = setTimeout(() => { + void this.evictIdle().catch((error) => { + this.log("warn", "harness_runtime_close_failed", { + provider: entry.driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + reason: "idle_timeout", + error: errorMessage(error), + }); + }); + }, entry.idleTimeoutMs); + entry.idleTimer.unref(); + } + + private clearIdleTimer(entry: RuntimeEntry): void { + if (!entry.idleTimer) return; + clearTimeout(entry.idleTimer); + entry.idleTimer = undefined; + } + + private async removeAndClose(entry: RuntimeEntry, reason: string): Promise { + if (this.entries.get(entry.key) === entry) this.entries.delete(entry.key); + await this.closeEntry(entry, reason); + } + + private async closeEntry(entry: RuntimeEntry, reason: string): Promise { + if (entry.closePromise) return entry.closePromise; + entry.closing = true; + this.clearIdleTimer(entry); + entry.closePromise = (async () => { + let runtime: LocalAgentRuntime; + try { + runtime = await entry.createPromise; + } catch { + return; + } + if (reason !== "server_shutdown" && reason !== "runtime_crashed" && reason !== "runtime_not_alive") { + await this.waitForNoActiveRuns(entry); + } + if (reason === "server_shutdown") { + // Shutdown is terminal for the provider runtime. Closing it first + // aborts stuck turns and avoids waiting forever before process cleanup. + // Do not start new individual releases here: that would race an active + // turn, and the provider runtime owns their final cleanup. Existing + // idle-release work is awaited so provider cleanup never overlaps it. + await this.waitForSessionReleases(entry); + try { + await runtime.close(); + this.log("info", "harness_runtime_closed", { + provider: entry.driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + reason, + }); + } catch (error) { + this.log("warn", "harness_runtime_close_failed", { + provider: entry.driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + reason, + error: errorMessage(error), + }); + } + entry.sessions.clear(); + return; + } + await this.releaseSessions(entry, runtime, reason); + try { + await runtime.close(); + this.log("info", "harness_runtime_closed", { + provider: entry.driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + reason, + }); + } catch (error) { + this.log("warn", "harness_runtime_close_failed", { + provider: entry.driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + reason, + error: errorMessage(error), + }); + throw error; + } + })(); + return entry.closePromise; + } + + private async releaseIdleSessions(entry: RuntimeEntry, now: number): Promise { + const releases: Promise[] = []; + for (const [providerSessionId, session] of entry.sessions) { + if (session.activeRuns > 0 || now - session.lastUsedAt < entry.sessionIdleTimeoutMs) continue; + releases.push(this.releaseSession(entry, providerSessionId)); + } + await Promise.all(releases); + } + + private async releaseSessions( + entry: RuntimeEntry, + runtime: LocalAgentRuntime, + reason: string, + ): Promise { + const releases = Array.from(entry.sessions.keys()).map((providerSessionId) => + this.releaseSession(entry, providerSessionId, runtime, reason)); + await Promise.all(releases); + entry.sessions.clear(); + } + + private async releaseSession( + entry: RuntimeEntry, + providerSessionId: string, + runtime = entry.runtime, + reason = "idle_timeout", + ): Promise { + const session = entry.sessions.get(providerSessionId); + if (!runtime || !session) return; + if (entry.closing && reason === "idle_timeout") return; + if (session.releasePromise) return session.releasePromise; + const releasePromise = (async () => { + try { + await runtime.releaseSession(providerSessionId); + if (entry.sessions.get(providerSessionId) === session && session.activeRuns === 0) { + entry.sessions.delete(providerSessionId); + } + } catch (error) { + this.log("warn", "harness_session_release_failed", { + provider: entry.driver.provider, + runtimeKeyHash: hashRuntimeKey(entry.key), + providerSessionIdPrefix: providerSessionId.slice(0, 8), + reason, + error: errorMessage(error), + }); + } + })(); + session.releasePromise = releasePromise; + try { + await releasePromise; + } finally { + if (entry.sessions.get(providerSessionId) === session) session.releasePromise = undefined; + } + } + + private async waitForNoActiveRuns(entry: RuntimeEntry): Promise { + if (entry.activeRuns === 0) return; + await new Promise((resolve) => entry.activeRunWaiters.add(resolve)); + } + + private async waitForSessionReleases(entry: RuntimeEntry): Promise { + const releases = Array.from(entry.sessions.values()) + .map((session) => session.releasePromise) + .filter((release): release is Promise => Boolean(release)); + await Promise.all(releases); + } + + private log( + level: "info" | "warn" | "error", + event: string, + fields: Record, + ): void { + this.logger?.(level, event, fields); + } +} + +function hashRuntimeKey(key: string): string { + return createHash("sha256").update(key).digest("hex").slice(0, 12); +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/src/local-agent-runtime.test.ts b/src/local-agent-runtime.test.ts index 1d45d166..5321d3b7 100644 --- a/src/local-agent-runtime.test.ts +++ b/src/local-agent-runtime.test.ts @@ -1,100 +1,207 @@ import assert from "node:assert/strict"; -import type { RunResult, ThreadOptions } from "@openai/codex-sdk"; -import { - CodexSdkLocalAgentRuntime, - createCodexSdkLocalAgentRuntime, +import { LocalAgentRuntimePool } from "./local-agent-runtime-pool.js"; +import type { + LocalAgentDriver, + LocalAgentRunInput, + LocalAgentRunResult, + LocalAgentRuntime, + LocalAgentRuntimeContext, } from "./local-agent-runtime.js"; -const emptyTurn = (finalResponse: string): RunResult => ({ - finalResponse, - items: [], - usage: null, -}); +const context: LocalAgentRuntimeContext = { + agentId: "agt_test", + provider: "codex", + workspace: "/tmp/project", +}; +const input: LocalAgentRunInput = { prompt: "inspect", workspace: "/tmp/project" }; -class FakeThread { - prompts: string[] = []; +class FakeRuntime implements LocalAgentRuntime { + readonly provider = "codex" as const; + alive = true; + closeCount = 0; + runCount = 0; + readonly releasedSessions: string[] = []; + private readonly pending: Array<() => void> = []; + releaseBlocked = false; + releaseStarted = false; + releaseInFlight = false; + private releaseResolve?: () => void; - constructor(readonly id: string | null) {} + releaseWait(): void { + for (const resolve of this.pending.splice(0)) resolve(); + } - async run(prompt: string): Promise { - this.prompts.push(prompt); - return emptyTurn(`response:${prompt}`); + finishSessionRelease(): void { + this.releaseResolve?.(); + this.releaseResolve = undefined; } -} -class FakeCodex { - started: ThreadOptions[] = []; - resumed: Array<{ id: string; options?: ThreadOptions }> = []; - readonly startThreadInstance = new FakeThread("new-thread"); - readonly resumeThreadInstance = new FakeThread("resumed-thread"); + async run(runInput: LocalAgentRunInput): Promise { + assert.equal(this.releaseInFlight, false, "a session turn must not overlap session release"); + this.runCount += 1; + if (runInput.prompt === "wait") await new Promise((resolve) => this.pending.push(resolve)); + return { + provider: this.provider, + providerSessionId: "thread_1", + finalResponse: `done:${runInput.prompt}`, + items: [], + }; + } + + async releaseSession(providerSessionId: string): Promise { + this.releaseInFlight = true; + this.releaseStarted = true; + this.releasedSessions.push(providerSessionId); + if (this.releaseBlocked) { + await new Promise((resolve) => { this.releaseResolve = resolve; }); + } + this.releaseInFlight = false; + } - startThread(options?: ThreadOptions): FakeThread { - this.started.push(options ?? {}); - return this.startThreadInstance; + isAlive(): boolean { + return this.alive; } - resumeThread(id: string, options?: ThreadOptions): FakeThread { - this.resumed.push({ id, options }); - return this.resumeThreadInstance; + async close(): Promise { + this.closeCount += 1; + this.alive = false; + this.releaseWait(); } } -const codex = new FakeCodex(); -const runtime = new CodexSdkLocalAgentRuntime(codex); -const readOnly = await runtime.run({ - prompt: "inspect only", - workspace: "/tmp/project", -}); +const runtime = new FakeRuntime(); +let createCount = 0; +const driver: LocalAgentDriver = { + provider: "codex", + idleTimeoutMs: Number.POSITIVE_INFINITY, + runtimeKey: () => "shared", + createRuntime: async () => { + createCount += 1; + await Promise.resolve(); + return runtime; + }, +}; -assert.equal(readOnly.provider, "codex"); -assert.equal(readOnly.providerSessionId, "new-thread"); -assert.equal(readOnly.finalResponse, "response:inspect only"); -assert.deepEqual(codex.startThreadInstance.prompts, ["inspect only"]); -assert.deepEqual(codex.started[0], { - workingDirectory: "/tmp/project", - sandboxMode: "read-only", - approvalPolicy: "never", - model: undefined, - modelReasoningEffort: undefined, -}); +const pool = new LocalAgentRuntimePool(); +const [first, second] = await Promise.all([ + pool.run(driver, context, input), + pool.run(driver, { ...context, agentId: "agt_other" }, { ...input, prompt: "second" }), +]); +assert.equal(createCount, 1, "runtime creation is single-flight per runtime key"); +assert.equal(first.finalResponse, "done:inspect"); +assert.equal(second.finalResponse, "done:second"); +assert.equal(runtime.runCount, 2); -await runtime.run({ - prompt: "make change", - workspace: "/tmp/project", - writeMode: "allowed", - model: "gpt-5.4", - thinking: "high", -}); +const running = pool.run(driver, context, { ...input, prompt: "wait", providerSessionId: "thread_1" }); +await new Promise((resolve) => setImmediate(resolve)); +await pool.evictIdle(Date.now() + 10_000_000); +assert.equal(runtime.closeCount, 0, "active runtimes are not evicted"); +runtime.releaseWait(); +await running; + +await pool.close(); +await pool.close(); +assert.equal(runtime.closeCount, 1, "runtime close is idempotent"); +assert.deepEqual(runtime.releasedSessions, [], "shutdown closes the runtime without racing session release"); +assert.equal(pool.size, 0); -assert.deepEqual(codex.started[1], { - workingDirectory: "/tmp/project", - sandboxMode: "workspace-write", - approvalPolicy: "never", - model: "gpt-5.4", - modelReasoningEffort: "high", +let clock = 0; +const sessionRuntime = new FakeRuntime(); +const sessionPool = new LocalAgentRuntimePool({ + now: () => clock, + sessionIdleTimeoutMs: 10, }); +const sessionDriver: LocalAgentDriver = { + provider: "codex", + idleTimeoutMs: Number.POSITIVE_INFINITY, + runtimeKey: () => "session-runtime", + createRuntime: async () => sessionRuntime, +}; +await sessionPool.run(sessionDriver, context, input); +clock = 11; +sessionRuntime.releaseBlocked = true; +const releasing = sessionPool.evictIdle(); +await waitFor(() => sessionRuntime.releaseStarted); +const reused = sessionPool.run(sessionDriver, context, { ...input, providerSessionId: "thread_1", prompt: "reuse" }); +await new Promise((resolve) => setImmediate(resolve)); +assert.equal(sessionRuntime.runCount, 1, "reuse waits for the in-flight session release"); +sessionRuntime.finishSessionRelease(); +await releasing; +await reused; +sessionRuntime.releaseBlocked = false; +assert.equal(sessionRuntime.releaseInFlight, false); +assert.deepEqual(sessionRuntime.releasedSessions, ["thread_1"]); +assert.equal(sessionPool.size, 1, "releasing an idle session does not close the runtime"); +await sessionPool.close(); -const resumed = await runtime.run({ - prompt: "continue", - workspace: "/tmp/project", - providerSessionId: "existing-thread", - writeMode: "full_access", +const shutdownReleaseRuntime = new FakeRuntime(); +const shutdownReleasePool = new LocalAgentRuntimePool({ + now: () => clock, + sessionIdleTimeoutMs: 10, }); +const shutdownReleaseDriver: LocalAgentDriver = { + provider: "codex", + idleTimeoutMs: Number.POSITIVE_INFINITY, + runtimeKey: () => "shutdown-release-runtime", + createRuntime: async () => shutdownReleaseRuntime, +}; +await shutdownReleasePool.run(shutdownReleaseDriver, context, input); +shutdownReleaseRuntime.releaseBlocked = true; +const shutdownRelease = shutdownReleasePool.evictIdle(30); +await waitFor(() => shutdownReleaseRuntime.releaseStarted); +const shutdown = shutdownReleasePool.close(); +await new Promise((resolve) => setImmediate(resolve)); +assert.equal(shutdownReleaseRuntime.closeCount, 0, "shutdown waits for an in-flight session release"); +shutdownReleaseRuntime.finishSessionRelease(); +await shutdownRelease; +await shutdown; +assert.equal(shutdownReleaseRuntime.closeCount, 1); -assert.equal(resumed.providerSessionId, "resumed-thread"); -assert.deepEqual(codex.resumeThreadInstance.prompts, ["continue"]); -assert.deepEqual(codex.resumed, [ - { - id: "existing-thread", - options: { - workingDirectory: "/tmp/project", - sandboxMode: "danger-full-access", - approvalPolicy: "never", - model: undefined, - modelReasoningEffort: undefined, - }, - }, -]); +class CleanupFailureRuntime extends FakeRuntime { + override async close(): Promise { + throw new Error("cleanup failed"); + } -const created = await createCodexSdkLocalAgentRuntime(undefined, () => new FakeCodex()); -assert.equal(created.provider, "codex"); + override async run(): Promise { + this.alive = false; + throw new Error("provider failed"); + } +} + +const cleanupPool = new LocalAgentRuntimePool(); +const cleanupRuntime = new CleanupFailureRuntime(); +const cleanupDriver: LocalAgentDriver = { + provider: "codex", + runtimeKey: () => "cleanup-runtime", + createRuntime: async () => cleanupRuntime, +}; +await assert.rejects( + cleanupPool.run(cleanupDriver, context, input), + /provider failed/, + "runtime cleanup must not replace the provider error", +); + +let resolveCreation!: (runtime: LocalAgentRuntime) => void; +const creating = new Promise((resolve) => { resolveCreation = resolve; }); +const raceRuntime = new FakeRuntime(); +const racePool = new LocalAgentRuntimePool(); +const raceDriver: LocalAgentDriver = { + provider: "codex", + runtimeKey: () => "creation-race", + createRuntime: async () => creating, +}; +const pendingRun = racePool.run(raceDriver, context, input); +await new Promise((resolve) => setImmediate(resolve)); +const pendingClose = racePool.close(); +resolveCreation(raceRuntime); +await pendingClose; +await assert.rejects(pendingRun, /closed/); +assert.equal(raceRuntime.closeCount, 1, "a runtime created during shutdown is closed"); + +async function waitFor(check: () => boolean): Promise { + const deadline = Date.now() + 2_000; + while (!check() && Date.now() < deadline) { + await new Promise((resolve) => setImmediate(resolve)); + } + assert.equal(check(), true, "condition did not become true before timeout"); +} diff --git a/src/local-agent-runtime.ts b/src/local-agent-runtime.ts index 54130c2e..0f2a27c9 100644 --- a/src/local-agent-runtime.ts +++ b/src/local-agent-runtime.ts @@ -1,11 +1,4 @@ -import type { - Codex, - CodexOptions, - ModelReasoningEffort, - RunResult, - SandboxMode, - ThreadOptions, -} from "@openai/codex-sdk"; +import type { LocalAgentProvider } from "./local-agent-profiles.js"; export type LocalAgentWriteMode = "read_only" | "allowed" | "full_access"; @@ -25,78 +18,42 @@ export interface LocalAgentRunResult { items: unknown[]; } -export interface LocalAgentRuntime { - readonly provider: string; - run(input: LocalAgentRunInput): Promise; -} - -interface CodexThreadLike { - readonly id: string | null; - run(prompt: string): Promise; +export interface LocalAgentRunCallbacks { + /** + * Called as soon as a provider creates or resolves a durable continuation + * identity. The callback is awaited before the provider starts work that + * could otherwise fail and lose that identity. + */ + onSessionId?: (providerSessionId: string) => void | Promise; } -interface CodexClientLike { - startThread(options?: ThreadOptions): CodexThreadLike; - resumeThread(id: string, options?: ThreadOptions): CodexThreadLike; -} - -type CodexFactory = (options?: CodexOptions) => CodexClientLike; - -function sandboxModeFor(writeMode: LocalAgentWriteMode | undefined): SandboxMode { - switch (writeMode) { - case "allowed": - return "workspace-write"; - case "full_access": - return "danger-full-access"; - case "read_only": - case undefined: - return "read-only"; - } -} - -function threadOptionsFor(input: LocalAgentRunInput): ThreadOptions { - return { - workingDirectory: input.workspace, - sandboxMode: sandboxModeFor(input.writeMode), - approvalPolicy: "never", - model: input.model, - modelReasoningEffort: input.thinking as ModelReasoningEffort | undefined, - }; -} - -export class CodexSdkLocalAgentRuntime implements LocalAgentRuntime { - readonly provider = "codex" as const; - private readonly codex: CodexClientLike; - - constructor(codex: CodexClientLike) { - this.codex = codex; - } - - async run(input: LocalAgentRunInput): Promise { - const options = threadOptionsFor(input); - const thread = input.providerSessionId - ? this.codex.resumeThread(input.providerSessionId, options) - : this.codex.startThread(options); - const turn = await thread.run(input.prompt); - - return { - provider: this.provider, - providerSessionId: thread.id, - finalResponse: turn.finalResponse, - items: turn.items, - }; - } -} - -export async function createCodexSdkLocalAgentRuntime( - options?: CodexOptions, - codexFactory?: CodexFactory, -): Promise { - const factory = codexFactory ?? (await defaultCodexFactory()); - return new CodexSdkLocalAgentRuntime(factory(options)); +export interface LocalAgentRuntimeContext { + agentId: string; + provider: LocalAgentProvider; + workspace: string; + providerSessionId?: string; + writeMode?: LocalAgentWriteMode; + model?: string; + thinking?: string; + agentDir?: string; } -async function defaultCodexFactory(): Promise { - const module = await import("@openai/codex-sdk"); - return (options) => new module.Codex(options) as Codex; +/** + * A runtime is deliberately disposable. Nothing from this interface is + * persisted; the provider session ID in LocalAgentStore is the continuation + * identity used when a later runtime is created. + */ +export interface LocalAgentRuntime { + readonly provider: LocalAgentProvider; + run(input: LocalAgentRunInput, callbacks?: LocalAgentRunCallbacks): Promise; + releaseSession(providerSessionId: string): Promise; + close(): Promise; + isAlive(): boolean; +} + +export interface LocalAgentDriver { + readonly provider: LocalAgentProvider; + runtimeKey(context: LocalAgentRuntimeContext): string; + createRuntime(context: LocalAgentRuntimeContext): Promise; + readonly idleTimeoutMs?: number; } diff --git a/src/local-agent-store.test.ts b/src/local-agent-store.test.ts index cf7265a9..b42dbe76 100644 --- a/src/local-agent-store.test.ts +++ b/src/local-agent-store.test.ts @@ -21,9 +21,9 @@ try { assert.match(created.id, /^agt_[a-f0-9]{8}$/); assert.equal(created.status, "starting"); - 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.getById(created.id)?.thinking, "high"); + assert.equal(store.getById(created.id)?.profileName, "reviewer"); + assert.equal(store.getById(created.id.slice(0, 7)), undefined); const updated = store.update(created.id, { status: "idle", @@ -34,16 +34,17 @@ try { assert.equal(updated.status, "idle"); assert.equal(updated.thinking, "medium"); - assert.equal(store.get("thread_123")?.id, created.id); - assert.equal(store.get(created.id)?.thinking, "medium"); + assert.equal(store.getById("thread_123"), undefined); + assert.equal(store.getById(created.id)?.thinking, "medium"); assert.equal(store.update(created.id, { latestResponse: undefined }).latestResponse, undefined); assert.deepEqual( store.list({ workspaceRoot: join(root, "project") }).map((agent) => agent.latestResponse), [undefined], ); - assert.deepEqual(store.list({ workspaceId: "ws_1" }).map((agent) => agent.id), [created.id]); - assert.deepEqual(store.list({ workspaceId: "ws_other" }), []); - assert.deepEqual(store.list({ workspaceRoot: join(root, "other") }), []); +assert.deepEqual(store.list({ workspaceId: "ws_1" }).map((agent) => agent.id), [created.id]); +assert.deepEqual(store.list({ workspaceId: "ws_other" }), []); +assert.deepEqual(store.list({ workspaceId: "ws_1", workspaceRoot: join(root, "other") }), []); +assert.deepEqual(store.list({ workspaceRoot: join(root, "other") }), []); const otherStore = new LocalAgentStore(root); stores.push(otherStore); diff --git a/src/local-agent-store.ts b/src/local-agent-store.ts index a850ca9f..504ebdc4 100644 --- a/src/local-agent-store.ts +++ b/src/local-agent-store.ts @@ -1,7 +1,6 @@ import { randomUUID } from "node:crypto"; import { resolve } from "node:path"; import { openDatabase, type DatabaseHandle } from "./db/client.js"; -import type { ServerConfig } from "./config.js"; export type LocalAgentStatus = "starting" | "running" | "idle" | "error" | "stopped"; @@ -30,6 +29,11 @@ export interface CreateLocalAgentRecordInput { thinking?: string; } +export interface LocalAgentWorkspaceScope { + workspaceId?: string; + workspaceRoot: string; +} + export interface LocalAgentListScope { workspaceId?: string; workspaceRoot?: string; @@ -60,7 +64,15 @@ export class LocalAgentStore { list(scope: LocalAgentListScope = {}): LocalAgentRecord[] { let rows: LocalAgentRow[]; - if (scope.workspaceId) { + if (scope.workspaceId && scope.workspaceRoot) { + rows = this.database.sqlite + .prepare( + `select * from local_agent_sessions + where workspace_id = ? and workspace_root = ? + order by updated_at desc`, + ) + .all(scope.workspaceId, resolve(scope.workspaceRoot)) as LocalAgentRow[]; + } else if (scope.workspaceId) { rows = this.database.sqlite .prepare( `select * from local_agent_sessions @@ -131,25 +143,23 @@ export class LocalAgentStore { return record; } - get(idOrPrefix: string): LocalAgentRecord | undefined { + getById(id: string): LocalAgentRecord | undefined { const exact = this.database.sqlite .prepare( `select * from local_agent_sessions - where id = ? or provider_session_id = ? + where id = ? limit 1`, ) - .get(idOrPrefix, idOrPrefix) as LocalAgentRow | undefined; - if (exact) return rowToLocalAgentRecord(exact); - - const matches = this.database.sqlite - .prepare( - `select * from local_agent_sessions - where id like ? escape '\\' or provider_session_id like ? escape '\\' - order by updated_at desc`, - ) - .all(`${escapeLike(idOrPrefix)}%`, `${escapeLike(idOrPrefix)}%`) as LocalAgentRow[]; + .get(id) as LocalAgentRow | undefined; + return exact ? rowToLocalAgentRecord(exact) : undefined; + } - return matches.length === 1 ? rowToLocalAgentRecord(matches[0]!) : undefined; + /** + * Compatibility alias for callers that already use the store directly. + * Identity lookup is exact and never falls back to provider session IDs. + */ + get(id: string): LocalAgentRecord | undefined { + return this.getById(id); } update(id: string, patch: Partial>): LocalAgentRecord { @@ -196,20 +206,26 @@ export class LocalAgentStore { return updated; } + reconcileActiveRuns(message = "DevSpace restarted while this agent turn was running."): number { + const now = new Date().toISOString(); + const result = this.database.sqlite + .prepare( + `update local_agent_sessions + set status = 'error', error = ?, updated_at = ? + where status in ('starting', 'running')`, + ) + .run(message, now); + return Number(result.changes); + } + close(): void { this.database.close(); } - private getById(id: string): LocalAgentRecord | undefined { - const row = this.database.sqlite - .prepare("select * from local_agent_sessions where id = ?") - .get(id) as LocalAgentRow | undefined; - return row ? rowToLocalAgentRecord(row) : undefined; - } } -export function createLocalAgentStore(config: ServerConfig): LocalAgentStore { - return new LocalAgentStore(config.stateDir); +export function createLocalAgentStore(stateDir: string): LocalAgentStore { + return new LocalAgentStore(stateDir); } function rowToLocalAgentRecord(row: LocalAgentRow): LocalAgentRecord { @@ -242,7 +258,3 @@ function readStatus(status: string): LocalAgentStatus { } return "error"; } - -function escapeLike(value: string): string { - return value.replaceAll("\\", "\\\\").replaceAll("%", "\\%").replaceAll("_", "\\_"); -} diff --git a/src/local-agent-targets.ts b/src/local-agent-targets.ts index 917e2804..ab6faaf0 100644 --- a/src/local-agent-targets.ts +++ b/src/local-agent-targets.ts @@ -12,6 +12,13 @@ export interface ParsedLocalAgentRunArgs { thinking?: string; } +export interface ParsedLocalAgentContinueArgs { + agentId: string; + prompt: string; + model?: string; + thinking?: string; +} + export type LocalAgentTarget = | { kind: "profile"; @@ -30,9 +37,28 @@ export type LocalAgentTarget = }; export function parseLocalAgentRunArgs(args: string[]): ParsedLocalAgentRunArgs { + const parsed = parseAgentPromptArgs( + args, + 'Usage: devspace agents run [--model ] [--thinking ] ""', + ); + return parsed; +} + +export function parseLocalAgentContinueArgs(args: string[]): ParsedLocalAgentContinueArgs { + const parsed = parseAgentPromptArgs( + args, + 'Usage: devspace agents continue [--model ] [--thinking ] ""', + ); + return { agentId: parsed.target, prompt: parsed.prompt, model: parsed.model, thinking: parsed.thinking }; +} + +function parseAgentPromptArgs( + args: string[], + usage: string, +): ParsedLocalAgentRunArgs { const [target, ...rest] = args; if (!target) { - throw new Error('Usage: devspace agents run [--model ] [--thinking ] ""'); + throw new Error(usage); } let model: string | undefined; @@ -71,7 +97,7 @@ export function parseLocalAgentRunArgs(args: string[]): ParsedLocalAgentRunArgs const prompt = promptParts.join(" ").trim(); if (!prompt) { - throw new Error('Usage: devspace agents run [--model ] [--thinking ] ""'); + throw new Error(usage); } return { target, prompt, model, thinking };