diff --git a/docs/tools/launch.md b/docs/tools/launch.md new file mode 100644 index 000000000..1a40a7fca --- /dev/null +++ b/docs/tools/launch.md @@ -0,0 +1,122 @@ +# launch + +> Launch and control long-running project processes shared by every omp instance in the same directory. + +## Source +- Tool: `packages/coding-agent/src/tools/launch.ts` +- Broker client: `packages/coding-agent/src/daemon/client.ts` +- Broker runtime: `packages/coding-agent/src/daemon/broker.ts` +- Omp process presence: `packages/coding-agent/src/daemon/presence.ts` +- Protocol: `packages/coding-agent/src/daemon/protocol.ts` +- Model-facing prompt: `packages/coding-agent/src/prompts/tools/launch.md` + +## When to use it +Use `launch` for processes that stay alive after one tool call or need later interaction: +- web development servers and file watchers +- debuggers such as lldb and gdb +- REPLs and interactive application consoles +- local services whose logs or readiness must be observed + +Use `bash` for commands that finish. Async bash remains appropriate for finite commands that need no later stdin; it is not a process supervisor. + +## Operations + +| Operation | Purpose | Main fields | +| --- | --- | --- | +| `start` | Launch a named process. | `name`, `application`, `args`, `env`, `cwd`, `pty`, `ready`, `restart`, `persist`, `detached` | +| `list` | Snapshot every managed process in the current project scope. | none | +| `logs` | Read, filter, or follow captured combined output. | `name`, `lines`, `head`, `grep`, `follow`, `cursor`, `timeout` | +| `wait` | Wait for readiness, exit, or an output regex. | `name`, `for`, `pattern`, `timeout` | +| `send` | Write stdin, terminal keys, or a process signal. | `name`, `text`, `enter`, `keys`, `signal` | +| `stop` | Gracefully terminate the managed process tree, then hard-kill if needed. | `name`, `timeout` | +| `restart` | Stop and relaunch using the retained launch specification. | `name` | +| `describe` | Show the retained launch specification and live state. | `name` | + +Names are stable and unique within one project directory. A live name must be stopped or restarted; starting a completed name creates a new launch and rotates its prior output log. + +## Starting and readiness +`application` and `args` are separate fields, so callers do not need shell quoting: + +```json +{ + "op": "start", + "name": "web", + "application": "bun", + "args": ["run", "dev"], + "ready": { + "log": "Local:.*http", + "port": 5173, + "timeout": 30 + } +} +``` + +Defaults: +- `cwd`: current coding-agent session directory +- `args`: `[]` +- `env`: `{}` over the broker's inherited environment +- `pty`: `true` +- `restart`: `no` +- `persist`: `false` +- `detached`: `false` +- readiness timeout: 30 seconds + +`detached: true` implies `persist: true`, forces `pty: false`, and disables stdin. Its process survives broker shutdown and every omp exit; a later broker reconnects to its records for logs and explicit `stop`. + +`ready.log` is a regular expression matched against captured output. `ready.port` probes TCP at `ready.host` (default `127.0.0.1`). When both are present, both must pass. A readiness timeout leaves the process running and returns its current state so the caller can inspect logs or stop it. + +Without a readiness condition, a successfully created process enters `running`. With readiness configured, it moves `starting` → `ready`; launch or nonzero-exit failures move to `failed`. + +## Logs and following +stdout and stderr are captured into one ordered stream when possible. PTY output is naturally combined. + +```json +{"op":"logs","name":"web","lines":100} +{"op":"logs","name":"web","grep":"error|warn","lines":50} +{"op":"logs","name":"web","follow":true,"cursor":1842,"timeout":30} +``` + +Each logs result returns a byte cursor. `follow: true` waits until output advances beyond the supplied cursor, the process exits, or the timeout elapses, then returns a fresh window. `head: true` reads from the beginning; the default reads the tail. + +The broker keeps a 25 MiB current log and one 25 MiB rotated log while it owns a process's output stream. A detached process writes directly to its disk log so it survives broker exit; output is not rotated while no broker is running. + +## Input and signals + +```json +{"op":"send","name":"debugger","text":"breakpoint set --name main"} +{"op":"send","name":"debugger","text":"run"} +{"op":"send","name":"debugger","keys":["CTRL_C"]} +``` + +`enter` defaults to true when `text` is present. Supported keys are `ENTER`, `TAB`, `ESCAPE`, `CTRL_C`, `CTRL_D`, `UP`, `DOWN`, `LEFT`, and `RIGHT`. Supported signals are `SIGINT`, `SIGTERM`, `SIGHUP`, `SIGQUIT`, and `SIGKILL`. + +All project clients may observe the same managed process. Input is one shared stream: each send operation is serialized, but two clients writing independently still address the same process stdin. + +## Cross-instance lifecycle +Every omp session registers its process in the canonical project scope. The first `launch` call starts a detached broker over a private socket; later `launch` calls from any registered omp process connect to the same broker and see the same names, logs, and state. + +Runtime data lives under `~/.omp/run/daemons//`: +- `broker.sock` (or a Windows named pipe) +- a mode-0600 authentication token +- broker PID metadata +- per-managed-process launch metadata and logs +- live omp process-presence records + +After the last tool socket disconnects, the broker checks the project-presence records. Live omp PIDs keep non-persistent managed processes running even when those omp instances have not called `launch`; dead PIDs are removed. Once no omp process remains, the broker waits three seconds, stops every non-persistent managed process, and exits. This PID check still works when an omp process is killed without JavaScript cleanup. + +`persist: true` explicitly opts a managed process out of last-client teardown. A broker with a live persistent process remains available without clients until another omp reconnects and stops it. Broker recovery terminates stale recorded children and preserves their records as exited instead of adopting an unknown process state. + +## Restart policies +- `no`: never restart automatically (default) +- `on-failure`: restart after a nonzero exit or runtime failure +- `always`: restart after any unexpected exit + +Automatic restarts use bounded exponential backoff up to 30 seconds. Explicit `stop` suppresses restart. `restart` always reuses the retained application, arguments, environment, working directory, PTY, readiness, persistence, and detached settings. + +## Errors and limits +- Names must be 1-48 letters, numbers, dots, underscores, or hyphens. +- `ready.port` must be an integer from 1 through 65535. +- Invalid readiness, wait, or log regular expressions are rejected before use. +- Sending to a stopped managed process or to unavailable stdin is an error. +- `logs`, `wait`, and `stop` timeouts are capped at one hour by the tool. +- PTY process-group signaling is POSIX-native. Windows ConPTY accepts input and Ctrl-C; other POSIX signals become hard termination because Windows has no equivalent signal model. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 621ab0ba9..cbd09bbd4 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -8,6 +8,8 @@ - Added a compact session-only model picker (Alt+P) for quick model switching without changing roles - Added `@` search to the Alt+P / `/switch` picker: it lists configured Ctrl+P quick roles in matching segment colors and applies the selected role's model and thinking for the current session. - Redesigned Agent Hub entries as two-line cards: identity (status glyph, name, agent type, parent when nested) on the left, active model + reasoning level and age right-aligned, with the task description on its own line; dropped the redundant `sub · of Main` noise +- Added a project-scoped `launch` tool for shared long-running services and debuggers, with readiness probes, bounded logs, PTY input, restart policies, and automatic teardown after the last omp instance exits. +- Added `detached` `launch` starts for standalone services that survive every omp instance and broker shutdown, then reconnect to the next broker for logs and explicit stop. ### Changed diff --git a/packages/coding-agent/src/cli.ts b/packages/coding-agent/src/cli.ts index 0de859a90..20ad9038c 100755 --- a/packages/coding-agent/src/cli.ts +++ b/packages/coding-agent/src/cli.ts @@ -27,6 +27,7 @@ import { import { declareWorkerHostEntry, installWorkerInbox } from "@oh-my-pi/pi-utils/worker-host"; import { installProfileAlias, resolveProfileAliasCommandFromProcess } from "./cli/profile-alias"; import { extractProfileFlags } from "./cli/profile-bootstrap"; +import { DAEMON_BROKER_WORKER_ARG } from "./launch/protocol"; if (Bun.semver.order(Bun.version, MIN_BUN_VERSION) < 0) { process.stderr.write( @@ -78,6 +79,8 @@ async function runSmokeTest(): Promise { const { smokeTestTtsWorker } = await import("./tts/tts-client"); const { smokeTestMnemopiEmbedWorker } = await import("./mnemopi/embed-client"); const { smokeTestJsEvalWorker } = await import("./eval/js/context-manager"); + // Smoke dependencies stay lazy so normal CLI startup does not load worker clients. + const { smokeTestDaemonBroker } = await import("./launch/client"); await smokeTestSyncWorker(); const statsServer = await startServer(0); @@ -97,6 +100,7 @@ async function runSmokeTest(): Promise { await smokeTestJsEvalWorker(); await smokeTestTtsWorker(); await smokeTestMnemopiEmbedWorker(); + await smokeTestDaemonBroker(); process.stdout.write("smoke-test: ok\n"); } @@ -167,6 +171,12 @@ async function runWorkerEntrypoint(arg: string | undefined): Promise { await runIpcSubprocessWorker(startMnemopiEmbedWorker); return true; } + if (arg === DAEMON_BROKER_WORKER_ARG) { + // Worker selectors must dispatch before the normal command graph loads. + const { startDaemonBrokerFromEnvironment } = await import("./launch/broker"); + await startDaemonBrokerFromEnvironment(); + return true; + } return false; } diff --git a/packages/coding-agent/src/config/settings-schema.ts b/packages/coding-agent/src/config/settings-schema.ts index 182ca9a38..9d73e3a96 100644 --- a/packages/coding-agent/src/config/settings-schema.ts +++ b/packages/coding-agent/src/config/settings-schema.ts @@ -331,6 +331,25 @@ export const DEFAULT_BASH_INTERCEPTOR_RULES: BashInterceptorRule[] = [ tool: "write", message: "Use the `write` tool instead of echo/cat redirection. It handles encoding and provides confirmation.", }, + { + pattern: "^\\s*nohup\\s+|(? = { + SIGINT: os.constants.signals.SIGINT, + SIGTERM: os.constants.signals.SIGTERM, + SIGHUP: os.constants.signals.SIGHUP, + SIGQUIT: os.constants.signals.SIGQUIT, + SIGKILL: os.constants.signals.SIGKILL, +}; + +interface ManagedProcess { + pid: number; + exited: Promise; + unref(): void; +} + +interface ManagedDaemon { + spec: DaemonSpec; + snapshot: DaemonSnapshot; + dir: string; + log?: DaemonLog; + process?: ManagedProcess; + input?: Bun.FileSink; + pty?: PtySession; + generation: number; + stopRequested: boolean; + logReady: boolean; + portReady: boolean; + readinessBuffer: string; + outputOffset: number; + readyPattern?: RegExp; + restartTimer?: NodeJS.Timeout; + consecutiveFailures: number; + persistQueue: Promise; +} + +interface BrokerLease { + path: string; + instanceId: string; +} + +function quoteShellArg(value: string): string { + return `'${value.replaceAll("'", `'\\''`)}'`; +} +function quoteCmdArg(value: string): string { + return `"${value.replaceAll('"', '""')}"`; +} + +function terminalState(state: DaemonSnapshot["state"]): boolean { + return state === "exited" || state === "failed"; +} + +async function fileTextSlice(filePath: string, head: boolean): Promise { + try { + const stat = await fs.stat(filePath); + const file = Bun.file(filePath); + if (stat.size <= LOG_READ_BYTES) return await file.text(); + return head + ? await file.slice(0, LOG_READ_BYTES).text() + : await file.slice(Math.max(0, stat.size - LOG_READ_BYTES)).text(); + } catch (error) { + if (isEnoent(error)) return ""; + throw error; + } +} + +class DaemonLog { + readonly #path: string; + readonly #previousPath: string; + readonly #file: Bun.BunFile; + #writer: Bun.FileSink; + #currentBytes = 0; + #queue: Promise = Promise.resolve(); + #closed = false; + + constructor(logPath: string, previousPath: string, file: Bun.BunFile, writer: Bun.FileSink) { + this.#path = logPath; + this.#previousPath = previousPath; + this.#file = file; + this.#writer = writer; + } + + static async open(dir: string): Promise { + await fs.mkdir(dir, { recursive: true, mode: 0o700 }); + const logPath = path.join(dir, LOG_FILE); + const previousPath = path.join(dir, PREVIOUS_LOG_FILE); + await fs.rm(previousPath, { force: true }); + try { + await fs.rename(logPath, previousPath); + } catch (error) { + if (!isEnoent(error)) throw error; + } + const file = Bun.file(logPath); + return new DaemonLog(logPath, previousPath, file, file.writer()); + } + + append(raw: string): string { + const text = sanitizeText(raw); + if (text.length === 0 || this.#closed) return text; + const bytes = Buffer.byteLength(text, "utf8"); + this.#queue = this.#queue.then(async () => { + if (this.#currentBytes > 0 && this.#currentBytes + bytes > MAX_LOG_BYTES) await this.#rotate(); + this.#writer.write(text); + this.#currentBytes += bytes; + await this.#writer.flush(); + }); + return text; + } + + async read(head: boolean, lines: number, grep?: string): Promise { + await this.#queue; + await this.#writer.flush(); + return DaemonLog.readFiles(this.#path, this.#previousPath, head, lines, grep); + } + + async close(): Promise { + if (this.#closed) return; + this.#closed = true; + await this.#queue; + await this.#writer.end(); + } + + static async readDir(dir: string, head: boolean, lines: number, grep?: string): Promise { + return DaemonLog.readFiles(path.join(dir, LOG_FILE), path.join(dir, PREVIOUS_LOG_FILE), head, lines, grep); + } + + static async readFiles( + logPath: string, + previousPath: string, + head: boolean, + lines: number, + grep?: string, + ): Promise { + const [previous, current] = await Promise.all([fileTextSlice(previousPath, head), fileTextSlice(logPath, head)]); + let text = sanitizeText(`${previous}${previous && current && !previous.endsWith("\n") ? "\n" : ""}${current}`); + if (grep) { + let pattern: RegExp; + try { + pattern = new RegExp(grep, "u"); + } catch (error) { + throw new Error(`Invalid log regex: ${error instanceof Error ? error.message : String(error)}`); + } + text = text + .split("\n") + .filter(line => pattern.test(line)) + .join("\n"); + } + const options = { maxLines: lines, maxBytes: 256 * 1024 }; + return head ? truncateHead(text, options).content : truncateTail(text, options).content; + } + + async #rotate(): Promise { + await this.#writer.end(); + await fs.rm(this.#previousPath, { force: true }); + await fs.rename(this.#path, this.#previousPath); + this.#writer = this.#file.writer(); + this.#currentBytes = 0; + } +} + +async function acquireBrokerLease(runtimeDir: string): Promise { + const pidPath = path.join(runtimeDir, PID_FILE); + for (let attempt = 0; attempt < 2; attempt++) { + try { + const handle = await fs.open(pidPath, "wx", 0o600); + const instanceId = crypto.randomUUID(); + try { + await handle.writeFile(JSON.stringify({ pid: process.pid, instanceId }), "utf8"); + } finally { + await handle.close(); + } + return { path: pidPath, instanceId }; + } catch (error) { + if (!isEexist(error)) throw error; + try { + const raw: unknown = await Bun.file(pidPath).json(); + if (typeof raw === "object" && raw !== null && "pid" in raw && typeof raw.pid === "number") { + try { + process.kill(raw.pid, 0); + return null; + } catch { + // Stale PID file; the next loop iteration claims it. + } + } + } catch { + // Malformed or partially-written PID files are stale. + } + await fs.rm(pidPath, { force: true }); + } + } + return null; +} + +async function releaseBrokerLease(lease: BrokerLease): Promise { + try { + const raw: unknown = await Bun.file(lease.path).json(); + if (typeof raw === "object" && raw !== null && "instanceId" in raw && raw.instanceId === lease.instanceId) { + await fs.rm(lease.path, { force: true }); + } + } catch (error) { + if (!isEnoent(error)) throw error; + } +} + +function connectPort(host: string, port: number): Promise { + const { promise, resolve } = Promise.withResolvers(); + const socket = net.createConnection({ host, port }); + let settled = false; + const finish = (connected: boolean): void => { + if (settled) return; + settled = true; + socket.destroy(); + resolve(connected); + }; + socket.once("connect", () => finish(true)); + socket.once("error", () => finish(false)); + socket.setTimeout(250, () => finish(false)); + return promise; +} + +class DaemonBroker { + readonly #projectDir: string; + readonly #runtimeDir: string; + readonly #endpoint: string; + readonly #token: string; + readonly #idleGraceMs: number; + readonly #records = new Map(); + readonly #clients = new Set(); + readonly #finished = Promise.withResolvers(); + readonly #sockets = new Set(); + #server: net.Server | undefined; + #idleTimer: NodeJS.Timeout | undefined; + #shuttingDown = false; + + constructor(projectDir: string, runtimeDir: string, token: string, idleGraceMs: number) { + this.#projectDir = projectDir; + this.#runtimeDir = runtimeDir; + this.#endpoint = daemonBrokerEndpoint(projectDir, runtimeDir); + this.#token = token; + this.#idleGraceMs = idleGraceMs; + } + + async run(): Promise { + await this.#recoverRecords(); + if (process.platform !== "win32") await fs.rm(this.#endpoint, { force: true }); + const server = net.createServer(socket => this.#accept(socket)); + this.#server = server; + const { promise: listening, resolve, reject } = Promise.withResolvers(); + server.once("listening", resolve); + server.once("error", reject); + server.listen(this.#endpoint); + await listening; + if (process.platform !== "win32") await fs.chmod(this.#endpoint, 0o600); + this.#scheduleIdleShutdown(); + await this.#finished.promise; + } + + async shutdown(): Promise { + if (this.#shuttingDown) return this.#finished.promise; + this.#shuttingDown = true; + clearTimeout(this.#idleTimer); + this.#idleTimer = undefined; + for (const record of this.#records.values()) { + const detached = record.spec.detached && !record.stopRequested && record.snapshot.pid !== undefined; + if (!detached && !terminalState(record.snapshot.state)) await this.#stopRecord(record, 2_000); + clearTimeout(record.restartTimer); + await record.log?.close(); + await record.persistQueue; + } + for (const socket of this.#sockets) socket.destroy(); + this.#sockets.clear(); + this.#clients.clear(); + if (this.#server) { + const { promise, resolve } = Promise.withResolvers(); + this.#server.close(() => resolve()); + await promise; + } + if (process.platform !== "win32") await fs.rm(this.#endpoint, { force: true }); + this.#finished.resolve(); + } + + #accept(socket: net.Socket): void { + this.#sockets.add(socket); + let authenticated = false; + let buffer = ""; + socket.setEncoding("utf8"); + socket.on("data", chunk => { + buffer += typeof chunk === "string" ? chunk : chunk.toString("utf8"); + if (Buffer.byteLength(buffer, "utf8") > MAX_REQUEST_BYTES) { + socket.destroy(new Error("Daemon broker request exceeds size limit")); + return; + } + for (;;) { + const newline = buffer.indexOf("\n"); + if (newline < 0) break; + const line = buffer.slice(0, newline); + buffer = buffer.slice(newline + 1); + if (!line) continue; + void this.#handleLine(socket, line, () => { + if (authenticated) return; + authenticated = true; + this.#clients.add(socket); + clearTimeout(this.#idleTimer); + this.#idleTimer = undefined; + }); + } + }); + socket.on("error", () => { + // Socket closure performs client accounting. + }); + socket.on("close", () => { + this.#sockets.delete(socket); + if (!authenticated) return; + this.#clients.delete(socket); + this.#scheduleIdleShutdown(); + }); + } + + async #handleLine(socket: net.Socket, line: string, onAuthenticated: () => void): Promise { + let id = "unknown"; + try { + const decoded: unknown = JSON.parse(line); + const request = parseDaemonWireRequest(decoded); + id = request.id; + if (request.token !== this.#token) throw new Error("Daemon broker authentication failed"); + onAuthenticated(); + const result = await this.#dispatch(request.operation); + socket.write(`${JSON.stringify({ id, ok: true, result })}\n`); + if (request.operation.op === "shutdown") setTimeout(() => void this.shutdown(), 10); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + socket.write(`${JSON.stringify({ id, ok: false, error: message })}\n`); + } + } + + async #dispatch(operation: DaemonOperation): Promise { + switch (operation.op) { + case "ping": + return { op: "ping", projectDir: this.#projectDir }; + case "start": + return this.#start(operation.spec, operation.owner); + case "list": { + await Promise.all([...this.#records.values()].map(record => this.#refreshDetached(record))); + return { + op: "list", + daemons: [...this.#records.values()] + .sort((left, right) => left.snapshot.createdAt - right.snapshot.createdAt) + .map(record => record.snapshot), + }; + } + case "logs": + return this.#logs(operation); + case "wait": + return this.#wait(operation); + case "send": + return this.#send(operation); + case "stop": { + const record = this.#record(operation.name); + await this.#stopRecord(record, operation.timeoutMs); + return { op: "stop", daemon: record.snapshot }; + } + case "restart": + return this.#restart(operation.name); + case "describe": { + const record = this.#record(operation.name); + await this.#refreshDetached(record); + return { op: "describe", daemon: record.snapshot, spec: record.spec }; + } + case "shutdown": + return { op: "shutdown" }; + } + } + + async #start(spec: DaemonSpec, owner?: string): Promise { + if (!/^[A-Za-z0-9][A-Za-z0-9._-]{0,47}$/.test(spec.name)) { + throw new Error("Daemon name must be 1-48 letters, numbers, dots, underscores, or hyphens"); + } + if (spec.detached && spec.pty) { + throw new Error("A detached daemon cannot allocate a PTY"); + } + const existing = this.#records.get(spec.name); + if (existing) await this.#refreshDetached(existing); + if (existing && !terminalState(existing.snapshot.state)) { + throw new Error(`Daemon ${spec.name} is already ${existing.snapshot.state}`); + } + if (spec.ready?.log) { + try { + new RegExp(spec.ready.log, "u"); + } catch (error) { + throw new Error(`Invalid readiness regex: ${error instanceof Error ? error.message : String(error)}`); + } + } + const stat = await fs.stat(spec.cwd); + if (!stat.isDirectory()) throw new Error(`Daemon cwd is not a directory: ${spec.cwd}`); + const dir = path.join(this.#runtimeDir, "daemons", spec.name); + const now = Date.now(); + const record: ManagedDaemon = { + spec, + snapshot: { + name: spec.name, + id: crypto.randomUUID(), + state: "starting", + createdAt: now, + startedAt: now, + restartCount: 0, + outputBytes: 0, + owner, + persist: spec.persist, + detached: spec.detached, + }, + dir, + log: await DaemonLog.open(dir), + generation: 0, + stopRequested: false, + logReady: !spec.ready?.log, + portReady: spec.ready?.port === undefined, + readinessBuffer: "", + outputOffset: 0, + readyPattern: spec.ready?.log ? new RegExp(spec.ready.log, "u") : undefined, + consecutiveFailures: 0, + persistQueue: Promise.resolve(), + }; + this.#records.set(spec.name, record); + await this.#launch(record); + let readyTimedOut = false; + if (spec.ready && !terminalState(record.snapshot.state)) { + const ready = await this.#waitUntil(record, () => record.snapshot.state === "ready", spec.ready.timeoutMs); + readyTimedOut = !ready && !terminalState(record.snapshot.state); + } + await record.persistQueue; + return { op: "start", daemon: record.snapshot, readyTimedOut }; + } + + async #launch(record: ManagedDaemon): Promise { + record.generation++; + const generation = record.generation; + record.stopRequested = false; + record.snapshot.state = record.spec.ready ? "starting" : "running"; + record.snapshot.startedAt = Date.now(); + record.snapshot.readyAt = undefined; + record.snapshot.exitedAt = undefined; + record.snapshot.exitCode = undefined; + record.snapshot.exitReason = undefined; + record.snapshot.pid = undefined; + record.snapshot.readyMatch = undefined; + record.logReady = !record.spec.ready?.log; + record.portReady = record.spec.ready?.port === undefined; + record.readinessBuffer = ""; + record.outputOffset = 0; + this.#persist(record); + try { + if (record.spec.detached) await this.#launchDetached(record, generation); + else if (record.spec.pty) await this.#launchPty(record, generation); + else this.#launchPipe(record, generation); + if (record.spec.ready?.port !== undefined) void this.#pollPort(record, generation, record.spec.ready); + this.#markReady(record); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + record.log?.append(`Daemon launch failed: ${message}\n`); + await this.#settle(record, generation, undefined, message); + } + } + + async #launchPty(record: ManagedDaemon, generation: number): Promise { + const pidPath = path.join(record.dir, "process.pid"); + await fs.rm(pidPath, { force: true }); + const argv = [record.spec.application, ...record.spec.args]; + const command = + process.platform === "win32" + ? argv.map(quoteCmdArg).join(" ") + : [`printf '%s' "$$" > ${quoteShellArg(pidPath)}`, `exec ${argv.map(quoteShellArg).join(" ")}`].join("; "); + const session = new PtySession(); + record.pty = session; + const shell = process.platform === "win32" ? process.env.COMSPEC : process.env.SHELL; + void session + .start( + { + command, + cwd: record.spec.cwd, + env: workerEnvFromParent({ TERM: "xterm-256color", ...record.spec.env }), + cols: 120, + rows: 40, + shell, + }, + (error, chunk) => { + if (generation !== record.generation) return; + if (error) record.log?.append(`PTY output error: ${error.message}\n`); + if (chunk) this.#onOutput(record, generation, chunk); + }, + ) + .then(result => this.#onPtyExit(record, generation, result)) + .catch(error => + this.#settle(record, generation, undefined, error instanceof Error ? error.message : String(error)), + ); + + if (process.platform === "win32") return; + const deadline = Date.now() + 5_000; + const pidFile = Bun.file(pidPath); + while (Date.now() < deadline && generation === record.generation) { + try { + const pid = Number.parseInt((await pidFile.text()).trim(), 10); + if (Number.isSafeInteger(pid) && pid > 0) { + record.snapshot.pid = pid; + this.#persist(record); + return; + } + } catch (error) { + if (!isEnoent(error)) throw error; + } + if (terminalState(record.snapshot.state)) return; + await Bun.sleep(20); + } + } + + #launchPipe(record: ManagedDaemon, generation: number): void { + const process = Bun.spawn([record.spec.application, ...record.spec.args], { + cwd: record.spec.cwd, + env: workerEnvFromParent(record.spec.env), + stdin: "pipe", + stdout: "pipe", + stderr: "pipe", + detached: true, + }); + record.process = process; + record.input = process.stdin; + record.snapshot.pid = process.pid; + this.#persist(record); + const stdout = this.#drain(record, generation, process.stdout); + const stderr = this.#drain(record, generation, process.stderr); + void Promise.all([stdout, stderr, process.exited]) + .then(([, , exitCode]) => this.#settle(record, generation, exitCode)) + .catch(error => + this.#settle(record, generation, undefined, error instanceof Error ? error.message : String(error)), + ); + } + + async #launchDetached(record: ManagedDaemon, generation: number): Promise { + const logPath = path.join(record.dir, LOG_FILE); + const output = await fs.open(logPath, "a", 0o600); + try { + const process = Bun.spawn([record.spec.application, ...record.spec.args], { + cwd: record.spec.cwd, + env: workerEnvFromParent(record.spec.env), + stdio: ["ignore", output.fd, output.fd], + detached: true, + }); + record.process = process; + record.snapshot.pid = process.pid; + this.#persist(record); + process.unref(); + void process.exited + .then(exitCode => this.#settle(record, generation, exitCode)) + .catch(error => + this.#settle(record, generation, undefined, error instanceof Error ? error.message : String(error)), + ); + } finally { + await output.close(); + } + } + + async #drain(record: ManagedDaemon, generation: number, stream: ReadableStream): Promise { + const reader = stream.getReader(); + const decoder = new TextDecoder(); + try { + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + if (generation === record.generation) + this.#onOutput(record, generation, decoder.decode(value, { stream: true })); + } + const tail = decoder.decode(); + if (tail && generation === record.generation) this.#onOutput(record, generation, tail); + } finally { + reader.releaseLock(); + } + } + + #onOutput(record: ManagedDaemon, generation: number, raw: string): void { + if (generation !== record.generation) return; + const text = record.log?.append(raw) ?? sanitizeText(raw); + record.snapshot.outputBytes += Buffer.byteLength(text, "utf8"); + this.#trackOutput(record, generation, text); + } + + async #readDetachedOutput(record: ManagedDaemon, generation: number): Promise { + if (!record.spec.detached || generation !== record.generation) return; + const logPath = path.join(record.dir, LOG_FILE); + let size: number; + try { + size = (await fs.stat(logPath)).size; + } catch (error) { + if (isEnoent(error)) return; + throw error; + } + if (size < record.outputOffset) record.outputOffset = 0; + if (size === record.outputOffset) return; + const file = Bun.file(logPath); + const raw = await file.slice(record.outputOffset, size).text(); + if (generation !== record.generation) return; + record.outputOffset = size; + record.snapshot.outputBytes = size; + this.#trackOutput(record, generation, sanitizeText(raw)); + } + + #trackOutput(record: ManagedDaemon, generation: number, text: string): void { + if (generation !== record.generation) return; + record.readinessBuffer = (record.readinessBuffer + text).slice(-READINESS_BUFFER_CHARS); + if (!record.logReady && record.readyPattern) { + const match = record.readyPattern.exec(record.readinessBuffer); + if (match) { + record.logReady = true; + record.snapshot.readyMatch = match[0].slice(0, 500); + } + } + this.#markReady(record); + } + + async #refreshDetached(record: ManagedDaemon): Promise { + if (!record.spec.detached || terminalState(record.snapshot.state)) return; + const generation = record.generation; + await this.#readDetachedOutput(record, generation); + if (generation !== record.generation || record.process) return; + const processRef = record.snapshot.pid === undefined ? null : Process.fromPid(record.snapshot.pid); + if (processRef?.status() === "running") return; + await this.#settle(record, generation); + } + + async #pollPort(record: ManagedDaemon, generation: number, ready: DaemonReadySpec): Promise { + const host = ready.host ?? "127.0.0.1"; + const port = ready.port; + if (port === undefined) return; + while (generation === record.generation && !terminalState(record.snapshot.state)) { + if (await connectPort(host, port)) { + record.portReady = true; + this.#markReady(record); + return; + } + await Bun.sleep(100); + } + } + + #markReady(record: ManagedDaemon): void { + if (!record.spec.ready || record.snapshot.state !== "starting") return; + if (!record.logReady || !record.portReady) return; + record.snapshot.state = "ready"; + record.snapshot.readyAt = Date.now(); + this.#persist(record); + } + + async #onPtyExit(record: ManagedDaemon, generation: number, result: PtyRunResult): Promise { + return this.#settle(record, generation, result.exitCode, result.timedOut ? "timed out" : undefined); + } + + async #settle(record: ManagedDaemon, generation: number, exitCode?: number, error?: string): Promise { + if (generation !== record.generation || terminalState(record.snapshot.state)) return; + await this.#readDetachedOutput(record, generation); + record.process = undefined; + record.input = undefined; + record.pty = undefined; + record.snapshot.pid = undefined; + record.snapshot.exitedAt = Date.now(); + record.snapshot.exitCode = exitCode; + record.snapshot.exitReason = error; + const failed = error !== undefined || (exitCode !== undefined && exitCode !== 0); + const shouldRestart = + !record.stopRequested && + (record.spec.restart === "always" || (record.spec.restart === "on-failure" && failed)); + if (shouldRestart && !this.#shuttingDown) { + const uptime = Date.now() - record.snapshot.startedAt; + record.consecutiveFailures = uptime >= 30_000 ? 0 : record.consecutiveFailures + 1; + record.snapshot.restartCount++; + record.snapshot.state = "restarting"; + const delay = Math.min(1_000 * 2 ** Math.min(record.consecutiveFailures, 5), RESTART_MAX_DELAY_MS); + record.log?.append( + `\n[daemon exited${exitCode === undefined ? "" : ` with code ${exitCode}`}; restarting in ${delay}ms]\n`, + ); + this.#persist(record); + record.restartTimer = setTimeout(() => { + record.restartTimer = undefined; + void this.#launch(record); + }, delay); + return; + } + record.snapshot.state = failed && !record.stopRequested ? "failed" : "exited"; + this.#persist(record); + await record.log?.close(); + record.log = undefined; + } + + async #logs(operation: Extract): Promise { + const record = this.#record(operation.name); + await this.#refreshDetached(record); + const cursor = operation.cursor ?? record.snapshot.outputBytes; + let timedOut = false; + if (operation.follow && record.snapshot.outputBytes <= cursor && !terminalState(record.snapshot.state)) { + const changed = await this.#waitUntil( + record, + () => record.snapshot.outputBytes > cursor || terminalState(record.snapshot.state), + operation.timeoutMs, + ); + timedOut = !changed; + } + const lines = Math.max(1, Math.min(1_000, Math.floor(operation.lines))); + const text = record.log + ? await record.log.read(operation.head, lines, operation.grep) + : await DaemonLog.readDir(record.dir, operation.head, lines, operation.grep); + return { + op: "logs", + name: record.snapshot.name, + text, + cursor: record.snapshot.outputBytes, + timedOut, + state: record.snapshot.state, + }; + } + + async #wait(operation: Extract): Promise { + const record = this.#record(operation.name); + await this.#refreshDetached(record); + let matched: string | undefined; + let pattern: RegExp | undefined; + if (operation.pattern) { + try { + pattern = new RegExp(operation.pattern, "u"); + } catch (error) { + throw new Error(`Invalid wait regex: ${error instanceof Error ? error.message : String(error)}`); + } + } + const condition = (): boolean => { + if (pattern) { + const match = pattern.exec(record.readinessBuffer); + if (!match) return false; + matched = match[0].slice(0, 500); + return true; + } + if (operation.for === "exit") return terminalState(record.snapshot.state); + return record.snapshot.state === "ready" || (record.snapshot.state === "running" && !record.spec.ready); + }; + const reached = condition() || (await this.#waitUntil(record, condition, operation.timeoutMs)); + return { op: "wait", daemon: record.snapshot, matched, timedOut: !reached }; + } + + async #send(operation: Extract): Promise { + const record = this.#record(operation.name); + await this.#refreshDetached(record); + if (terminalState(record.snapshot.state) || record.snapshot.state === "stopping") { + throw new Error(`Daemon ${operation.name} is ${record.snapshot.state}`); + } + if (operation.data === undefined && operation.signal === undefined) { + throw new Error("send requires data or signal"); + } + if (operation.data !== undefined) { + if (record.pty) record.pty.write(operation.data); + else if (record.input) { + record.input.write(operation.data); + await record.input.flush(); + } else throw new Error(`Daemon ${operation.name} stdin is unavailable`); + } + if (operation.signal) { + if (process.platform === "win32" && record.pty) { + if (operation.signal === "SIGINT") record.pty.write("\u0003"); + else record.pty.kill(); + } else { + const processRef = record.snapshot.pid === undefined ? null : Process.fromPid(record.snapshot.pid); + if (!processRef) throw new Error(`Daemon ${operation.name} process is unavailable`); + processRef.killTree(SIGNAL_NUMBER[operation.signal]); + } + } + return { op: "send", daemon: record.snapshot }; + } + + async #stopRecord(record: ManagedDaemon, timeoutMs: number): Promise { + await this.#refreshDetached(record); + if (terminalState(record.snapshot.state)) return; + record.stopRequested = true; + if (record.restartTimer) { + clearTimeout(record.restartTimer); + record.restartTimer = undefined; + record.snapshot.state = "exited"; + record.snapshot.exitedAt = Date.now(); + this.#persist(record); + await record.log?.close(); + record.log = undefined; + return; + } + record.snapshot.state = "stopping"; + this.#persist(record); + const processRef = record.snapshot.pid === undefined ? null : Process.fromPid(record.snapshot.pid); + if (processRef) await processRef.terminate({ group: true, gracefulMs: timeoutMs, timeoutMs: timeoutMs + 1_000 }); + else record.pty?.kill(); + const settled = await this.#waitUntil(record, () => terminalState(record.snapshot.state), timeoutMs + 1_000); + if (!settled && record.pty) record.pty.kill(); + } + + async #restart(name: string): Promise { + const record = this.#record(name); + await this.#stopRecord(record, 2_000); + await record.log?.close(); + record.log = await DaemonLog.open(record.dir); + record.stopRequested = false; + await this.#launch(record); + await record.persistQueue; + return { op: "restart", daemon: record.snapshot }; + } + + async #waitUntil(record: ManagedDaemon, condition: () => boolean, timeoutMs: number): Promise { + const deadline = Date.now() + Math.max(0, timeoutMs); + while (Date.now() < deadline) { + await this.#refreshDetached(record); + if (condition()) return true; + if (this.#shuttingDown && terminalState(record.snapshot.state)) return condition(); + await Bun.sleep(50); + } + await this.#refreshDetached(record); + return condition(); + } + + #record(name: string): ManagedDaemon { + const record = this.#records.get(name); + if (record) return record; + const names = [...this.#records.keys()]; + throw new Error(`Unknown daemon ${name}${names.length ? `. Available: ${names.join(", ")}` : ""}`); + } + + #persist(record: ManagedDaemon): void { + const metaPath = path.join(record.dir, META_FILE); + const tempPath = `${metaPath}.${process.pid}.tmp`; + record.persistQueue = record.persistQueue + .then(async () => { + await Bun.write(tempPath, JSON.stringify({ daemon: record.snapshot, spec: record.spec })); + await fs.rename(tempPath, metaPath); + }) + .catch(error => { + logger.warn("Failed to persist daemon metadata", { + name: record.snapshot.name, + error: error instanceof Error ? error.message : String(error), + }); + }); + } + + async #recoverRecords(): Promise { + const root = path.join(this.#runtimeDir, "daemons"); + const entries = await fs.readdir(root, { withFileTypes: true }).catch(error => { + if (isEnoent(error)) return []; + throw error; + }); + for (const entry of entries) { + if (!entry.isDirectory()) continue; + const dir = path.join(root, entry.name); + try { + const decoded: unknown = await Bun.file(path.join(dir, META_FILE)).json(); + if (typeof decoded !== "object" || decoded === null || !("daemon" in decoded) || !("spec" in decoded)) { + continue; + } + const snapshot = parseDaemonSnapshot(decoded.daemon); + const spec = parseDaemonSpec(decoded.spec); + const processRef = snapshot.pid === undefined ? null : Process.fromPid(snapshot.pid); + const detached = + spec.detached && + !terminalState(snapshot.state) && + snapshot.state !== "stopping" && + processRef?.status() === "running"; + if (!detached) { + if (processRef) await processRef.terminate({ group: true, gracefulMs: 500, timeoutMs: 2_000 }); + snapshot.pid = undefined; + snapshot.state = "exited"; + snapshot.exitedAt = Date.now(); + snapshot.exitReason = "previous broker exited"; + } else if (snapshot.state === "restarting") { + snapshot.state = spec.ready ? "starting" : "running"; + } + snapshot.persist = spec.persist; + snapshot.detached = spec.detached; + const record: ManagedDaemon = { + spec, + snapshot, + dir, + generation: 0, + stopRequested: !detached || snapshot.state === "stopping", + logReady: detached && (!spec.ready?.log || snapshot.state === "ready"), + portReady: detached && (spec.ready?.port === undefined || snapshot.state === "ready"), + readinessBuffer: "", + outputOffset: detached ? snapshot.outputBytes : 0, + readyPattern: spec.ready?.log ? new RegExp(spec.ready.log, "u") : undefined, + consecutiveFailures: 0, + persistQueue: Promise.resolve(), + }; + this.#records.set(snapshot.name, record); + if (detached && spec.ready?.port !== undefined && snapshot.state !== "ready") { + void this.#pollPort(record, record.generation, spec.ready); + } + this.#persist(record); + } catch (error) { + logger.warn("Failed to recover daemon record", { + name: entry.name, + error: error instanceof Error ? error.message : String(error), + }); + } + } + } + + #scheduleIdleShutdown(): void { + if (this.#shuttingDown || this.#clients.size > 0) return; + clearTimeout(this.#idleTimer); + this.#idleTimer = setTimeout(() => { + this.#idleTimer = undefined; + void (async () => { + const livePersistent = [...this.#records.values()].some( + record => record.spec.persist && !terminalState(record.snapshot.state), + ); + if (this.#clients.size > 0 || livePersistent) return; + if (await hasLiveDaemonProjectPresence(this.#runtimeDir)) { + this.#scheduleIdleShutdown(); + return; + } + if (this.#clients.size === 0) await this.shutdown(); + })(); + }, this.#idleGraceMs); + } +} + +/** Start the detached per-project daemon broker selected by the CLI worker host. */ +export async function startDaemonBrokerFromEnvironment(): Promise { + const projectDir = process.env[DAEMON_PROJECT_DIR_ENV]; + const runtimeDir = process.env[DAEMON_RUNTIME_DIR_ENV]; + if (!projectDir || !runtimeDir) throw new Error("Daemon broker environment is incomplete"); + delete process.env[DAEMON_PROJECT_DIR_ENV]; + delete process.env[DAEMON_RUNTIME_DIR_ENV]; + const rawGrace = process.env[DAEMON_IDLE_GRACE_ENV]; + delete process.env[DAEMON_IDLE_GRACE_ENV]; + const parsedGrace = rawGrace === undefined ? DEFAULT_IDLE_GRACE_MS : Number.parseInt(rawGrace, 10); + const idleGraceMs = Number.isFinite(parsedGrace) && parsedGrace >= 0 ? parsedGrace : DEFAULT_IDLE_GRACE_MS; + await fs.mkdir(runtimeDir, { recursive: true, mode: 0o700 }); + const lease = await acquireBrokerLease(runtimeDir); + if (!lease) return; + process.title = "omp daemon broker"; + const token = (await Bun.file(path.join(runtimeDir, TOKEN_FILE)).text()).trim(); + if (!token) throw new Error("Daemon broker token is empty"); + const broker = new DaemonBroker(projectDir, runtimeDir, token, idleGraceMs); + const cancelCleanup = postmortem.register("daemon-broker", () => broker.shutdown()); + try { + await broker.run(); + } finally { + cancelCleanup(); + await releaseBrokerLease(lease); + } +} diff --git a/packages/coding-agent/src/launch/client.ts b/packages/coding-agent/src/launch/client.ts new file mode 100644 index 000000000..0e06bf06c --- /dev/null +++ b/packages/coding-agent/src/launch/client.ts @@ -0,0 +1,344 @@ +import * as fs from "node:fs/promises"; +import * as net from "node:net"; +import * as os from "node:os"; +import * as path from "node:path"; +import { isEexist, isEnoent, postmortem } from "@oh-my-pi/pi-utils"; +import { resolveWorkerSpawnCmd, workerEnvFromParent } from "../subprocess/worker-client"; +import { daemonBrokerEndpoint, daemonRuntimeDir } from "./paths"; +import { + DAEMON_BROKER_WORKER_ARG, + DAEMON_IDLE_GRACE_ENV, + DAEMON_PROJECT_DIR_ENV, + DAEMON_RUNTIME_DIR_ENV, + type DaemonOperation, + type DaemonRpcResult, + type DaemonWireResponse, + parseDaemonRpcResult, + parseDaemonWireResponse, +} from "./protocol"; + +const CONNECT_TIMEOUT_MS = 10_000; +const CONNECT_RETRY_MS = 50; +const TOKEN_FILE = "broker.token"; + +interface PendingRequest { + operation: DaemonOperation; + resolve: (result: DaemonRpcResult) => void; + reject: (error: Error) => void; + timer: NodeJS.Timeout; + removeAbort?: () => void; +} + +/** Broker location and lifecycle overrides used by smoke tests and isolated consumers. */ +export interface DaemonBrokerClientOptions { + /** Runtime directory override; defaults to the project-scoped config path. */ + runtimeDir?: string; + /** Last-client shutdown grace override in milliseconds. */ + idleGraceMs?: number; +} + +/** Persistent per-process connection to one project's daemon broker. */ +export interface DaemonBrokerClient { + readonly projectDir: string; + request(operation: DaemonOperation, signal?: AbortSignal): Promise; + close(): void; +} + +async function canonicalProjectDir(projectDir: string): Promise { + const resolved = path.resolve(projectDir); + try { + return await fs.realpath(resolved); + } catch (error) { + if (isEnoent(error)) return resolved; + throw error; + } +} + +async function readOrCreateToken(runtimeDir: string): Promise { + await fs.mkdir(runtimeDir, { recursive: true, mode: 0o700 }); + const tokenPath = path.join(runtimeDir, TOKEN_FILE); + const tokenFile = Bun.file(tokenPath); + for (let attempt = 0; attempt < 100; attempt++) { + try { + const token = (await tokenFile.text()).trim(); + if (token.length > 0) return token; + } catch (error) { + if (!isEnoent(error)) throw error; + } + + try { + const handle = await fs.open(tokenPath, "wx", 0o600); + try { + const token = crypto.randomUUID().replaceAll("-", "") + crypto.randomUUID().replaceAll("-", ""); + await handle.writeFile(token, "utf8"); + return token; + } finally { + await handle.close(); + } + } catch (error) { + if (!isEexist(error)) throw error; + } + await Bun.sleep(10); + } + throw new Error(`Timed out initializing daemon broker token in ${runtimeDir}`); +} + +function requestTimeoutMs(operation: DaemonOperation): number { + switch (operation.op) { + case "start": + return (operation.spec.ready?.timeoutMs ?? CONNECT_TIMEOUT_MS) + 5_000; + case "wait": + case "logs": + case "stop": + return operation.timeoutMs + 5_000; + default: + return 30_000; + } +} + +function openSocket(endpoint: string, timeoutMs: number): Promise { + const { promise, resolve, reject } = Promise.withResolvers(); + const socket = net.createConnection({ path: endpoint }); + const timer = setTimeout(() => { + socket.destroy(); + reject(new Error(`Timed out connecting to daemon broker at ${endpoint}`)); + }, timeoutMs); + const cleanup = (): void => { + clearTimeout(timer); + socket.off("connect", onConnect); + socket.off("error", onError); + }; + const onConnect = (): void => { + cleanup(); + resolve(socket); + }; + const onError = (error: Error): void => { + cleanup(); + socket.destroy(); + reject(error); + }; + socket.once("connect", onConnect); + socket.once("error", onError); + return promise; +} + +class SocketDaemonClient implements DaemonBrokerClient { + readonly projectDir: string; + readonly #runtimeDir: string; + readonly #endpoint: string; + readonly #token: string; + readonly #idleGraceMs: number | undefined; + readonly #pending = new Map(); + #socket: net.Socket | undefined; + #connectPromise: Promise | undefined; + #buffer = ""; + #closed = false; + + constructor(projectDir: string, runtimeDir: string, token: string, options: DaemonBrokerClientOptions) { + this.projectDir = projectDir; + this.#runtimeDir = runtimeDir; + this.#endpoint = daemonBrokerEndpoint(projectDir, runtimeDir); + this.#token = token; + this.#idleGraceMs = options.idleGraceMs; + } + + async request(operation: DaemonOperation, signal?: AbortSignal): Promise { + if (this.#closed) throw new Error("Daemon broker client is closed"); + if (signal?.aborted) throw new Error("Daemon broker request aborted"); + await this.#connect(); + const socket = this.#socket; + if (!socket || socket.destroyed) throw new Error("Daemon broker socket is unavailable"); + + const id = crypto.randomUUID(); + const { promise, resolve, reject } = Promise.withResolvers(); + const timer = setTimeout(() => { + const pending = this.#pending.get(id); + if (!pending) return; + this.#pending.delete(id); + pending.removeAbort?.(); + reject(new Error(`Daemon ${operation.op} request timed out`)); + }, requestTimeoutMs(operation)); + const pending: PendingRequest = { operation, resolve, reject, timer }; + if (signal) { + const abort = (): void => { + if (!this.#pending.delete(id)) return; + clearTimeout(timer); + reject(new Error("Daemon broker request aborted")); + }; + signal.addEventListener("abort", abort, { once: true }); + pending.removeAbort = () => signal.removeEventListener("abort", abort); + } + this.#pending.set(id, pending); + socket.write(`${JSON.stringify({ id, token: this.#token, operation })}\n`); + return promise; + } + + close(): void { + if (this.#closed) return; + this.#closed = true; + this.#socket?.destroy(); + this.#socket = undefined; + this.#rejectPending(new Error("Daemon broker client closed")); + } + + async #connect(): Promise { + if (this.#socket && !this.#socket.destroyed) return; + if (this.#connectPromise) return this.#connectPromise; + this.#connectPromise = this.#connectOnce(); + try { + await this.#connectPromise; + } finally { + this.#connectPromise = undefined; + } + } + + async #connectOnce(): Promise { + try { + this.#bindSocket(await openSocket(this.#endpoint, 250)); + return; + } catch { + // No live broker. Multiple clients may race to spawn; the broker's PID + // lease selects one winner before any candidate touches the socket. + } + this.#spawnBroker(); + const deadline = Date.now() + CONNECT_TIMEOUT_MS; + let lastError: Error | undefined; + while (Date.now() < deadline) { + try { + this.#bindSocket(await openSocket(this.#endpoint, 250)); + return; + } catch (error) { + lastError = error instanceof Error ? error : new Error(String(error)); + await Bun.sleep(CONNECT_RETRY_MS); + } + } + throw new Error(`Failed to start daemon broker: ${lastError?.message ?? "socket unavailable"}`); + } + + #spawnBroker(): void { + const spawn = resolveWorkerSpawnCmd(DAEMON_BROKER_WORKER_ARG); + const overlay: Record = { + [DAEMON_PROJECT_DIR_ENV]: this.projectDir, + [DAEMON_RUNTIME_DIR_ENV]: this.#runtimeDir, + }; + if (this.#idleGraceMs !== undefined) overlay[DAEMON_IDLE_GRACE_ENV] = String(this.#idleGraceMs); + const child = Bun.spawn(spawn.cmd, { + cwd: spawn.cwd, + env: workerEnvFromParent(overlay), + stdin: "ignore", + stdout: "ignore", + stderr: "ignore", + detached: true, + }); + child.unref(); + } + + #bindSocket(socket: net.Socket): void { + this.#socket = socket; + this.#buffer = ""; + socket.setEncoding("utf8"); + socket.on("data", chunk => this.#onData(chunk)); + socket.on("error", () => { + // The close handler rejects pending requests with one stable error. + }); + socket.on("close", () => { + if (this.#socket === socket) this.#socket = undefined; + this.#rejectPending(new Error("Daemon broker connection closed")); + }); + } + + #onData(chunk: string | Buffer): void { + this.#buffer += typeof chunk === "string" ? chunk : chunk.toString("utf8"); + for (;;) { + const newline = this.#buffer.indexOf("\n"); + if (newline < 0) return; + const line = this.#buffer.slice(0, newline); + this.#buffer = this.#buffer.slice(newline + 1); + if (line.length === 0) continue; + let response: DaemonWireResponse; + try { + const decoded: unknown = JSON.parse(line); + response = parseDaemonWireResponse(decoded); + } catch (error) { + this.#rejectPending(error instanceof Error ? error : new Error(String(error))); + continue; + } + const pending = this.#pending.get(response.id); + if (!pending) continue; + this.#pending.delete(response.id); + clearTimeout(pending.timer); + pending.removeAbort?.(); + if (!response.ok) { + pending.reject(new Error(response.error)); + continue; + } + try { + pending.resolve(parseDaemonRpcResult(pending.operation, response.result)); + } catch (error) { + pending.reject(error instanceof Error ? error : new Error(String(error))); + } + } + } + + #rejectPending(error: Error): void { + for (const pending of this.#pending.values()) { + clearTimeout(pending.timer); + pending.removeAbort?.(); + pending.reject(error); + } + this.#pending.clear(); + } +} + +const sharedClients = new Map>(); +let cancelExitCleanup: (() => void) | undefined; + +/** Create an independent socket connection to one project's shared daemon broker. */ +export async function createDaemonBrokerClient( + projectDir: string, + options: DaemonBrokerClientOptions = {}, +): Promise { + const canonical = await canonicalProjectDir(projectDir); + const runtimeDir = options.runtimeDir ?? daemonRuntimeDir(canonical); + const token = await readOrCreateToken(runtimeDir); + return new SocketDaemonClient(canonical, runtimeDir, token, options); +} + +/** Get the process-shared daemon broker client for one canonical project directory. */ +export async function daemonClientForProject(projectDir: string): Promise { + const canonical = await canonicalProjectDir(projectDir); + let pending = sharedClients.get(canonical); + if (!pending) { + pending = createDaemonBrokerClient(canonical); + sharedClients.set(canonical, pending); + if (!cancelExitCleanup) { + cancelExitCleanup = postmortem.register("daemon-broker-clients", () => closeDaemonClients()); + } + } + return pending; +} + +/** Close every project broker connection held by this omp process. */ +export async function closeDaemonClients(): Promise { + const pending = [...sharedClients.values()]; + sharedClients.clear(); + for (const client of await Promise.all(pending)) client.close(); + cancelExitCleanup?.(); + cancelExitCleanup = undefined; +} + +/** Exercise worker-host broker startup and authenticated RPC for distribution smoke tests. */ +export async function smokeTestDaemonBroker(): Promise { + const projectDir = await fs.mkdtemp(path.join(os.tmpdir(), "omp-daemon-smoke-project-")); + const runtimeDir = await fs.mkdtemp(path.join(os.tmpdir(), "omp-daemon-smoke-run-")); + const client = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 5_000 }); + try { + const ping = await client.request({ op: "ping" }); + if (ping.op !== "ping" || ping.projectDir !== client.projectDir) throw new Error("daemon broker ping mismatch"); + await client.request({ op: "shutdown" }); + } finally { + client.close(); + await fs.rm(projectDir, { recursive: true, force: true }); + await fs.rm(runtimeDir, { recursive: true, force: true }); + } +} diff --git a/packages/coding-agent/src/launch/paths.ts b/packages/coding-agent/src/launch/paths.ts new file mode 100644 index 000000000..2426da642 --- /dev/null +++ b/packages/coding-agent/src/launch/paths.ts @@ -0,0 +1,17 @@ +import * as path from "node:path"; +import { getConfigRootDir } from "@oh-my-pi/pi-utils"; + +/** Resolve the private runtime directory shared by omp processes in one project directory. */ +export function daemonRuntimeDir(projectDir: string, configRoot: string = getConfigRootDir()): string { + const key = Bun.hash.wyhash(path.resolve(projectDir)).toString(16).padStart(16, "0"); + return path.join(configRoot, "run", "daemons", key); +} + +/** Resolve the Unix socket or Windows named pipe used by one project broker. */ +export function daemonBrokerEndpoint(projectDir: string, runtimeDir: string): string { + if (process.platform === "win32") { + const key = Bun.hash.wyhash(path.resolve(projectDir)).toString(16).padStart(16, "0"); + return `\\\\.\\pipe\\omp-daemon-${key}`; + } + return path.join(runtimeDir, "broker.sock"); +} diff --git a/packages/coding-agent/src/launch/presence.ts b/packages/coding-agent/src/launch/presence.ts new file mode 100644 index 000000000..afce0d88c --- /dev/null +++ b/packages/coding-agent/src/launch/presence.ts @@ -0,0 +1,82 @@ +import * as fs from "node:fs/promises"; +import * as path from "node:path"; +import { isEnoent, postmortem } from "@oh-my-pi/pi-utils"; +import { daemonRuntimeDir } from "./paths"; + +const CLIENTS_DIR = "clients"; + +/** Handle keeping one omp process registered in a project daemon scope. */ +export interface DaemonProjectPresence { + close(): Promise; +} + +async function canonicalProjectDir(projectDir: string): Promise { + const resolved = path.resolve(projectDir); + try { + return await fs.realpath(resolved); + } catch (error) { + if (isEnoent(error)) return resolved; + throw error; + } +} + +/** Register this omp process so project daemons survive while it remains alive. */ +export async function registerDaemonProjectPresence( + projectDir: string, + runtimeOverride?: string, +): Promise { + const canonical = await canonicalProjectDir(projectDir); + const runtimeDir = runtimeOverride ?? daemonRuntimeDir(canonical); + const clientsDir = path.join(runtimeDir, CLIENTS_DIR); + await fs.mkdir(clientsDir, { recursive: true, mode: 0o700 }); + const id = `${process.pid}-${crypto.randomUUID()}`; + const presencePath = path.join(clientsDir, `${id}.json`); + await Bun.write(presencePath, JSON.stringify({ pid: process.pid, id, projectDir: canonical })); + await fs.chmod(presencePath, 0o600); + let closed = false; + const close = async (): Promise => { + if (closed) return; + closed = true; + cancelCleanup(); + await fs.rm(presencePath, { force: true }); + }; + const cancelCleanup = postmortem.register(`daemon-presence:${id}`, () => close()); + return { close }; +} + +/** Return whether a registered omp process in this runtime directory is still alive. */ +export async function hasLiveDaemonProjectPresence(runtimeDir: string): Promise { + const clientsDir = path.join(runtimeDir, CLIENTS_DIR); + let entries: string[]; + try { + entries = await fs.readdir(clientsDir); + } catch (error) { + if (isEnoent(error)) return false; + throw error; + } + let live = false; + for (const entry of entries) { + const presencePath = path.join(clientsDir, entry); + try { + const decoded: unknown = await Bun.file(presencePath).json(); + if ( + typeof decoded !== "object" || + decoded === null || + !("pid" in decoded) || + typeof decoded.pid !== "number" + ) { + await fs.rm(presencePath, { force: true }); + continue; + } + try { + process.kill(decoded.pid, 0); + live = true; + } catch { + await fs.rm(presencePath, { force: true }); + } + } catch (error) { + if (!isEnoent(error)) await fs.rm(presencePath, { force: true }); + } + } + return live; +} diff --git a/packages/coding-agent/src/launch/protocol.ts b/packages/coding-agent/src/launch/protocol.ts new file mode 100644 index 000000000..bb61cb7ee --- /dev/null +++ b/packages/coding-agent/src/launch/protocol.ts @@ -0,0 +1,365 @@ +/** + * Cross-process daemon broker protocol shared by the tool, client, and broker. + */ +/** Hidden CLI selector used to re-enter the daemon broker worker. */ +export const DAEMON_BROKER_WORKER_ARG = "__omp_worker_daemon_broker"; + +/** Environment key carrying the broker's canonical project directory. */ +export const DAEMON_PROJECT_DIR_ENV = "OMP_DAEMON_PROJECT_DIR"; + +/** Environment key carrying the broker's private runtime directory. */ +export const DAEMON_RUNTIME_DIR_ENV = "OMP_DAEMON_RUNTIME_DIR"; + +/** Optional environment key overriding last-client shutdown grace. */ +export const DAEMON_IDLE_GRACE_ENV = "OMP_DAEMON_IDLE_GRACE_MS"; + +/** Stable lifecycle states exposed by the launch tool. */ +export type DaemonState = "starting" | "running" | "ready" | "restarting" | "stopping" | "exited" | "failed"; + +/** Restart behavior applied after an unexpected daemon exit. */ +export type DaemonRestartPolicy = "no" | "on-failure" | "always"; + +/** Readiness conditions; every configured condition must pass. */ +export interface DaemonReadySpec { + log?: string; + port?: number; + host?: string; + timeoutMs: number; +} + +/** Immutable launch specification retained for restart and inspection. */ +export interface DaemonSpec { + name: string; + application: string; + args: string[]; + env: Record; + cwd: string; + pty: boolean; + ready?: DaemonReadySpec; + restart: DaemonRestartPolicy; + persist: boolean; + detached: boolean; +} + +/** Serializable daemon state visible to every client in one project directory. */ +export interface DaemonSnapshot { + name: string; + id: string; + state: DaemonState; + pid?: number; + createdAt: number; + startedAt: number; + readyAt?: number; + exitedAt?: number; + exitCode?: number; + exitReason?: string; + restartCount: number; + outputBytes: number; + owner?: string; + readyMatch?: string; + persist: boolean; + detached: boolean; +} + +/** Signals accepted by daemon input operations. */ +export type DaemonSignal = "SIGINT" | "SIGTERM" | "SIGHUP" | "SIGQUIT" | "SIGKILL"; + +/** Typed broker operation sent over the authenticated socket. */ +export type DaemonOperation = + | { op: "ping" } + | { op: "start"; spec: DaemonSpec; owner?: string } + | { op: "list" } + | { + op: "logs"; + name: string; + lines: number; + head: boolean; + grep?: string; + follow: boolean; + cursor?: number; + timeoutMs: number; + } + | { op: "wait"; name: string; for: "ready" | "exit"; pattern?: string; timeoutMs: number } + | { op: "send"; name: string; data?: string; signal?: DaemonSignal } + | { op: "stop"; name: string; timeoutMs: number } + | { op: "restart"; name: string } + | { op: "describe"; name: string } + | { op: "shutdown" }; + +/** Typed broker result decoded before it reaches tool code. */ +export type DaemonRpcResult = + | { op: "ping"; projectDir: string } + | { op: "start"; daemon: DaemonSnapshot; readyTimedOut: boolean } + | { op: "list"; daemons: DaemonSnapshot[] } + | { + op: "logs"; + name: string; + text: string; + cursor: number; + timedOut: boolean; + state: DaemonState; + } + | { op: "wait"; daemon: DaemonSnapshot; matched?: string; timedOut: boolean } + | { op: "send"; daemon: DaemonSnapshot } + | { op: "stop"; daemon: DaemonSnapshot } + | { op: "restart"; daemon: DaemonSnapshot } + | { op: "describe"; daemon: DaemonSnapshot; spec: DaemonSpec } + | { op: "shutdown" }; + +/** Authenticated request envelope used by socket clients. */ +export interface DaemonWireRequest { + id: string; + token: string; + operation: DaemonOperation; +} + +/** Response envelope kept raw until matched with its pending operation. */ +export type DaemonWireResponse = { id: string; ok: true; result: unknown } | { id: string; ok: false; error: string }; + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +function record(value: unknown, label: string): Record { + if (!isRecord(value)) throw new Error(`${label} must be an object`); + return value; +} + +function stringValue(value: unknown, label: string): string { + if (typeof value !== "string" || value.length === 0) throw new Error(`${label} must be a non-empty string`); + return value; +} +function rawString(value: unknown, label: string): string { + if (typeof value !== "string") throw new Error(`${label} must be a string`); + return value; +} + +function optionalString(value: unknown, label: string): string | undefined { + if (value === undefined) return undefined; + return stringValue(value, label); +} + +function booleanValue(value: unknown, label: string): boolean { + if (typeof value !== "boolean") throw new Error(`${label} must be a boolean`); + return value; +} + +function numberValue(value: unknown, label: string): number { + if (typeof value !== "number" || !Number.isFinite(value)) throw new Error(`${label} must be a finite number`); + return value; +} + +function optionalNumber(value: unknown, label: string): number | undefined { + if (value === undefined) return undefined; + return numberValue(value, label); +} + +function stringArray(value: unknown, label: string): string[] { + if (!Array.isArray(value)) throw new Error(`${label} must be an array of strings`); + const result: string[] = []; + for (const item of value) result.push(rawString(item, `${label} item`)); + return result; +} + +function stringRecord(value: unknown, label: string): Record { + const source = record(value, label); + const result: Record = {}; + for (const key in source) result[key] = rawString(source[key], `${label}.${key}`); + return result; +} + +function daemonState(value: unknown): DaemonState { + const state = stringValue(value, "daemon state"); + if (state === "starting" || state === "running" || state === "ready" || state === "restarting") return state; + if (state === "stopping" || state === "exited" || state === "failed") return state; + throw new Error(`Unknown daemon state: ${state}`); +} + +function restartPolicy(value: unknown): DaemonRestartPolicy { + const policy = stringValue(value, "restart policy"); + if (policy === "no" || policy === "on-failure" || policy === "always") return policy; + throw new Error(`Unknown restart policy: ${policy}`); +} + +function daemonSignal(value: unknown): DaemonSignal { + const signal = stringValue(value, "signal"); + if (signal === "SIGINT" || signal === "SIGTERM" || signal === "SIGHUP") return signal; + if (signal === "SIGQUIT" || signal === "SIGKILL") return signal; + throw new Error(`Unknown daemon signal: ${signal}`); +} + +function readySpec(value: unknown): DaemonReadySpec { + const source = record(value, "ready"); + const log = optionalString(source.log, "ready.log"); + const port = optionalNumber(source.port, "ready.port"); + const host = optionalString(source.host, "ready.host"); + const timeoutMs = numberValue(source.timeoutMs, "ready.timeoutMs"); + if (!log && port === undefined) throw new Error("ready requires log or port"); + return { log, port, host, timeoutMs }; +} + +/** Decode and validate a daemon launch specification. */ +export function parseDaemonSpec(value: unknown): DaemonSpec { + const source = record(value, "daemon spec"); + const detached = source.detached === undefined ? false : booleanValue(source.detached, "spec.detached"); + return { + name: stringValue(source.name, "spec.name"), + application: stringValue(source.application, "spec.application"), + args: stringArray(source.args, "spec.args"), + env: stringRecord(source.env, "spec.env"), + cwd: stringValue(source.cwd, "spec.cwd"), + pty: booleanValue(source.pty, "spec.pty"), + ready: source.ready === undefined ? undefined : readySpec(source.ready), + restart: restartPolicy(source.restart), + persist: booleanValue(source.persist, "spec.persist") || detached, + detached, + }; +} + +/** Decode and validate one daemon snapshot. */ +export function parseDaemonSnapshot(value: unknown): DaemonSnapshot { + const source = record(value, "daemon snapshot"); + return { + name: stringValue(source.name, "daemon.name"), + id: stringValue(source.id, "daemon.id"), + state: daemonState(source.state), + pid: optionalNumber(source.pid, "daemon.pid"), + createdAt: numberValue(source.createdAt, "daemon.createdAt"), + startedAt: numberValue(source.startedAt, "daemon.startedAt"), + readyAt: optionalNumber(source.readyAt, "daemon.readyAt"), + exitedAt: optionalNumber(source.exitedAt, "daemon.exitedAt"), + exitCode: optionalNumber(source.exitCode, "daemon.exitCode"), + exitReason: optionalString(source.exitReason, "daemon.exitReason"), + restartCount: numberValue(source.restartCount, "daemon.restartCount"), + outputBytes: numberValue(source.outputBytes, "daemon.outputBytes"), + owner: optionalString(source.owner, "daemon.owner"), + readyMatch: optionalString(source.readyMatch, "daemon.readyMatch"), + persist: booleanValue(source.persist, "daemon.persist"), + detached: source.detached === undefined ? false : booleanValue(source.detached, "daemon.detached"), + }; +} + +/** Decode a socket request before the broker acts on it. */ +export function parseDaemonWireRequest(value: unknown): DaemonWireRequest { + const source = record(value, "daemon request"); + return { + id: stringValue(source.id, "request.id"), + token: stringValue(source.token, "request.token"), + operation: parseDaemonOperation(source.operation), + }; +} + +/** Decode a socket response envelope before resolving a pending call. */ +export function parseDaemonWireResponse(value: unknown): DaemonWireResponse { + const source = record(value, "daemon response"); + const id = stringValue(source.id, "response.id"); + if (source.ok === true) return { id, ok: true, result: source.result }; + if (source.ok === false) return { id, ok: false, error: stringValue(source.error, "response.error") }; + throw new Error("response.ok must be a boolean"); +} + +function parseDaemonOperation(value: unknown): DaemonOperation { + const source = record(value, "daemon operation"); + const op = stringValue(source.op, "operation.op"); + switch (op) { + case "ping": + case "list": + case "shutdown": + return { op }; + case "start": + return { + op, + spec: parseDaemonSpec(source.spec), + owner: optionalString(source.owner, "operation.owner"), + }; + case "logs": + return { + op, + name: stringValue(source.name, "operation.name"), + lines: numberValue(source.lines, "operation.lines"), + head: booleanValue(source.head, "operation.head"), + grep: optionalString(source.grep, "operation.grep"), + follow: booleanValue(source.follow, "operation.follow"), + cursor: optionalNumber(source.cursor, "operation.cursor"), + timeoutMs: numberValue(source.timeoutMs, "operation.timeoutMs"), + }; + case "wait": { + const target = stringValue(source.for, "operation.for"); + if (target !== "ready" && target !== "exit") throw new Error("operation.for must be ready or exit"); + return { + op, + name: stringValue(source.name, "operation.name"), + for: target, + pattern: optionalString(source.pattern, "operation.pattern"), + timeoutMs: numberValue(source.timeoutMs, "operation.timeoutMs"), + }; + } + case "send": + return { + op, + name: stringValue(source.name, "operation.name"), + data: optionalString(source.data, "operation.data"), + signal: source.signal === undefined ? undefined : daemonSignal(source.signal), + }; + case "stop": + return { + op, + name: stringValue(source.name, "operation.name"), + timeoutMs: numberValue(source.timeoutMs, "operation.timeoutMs"), + }; + case "restart": + case "describe": + return { op, name: stringValue(source.name, "operation.name") }; + default: + throw new Error(`Unknown daemon operation: ${op}`); + } +} + +/** Decode a broker result using its pending operation as the discriminator. */ +export function parseDaemonRpcResult(operation: DaemonOperation, value: unknown): DaemonRpcResult { + const source = record(value, `${operation.op} result`); + switch (operation.op) { + case "ping": + return { op: "ping", projectDir: stringValue(source.projectDir, "result.projectDir") }; + case "start": + return { + op: "start", + daemon: parseDaemonSnapshot(source.daemon), + readyTimedOut: booleanValue(source.readyTimedOut, "result.readyTimedOut"), + }; + case "list": { + if (!Array.isArray(source.daemons)) throw new Error("result.daemons must be an array"); + return { op: "list", daemons: source.daemons.map(parseDaemonSnapshot) }; + } + case "logs": + return { + op: "logs", + name: stringValue(source.name, "result.name"), + text: typeof source.text === "string" ? source.text : "", + cursor: numberValue(source.cursor, "result.cursor"), + timedOut: booleanValue(source.timedOut, "result.timedOut"), + state: daemonState(source.state), + }; + case "wait": + return { + op: "wait", + daemon: parseDaemonSnapshot(source.daemon), + matched: optionalString(source.matched, "result.matched"), + timedOut: booleanValue(source.timedOut, "result.timedOut"), + }; + case "send": + return { op: "send", daemon: parseDaemonSnapshot(source.daemon) }; + case "stop": + return { op: "stop", daemon: parseDaemonSnapshot(source.daemon) }; + case "restart": + return { op: "restart", daemon: parseDaemonSnapshot(source.daemon) }; + case "describe": + return { + op: "describe", + daemon: parseDaemonSnapshot(source.daemon), + spec: parseDaemonSpec(source.spec), + }; + case "shutdown": + return { op: "shutdown" }; + } +} diff --git a/packages/coding-agent/src/main.ts b/packages/coding-agent/src/main.ts index 8cb37eb3d..cf49f150d 100644 --- a/packages/coding-agent/src/main.ts +++ b/packages/coding-agent/src/main.ts @@ -50,6 +50,7 @@ import { injectOmpExtensionCliRoots } from "./discovery/omp-extension-roots"; import { ExtensionRunner } from "./extensibility/extensions/runner"; import type { ExtensionUIContext } from "./extensibility/extensions/types"; import { scheduleMarketplaceAutoUpdate } from "./extensibility/plugins/marketplace-auto-update"; +import { registerDaemonProjectPresence } from "./launch/presence"; import type { MCPManager } from "./mcp"; import { InteractiveMode } from "./modes/interactive-mode"; import type { PrintModeOptions } from "./modes/print-mode"; @@ -1047,11 +1048,12 @@ interface RunRootCommandDependencies { settings?: Settings; forceSetupWizard?: boolean; } +const DEFAULT_RUN_ROOT_DEPENDENCIES: RunRootCommandDependencies = {}; export async function runRootCommand( parsed: Args, rawArgs: string[], - deps: RunRootCommandDependencies = {}, + deps: RunRootCommandDependencies = DEFAULT_RUN_ROOT_DEPENDENCIES, ): Promise { logger.startTiming(); startStartupWatchdog(); @@ -1308,6 +1310,9 @@ export async function runRootCommand( } await pluginPreloadPromise; + if (deps === DEFAULT_RUN_ROOT_DEPENDENCIES) { + await logger.time("registerDaemonProjectPresence", registerDaemonProjectPresence, cwd); + } scheduleMarketplaceAutoUpdate({ autoUpdate: settingsInstance.get("marketplace.autoUpdate"), diff --git a/packages/coding-agent/src/prompts/tools/bash.md b/packages/coding-agent/src/prompts/tools/bash.md index 83ad3ba0a..fd5b6921b 100644 --- a/packages/coding-agent/src/prompts/tools/bash.md +++ b/packages/coding-agent/src/prompts/tools/bash.md @@ -5,6 +5,7 @@ Runs commands in the embedded shell — terminal ops: git, bun, cargo, python. The shell invokes **real binaries** with simple args. It is NOT full GNU Bash. Use bash ONLY for: a single binary call, or one short pipeline that COMPUTES a fact and does not depend on shell-specific regex/quoting (`wc -l`, `sort | uniq -c`, `comm`, `diff`, a checksum, `git status`). +{{#if hasLaunch}}Long-running service, watcher, debugger, REPL, or process needing later input? MUST use `launch`, not bash.{{/if}} {{#if hasEval}}Anything below → `eval` cell, not bash: - Inline interpreter scripts (`-e`/`-c`/`--eval`) when an eval runtime exists for that language @@ -33,7 +34,7 @@ Use bash ONLY for: a single binary call, or one short pipeline that COMPUTES a f - Internal URIs (`skill://`, `agent://`, …) auto-resolve to FS paths {{#if hasEval}}- Need exact pipeline semantics (`cmd | head`, multi-stage filtering) or output truncation? Prefer `eval` and process the stream directly.{{else}}- Need exact pipeline semantics (`cmd | head`, multi-stage filtering) or output truncation? Use a checked-in script, purpose-built tool, or single command that owns the output shape.{{/if}} {{#if asyncEnabled}} -- `async: true` for long-running commands when you don't need immediate output: returns a background job ID; result delivered as a follow-up. +- `async: true` defers reporting for finite commands that need no later input; completion arrives as a follow-up. {{/if}} @@ -42,6 +43,7 @@ Use bash ONLY for: a single binary call, or one short pipeline that COMPUTES a f {{#if hasGrep}}- NEVER shell out to search content or files: `grep/rg` → `grep`.{{else}}- Avoid shelling out for broad content search; use an active search/read tool when one is available.{{/if}} {{#if hasRead}}{{#if hasGlob}}- NEVER use `ls` or `find` to list or locate files — `ls` → `read` (a directory path lists entries), `find` → the `glob` tool (globbing). This is non-negotiable, even for a single quick listing.{{else}}- Prefer `read` for known file and directory reads. Only use shell listing when no file-listing tool is active.{{/if}}{{else}}{{#if hasGlob}}- Prefer `glob` for file discovery; avoid `find` when `glob` is active.{{else}}- If no file read/listing tool is active, keep shell inspection narrow and state that limitation.{{/if}}{{/if}} - Avoid head/tail/redirections: stderr already merged; long output auto-truncated, FULL capture kept at `artifact://`. +{{#if hasLaunch}}- NEVER launch daemons, watchers, dev servers, debuggers, or REPLs through bash/background shell syntax — use `launch`.{{/if}} @@ -52,9 +54,9 @@ Use bash ONLY for: a single binary call, or one short pipeline that COMPUTES a f {{#if asyncEnabled}} # Timeout and async -- `timeout` is seconds; nonzero values are clamped to `1..3600` and the process is killed on elapse. Set `timeout: 0` only for commands that must run until completion or explicit cancellation. -- `async: true` defers only reporting — it does NOT extend a nonzero timeout; use `timeout: 0` when a daemon or watcher must be cancellation-owned. -- Need a daemon or >3600s run? Use `async: true` with `timeout: 0` when the harness should keep it alive until cancellation, or detach/manage lifecycle yourself (`cmd &`, supervisor, self-restarting script). The shell session persists across calls. +- `timeout` is seconds; nonzero values are clamped to `1..3600` and the process is killed on elapse. Set `timeout: 0` only for finite commands whose completion is cancellation-owned. +- `async: true` defers only reporting; it does NOT extend a nonzero timeout. +{{#if hasLaunch}}- Need a service, watcher, debugger, REPL, or later stdin? MUST use `launch`. NEVER use `cmd &`, `nohup`, or async bash as a process supervisor.{{else}}- Need a long-running process or >3600s run? Use an external process supervisor; avoid detached shell jobs you cannot later observe or stop.{{/if}} {{/if}} {{#if autoBackgroundEnabled}} diff --git a/packages/coding-agent/src/prompts/tools/launch.md b/packages/coding-agent/src/prompts/tools/launch.md new file mode 100644 index 000000000..615e4f9e5 --- /dev/null +++ b/packages/coding-agent/src/prompts/tools/launch.md @@ -0,0 +1,25 @@ +Launches and controls project-scoped long-running processes shared by every omp instance in the same directory. + + +- Long-running service, watcher, debugger, REPL, or process needing later input? MUST use `launch`, not `bash`. +- `start` launches `application` + `args` directly. `cwd` defaults to the session directory; `pty` defaults true. +- `ready.log` is a regex; `ready.port` is a TCP port. Both supplied? BOTH MUST pass. `ready.timeout` is seconds. +- Names are unique per project directory. A completed name MAY be started again; a live name MUST be stopped or restarted. +- `list`, `logs`, `wait`, `send`, `stop`, `restart`, and `describe` address the stable `name`. +- `logs` defaults to the last 100 lines. `head: true` reads the beginning. `grep` is a regex. +- `logs` with `follow: true` waits for output after `cursor`; reuse the returned cursor on the next call. +- `wait` blocks until readiness/exit/pattern or timeout. Use it only when blocked; do useful work instead of tight polling. +- `send.text` writes stdin; `enter` defaults true. `keys` supports ENTER, TAB, ESCAPE, CTRL_C, CTRL_D, UP, DOWN, LEFT, RIGHT. +- `send.signal` supports SIGINT, SIGTERM, SIGHUP, SIGQUIT, SIGKILL. PTY input is serialized; many clients MAY observe, but writes share one input stream. +- `stop` performs graceful process-tree termination before hard-kill. `restart` reuses the retained launch spec. +- `restart` policy defaults `no`; `on-failure` and `always` use bounded backoff. +- `persist: true` opts out of last-omp teardown. Otherwise the broker stops every non-persistent supervised process after the last omp in this directory exits. +- `detached: true` survives broker shutdown and all omp exits. It implies `persist` and disables PTY/stdin. + + + +- Long-running work MUST use `launch`, not async/background bash. +- Readiness MUST be observed; process creation alone is not readiness. +- Omit `persist` and `detached` unless their survival guarantees are required. +- Use `stop`; NEVER kill an unverified PID through bash. + diff --git a/packages/coding-agent/src/tools/bash.ts b/packages/coding-agent/src/tools/bash.ts index d4dfe2458..dfd150cf0 100644 --- a/packages/coding-agent/src/tools/bash.ts +++ b/packages/coding-agent/src/tools/bash.ts @@ -397,6 +397,7 @@ export class BashTool implements AgentTool = { read: s => new ReadTool(s), bash: s => new BashTool(s), + launch: s => new LaunchTool(s), edit: s => new EditTool(s), ast_grep: s => new AstGrepTool(s), ast_edit: s => new AstEditTool(s), diff --git a/packages/coding-agent/src/tools/launch.ts b/packages/coding-agent/src/tools/launch.ts new file mode 100644 index 000000000..2491f8bf6 --- /dev/null +++ b/packages/coding-agent/src/tools/launch.ts @@ -0,0 +1,356 @@ +import type { + AgentTool, + AgentToolContext, + AgentToolResult, + AgentToolUpdateCallback, + ToolApprovalDecision, +} from "@oh-my-pi/pi-agent-core"; +import type { ToolExample } from "@oh-my-pi/pi-ai"; +import type { Component } from "@oh-my-pi/pi-tui"; +import { Text } from "@oh-my-pi/pi-tui"; +import { prompt } from "@oh-my-pi/pi-utils"; +import { type } from "arktype"; +import type { RenderResultOptions } from "../extensibility/custom-tools/types"; +import { daemonClientForProject } from "../launch/client"; +import type { DaemonOperation, DaemonRpcResult, DaemonSnapshot, DaemonSpec } from "../launch/protocol"; +import type { Theme } from "../modes/theme/theme"; +import launchDescription from "../prompts/tools/launch.md" with { type: "text" }; +import { renderStatusLine } from "../tui"; +import type { ToolSession } from "."; +import { resolveToCwd } from "./path-utils"; +import { formatDuration, replaceTabs, shortenPath, TRUNCATE_LENGTHS, truncateToWidth } from "./render-utils"; +import { ToolError } from "./tool-errors"; + +const launchSchema = type({ + op: type("'start' | 'list' | 'logs' | 'wait' | 'send' | 'stop' | 'restart' | 'describe'").describe( + "launch operation", + ), + "name?": type("string <= 48").describe("stable project-scoped launch name"), + "application?": type("string > 0").describe("start: executable or application path"), + "args?": type("string[]").describe("start: argv passed directly to the application"), + "env?": type({ "[string]": "string" }).describe("start: extra environment variables"), + "cwd?": type("string").describe("start: working directory; defaults to the session directory"), + "pty?": type("boolean").describe("start: allocate an interactive PTY; default true"), + "ready?": type({ + "log?": type("string > 0").describe("regex matched against output"), + "port?": type("number").describe("TCP port that must accept connections"), + "host?": type("string > 0").describe("TCP readiness host; default 127.0.0.1"), + "timeout?": type("number > 0").describe("seconds to wait; default 30"), + }).describe("start: readiness conditions; all supplied conditions must pass"), + "restart?": type("'no' | 'on-failure' | 'always'").describe("start: restart policy; default no"), + "persist?": type("boolean").describe("start: survive the last omp client exiting; default false"), + "detached?": type("boolean").describe( + "start: survive every omp and broker exit; implies persist and disables PTY input", + ), + "lines?": type("number > 0").describe("logs: output lines; default 100, max 1000"), + "head?": type("boolean").describe("logs: read from the beginning instead of the tail"), + "grep?": type("string > 0").describe("logs: regex filter"), + "follow?": type("boolean").describe("logs: wait for output newer than cursor"), + "cursor?": type("number >= 0").describe("logs: output cursor returned by an earlier call"), + "for?": type("'ready' | 'exit'").describe("wait: lifecycle condition; default exit"), + "pattern?": type("string > 0").describe("wait: output regex; takes precedence over for"), + "text?": type("string > 0").describe("send: stdin text"), + "enter?": type("boolean").describe("send: append Enter after text; default true"), + "keys?": type("string[]").describe("send: terminal keys after text"), + "signal?": type("'SIGINT' | 'SIGTERM' | 'SIGHUP' | 'SIGQUIT' | 'SIGKILL'").describe("send: process-tree signal"), + "timeout?": type("number > 0").describe("logs/wait/stop: max seconds; default 30 (stop: 5)"), +}); + +type LaunchParams = typeof launchSchema.infer; + +const KEY_INPUT: Record = { + ENTER: "\r", + TAB: "\t", + ESCAPE: "\u001b", + CTRL_C: "\u0003", + CTRL_D: "\u0004", + UP: "\u001b[A", + DOWN: "\u001b[B", + RIGHT: "\u001b[C", + LEFT: "\u001b[D", +}; + +/** Structured launch state retained for compact TUI rendering. */ +export interface LaunchToolDetails { + op: LaunchParams["op"]; + daemon?: DaemonSnapshot; + daemons?: DaemonSnapshot[]; + cursor?: number; + timedOut?: boolean; +} + +function requiredName(params: LaunchParams): string { + if (!params.name) throw new ToolError(`${params.op} requires name`); + return params.name; +} + +function timeoutMs(value: number | undefined, fallbackSeconds: number): number { + const seconds = Math.max(0.05, Math.min(3_600, value ?? fallbackSeconds)); + return Math.round(seconds * 1_000); +} + +function commandSpec(params: LaunchParams, session: ToolSession): DaemonSpec { + const name = requiredName(params); + if (!params.application) throw new ToolError("start requires application"); + const ready = params.ready; + const detached = params.detached ?? false; + if (ready?.port !== undefined && (!Number.isInteger(ready.port) || ready.port < 1 || ready.port > 65_535)) { + throw new ToolError("ready.port must be an integer from 1 to 65535"); + } + if (ready && !ready.log && ready.port === undefined) throw new ToolError("ready requires log or port"); + return { + name, + application: params.application, + args: params.args ?? [], + env: params.env ?? {}, + cwd: resolveToCwd(params.cwd ?? session.cwd, session.cwd), + pty: detached ? false : (params.pty ?? true), + ready: ready + ? { + log: ready.log, + port: ready.port, + host: ready.host, + timeoutMs: timeoutMs(ready.timeout, 30), + } + : undefined, + restart: params.restart ?? "no", + persist: (params.persist ?? false) || detached, + detached, + }; +} + +function sendData(params: LaunchParams): string | undefined { + let data = params.text ?? ""; + if (params.text && (params.enter ?? true)) data += KEY_INPUT.ENTER; + for (const rawKey of params.keys ?? []) { + const key = rawKey.trim().toUpperCase(); + const input = KEY_INPUT[key]; + if (input === undefined) throw new ToolError(`Unsupported launch key ${rawKey}`); + data += input; + } + return data || undefined; +} + +function operationFor(params: LaunchParams, session: ToolSession): DaemonOperation { + switch (params.op) { + case "start": + return { op: "start", spec: commandSpec(params, session), owner: session.getSessionId?.() ?? undefined }; + case "list": + return { op: "list" }; + case "logs": + return { + op: "logs", + name: requiredName(params), + lines: Math.min(1_000, Math.floor(params.lines ?? 100)), + head: params.head ?? false, + grep: params.grep, + follow: params.follow ?? false, + cursor: params.cursor, + timeoutMs: timeoutMs(params.timeout, 30), + }; + case "wait": + return { + op: "wait", + name: requiredName(params), + for: params.for ?? "exit", + pattern: params.pattern, + timeoutMs: timeoutMs(params.timeout, 30), + }; + case "send": + return { + op: "send", + name: requiredName(params), + data: sendData(params), + signal: params.signal, + }; + case "stop": + return { op: "stop", name: requiredName(params), timeoutMs: timeoutMs(params.timeout, 5) }; + case "restart": + return { op: "restart", name: requiredName(params) }; + case "describe": + return { op: "describe", name: requiredName(params) }; + } +} + +function daemonLabel(daemon: DaemonSnapshot): string { + const pid = daemon.pid === undefined ? "" : ` pid=${daemon.pid}`; + const exit = daemon.exitCode === undefined ? "" : ` exit=${daemon.exitCode}`; + return `${daemon.name}: ${daemon.state}${pid}${exit} uptime=${formatDuration( + (daemon.exitedAt ?? Date.now()) - daemon.startedAt, + )} restarts=${daemon.restartCount}${daemon.detached ? " detached" : daemon.persist ? " persistent" : ""}`; +} + +function toolContent(result: DaemonRpcResult): string { + switch (result.op) { + case "ping": + case "shutdown": + throw new ToolError(`Internal daemon result ${result.op} is not tool-visible`); + case "start": { + const lines = [ + `${result.daemon.state === "failed" ? "Failed to launch" : "Started"} ${daemonLabel(result.daemon)}`, + ]; + if (result.daemon.readyMatch) lines.push(`Ready: ${result.daemon.readyMatch}`); + if (result.readyTimedOut) + lines.push("Readiness timed out; the daemon remains running. Inspect logs or stop it."); + return lines.join("\n"); + } + case "list": + return result.daemons.length + ? result.daemons.map(daemon => `- ${daemonLabel(daemon)}`).join("\n") + : "No daemons."; + case "logs": + return `${result.text}${result.text && !result.text.endsWith("\n") ? "\n" : ""}[${result.name}: ${result.state}; cursor=${result.cursor}${result.timedOut ? "; follow timed out" : ""}]`; + case "wait": + return `${daemonLabel(result.daemon)}${result.matched ? `\nMatched: ${result.matched}` : ""}${result.timedOut ? "\nWait timed out." : ""}`; + case "send": + return `Sent input to ${daemonLabel(result.daemon)}`; + case "stop": + return `Stopped ${daemonLabel(result.daemon)}`; + case "restart": + return `Restarted ${daemonLabel(result.daemon)}`; + case "describe": + return [ + daemonLabel(result.daemon), + `Command: ${[result.spec.application, ...result.spec.args].join(" ")}`, + `Cwd: ${shortenPath(result.spec.cwd)}`, + `PTY: ${result.spec.pty}; restart=${result.spec.restart}; persist=${result.spec.persist}; detached=${result.spec.detached}`, + ].join("\n"); + } +} + +function toolDetails(result: DaemonRpcResult): LaunchToolDetails { + switch (result.op) { + case "start": + return { op: "start", daemon: result.daemon, timedOut: result.readyTimedOut }; + case "list": + return { op: "list", daemons: result.daemons }; + case "logs": + return { op: "logs", cursor: result.cursor, timedOut: result.timedOut }; + case "wait": + return { op: "wait", daemon: result.daemon, timedOut: result.timedOut }; + case "send": + return { op: "send", daemon: result.daemon }; + case "stop": + return { op: "stop", daemon: result.daemon }; + case "restart": + return { op: "restart", daemon: result.daemon }; + case "describe": + return { op: "describe", daemon: result.daemon }; + case "ping": + case "shutdown": + throw new ToolError(`Internal daemon result ${result.op} is not tool-visible`); + } +} +function approvalFor(params: unknown): ToolApprovalDecision { + if (typeof params !== "object" || params === null || !("op" in params)) return "exec"; + switch (params.op) { + case "list": + case "logs": + case "wait": + case "describe": + return "read"; + default: + return "exec"; + } +} + +/** Project-scoped launch tool for supervising processes in every coding-agent session. */ +export class LaunchTool implements AgentTool { + readonly name = "launch"; + readonly label = "Launch"; + readonly loadMode = "essential"; + readonly summary = "Launch and control shared long-running project processes"; + readonly description = prompt.render(launchDescription); + readonly parameters = launchSchema; + readonly strict = true; + readonly examples: readonly ToolExample[] = [ + { + caption: "Start a dev server and wait for its log banner and port", + call: { + op: "start", + name: "web", + application: "bun", + args: ["run", "dev"], + ready: { log: "Local:.*http", port: 5173, timeout: 30 }, + }, + }, + { + caption: "Run a noninteractive service beyond broker lifetime", + call: { + op: "start", + name: "worker", + application: "worker", + args: ["serve"], + detached: true, + }, + }, + { + caption: "Inspect recent output", + call: { op: "logs", name: "web", lines: 100 }, + }, + { + caption: "Follow output after a cursor", + call: { op: "logs", name: "web", follow: true, cursor: 1842, timeout: 30 }, + }, + { + caption: "Set a debugger breakpoint", + call: { op: "send", name: "debugger", text: "breakpoint set --name main" }, + }, + { + caption: "Run a debugger command", + call: { op: "send", name: "debugger", text: "run" }, + }, + { + caption: "Interrupt a debugger", + call: { op: "send", name: "debugger", keys: ["CTRL_C"] }, + }, + ]; + readonly approval = approvalFor; + + constructor(private readonly session: ToolSession) {} + + async execute( + _toolCallId: string, + params: LaunchParams, + signal?: AbortSignal, + _onUpdate?: AgentToolUpdateCallback, + _context?: AgentToolContext, + ): Promise> { + const client = await daemonClientForProject(this.session.cwd); + const result = await client.request(operationFor(params, this.session), signal); + return { + content: [{ type: "text", text: replaceTabs(toolContent(result)) }], + details: toolDetails(result), + }; + } + + renderCall(args: LaunchParams, _options: RenderResultOptions, theme: Theme): Component { + const target = args.name ?? args.application; + return new Text( + renderStatusLine( + { + icon: "pending", + title: `Launch ${args.op ?? "…"}`, + description: target ? replaceTabs(target) : undefined, + }, + theme, + ), + 0, + 0, + ); + } + + renderResult( + result: AgentToolResult, + _options: RenderResultOptions, + theme: Theme, + ): Component { + const raw = result.content.find(item => item.type === "text")?.text ?? ""; + const text = replaceTabs(raw) + .split("\n") + .map(line => truncateToWidth(line, TRUNCATE_LENGTHS.CONTENT)) + .join("\n"); + const status = renderStatusLine({ icon: result.isError ? "error" : "success", title: "Launch" }, theme); + return new Text(`${status}${text ? `\n${text}` : ""}`, 0, 0); + } +} diff --git a/packages/coding-agent/test/tool-discovery/initial-tools.test.ts b/packages/coding-agent/test/tool-discovery/initial-tools.test.ts index 303023cc4..0d2e7881b 100644 --- a/packages/coding-agent/test/tool-discovery/initial-tools.test.ts +++ b/packages/coding-agent/test/tool-discovery/initial-tools.test.ts @@ -63,6 +63,11 @@ describe("BUILTIN_TOOLS public factory map", () => { const missing = Object.keys(BUILTIN_TOOLS).filter(name => metadata.get(name)?.loadMode === undefined); expect(missing).toEqual([]); }); + it("exposes launch instead of daemon", async () => { + const launch = await BUILTIN_TOOLS.launch(toolSession); + expect(launch?.name).toBe("launch"); + expect(Object.hasOwn(BUILTIN_TOOLS, "daemon")).toBeFalse(); + }); }); describe("built-in tool loadMode annotations", () => { diff --git a/packages/coding-agent/test/tools/bash-interceptor.test.ts b/packages/coding-agent/test/tools/bash-interceptor.test.ts index 5b067aff8..649b4184e 100644 --- a/packages/coding-agent/test/tools/bash-interceptor.test.ts +++ b/packages/coding-agent/test/tools/bash-interceptor.test.ts @@ -111,6 +111,32 @@ describe("default echo/printf redirect rule", () => { }); }); +describe("default launch rules", () => { + const tools = ["launch"]; + + it.each([ + "bun run dev", + "vite --host 0.0.0.0", + "lldb ./app", + "bun test --watch", + "nohup server", + "server &", + ])("routes %s to launch", command => { + const result = checkBashInterception(command, tools, DEFAULT_BASH_INTERCEPTOR_RULES); + expect(result.block).toBe(true); + expect(result.suggestedTool).toBe("launch"); + }); + + it.each([ + "git diff -w", + "docker compose up -d", + "bun test", + "printf 'server &'", + ])("does not misclassify finite command %s", command => { + expect(checkBashInterception(command, tools, DEFAULT_BASH_INTERCEPTOR_RULES).block).toBe(false); + }); +}); + describe("BashTool argument validation", () => { it("preserves async requests so disabled async mode returns the explicit error", async () => { const tool = createBashTool([]); diff --git a/packages/coding-agent/test/tools/launch.test.ts b/packages/coding-agent/test/tools/launch.test.ts new file mode 100644 index 000000000..76dc68c2b --- /dev/null +++ b/packages/coding-agent/test/tools/launch.test.ts @@ -0,0 +1,255 @@ +import { afterEach, describe, expect, it } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as os from "node:os"; +import * as path from "node:path"; +import { createDaemonBrokerClient, type DaemonBrokerClient } from "../../src/launch/client"; +import { registerDaemonProjectPresence } from "../../src/launch/presence"; +import type { DaemonSpec } from "../../src/launch/protocol"; + +const cleanupDirs: string[] = []; + +async function tempDir(prefix: string): Promise { + const dir = await fs.mkdtemp(path.join(os.tmpdir(), prefix)); + cleanupDirs.push(dir); + return dir; +} + +// Cross-process integration: fake timers cannot advance a detached broker or OS process table. +async function waitUntil(condition: () => boolean | Promise, timeoutMs: number): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (await condition()) return true; + await Bun.sleep(50); + } + return condition(); +} + +function processExists(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +async function shutdown(client: DaemonBrokerClient): Promise { + try { + await client.request({ op: "shutdown" }); + } catch { + // A last-client shutdown may already have closed the broker. + } + client.close(); +} + +afterEach(async () => { + while (cleanupDirs.length > 0) { + const dir = cleanupDirs.pop(); + if (dir) await fs.rm(dir, { recursive: true, force: true }); + } +}); + +describe("daemon broker", () => { + it("shares PTY output and input across project clients", async () => { + const projectDir = await tempDir("omp-daemon-project-"); + const runtimeDir = await tempDir("omp-daemon-runtime-"); + const scriptPath = path.join(projectDir, "service.ts"); + await Bun.write( + scriptPath, + `process.stdin.setRawMode?.(true); +process.stdin.setEncoding("utf8"); +process.stdin.resume(); +process.stdout.write("READY\\n"); +process.stdin.on("data", chunk => process.stdout.write("INPUT:" + JSON.stringify(chunk) + "\\n")); +setInterval(() => {}, 1000); +`, + ); + const first = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 5_000 }); + const second = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 5_000 }); + try { + const spec: DaemonSpec = { + name: "debugger", + application: process.execPath, + args: [scriptPath], + env: {}, + cwd: projectDir, + pty: true, + ready: { log: "READY", timeoutMs: 5_000 }, + restart: "no", + persist: false, + detached: false, + }; + const started = await first.request({ op: "start", spec, owner: "first-client" }); + expect(started.op).toBe("start"); + if (started.op !== "start") throw new Error("unexpected start result"); + expect(started.readyTimedOut).toBeFalse(); + expect(started.daemon.state).toBe("ready"); + + const listed = await second.request({ op: "list" }); + expect(listed.op).toBe("list"); + if (listed.op !== "list") throw new Error("unexpected list result"); + expect(listed.daemons.map(daemon => daemon.name)).toEqual(["debugger"]); + + await second.request({ op: "send", name: "debugger", data: "run\r" }); + const waited = await first.request({ + op: "wait", + name: "debugger", + for: "exit", + pattern: "INPUT", + timeoutMs: 3_000, + }); + expect(waited.op).toBe("wait"); + if (waited.op !== "wait") throw new Error("unexpected wait result"); + expect(waited.timedOut).toBeFalse(); + expect(waited.matched).toBe("INPUT"); + + const logs = await second.request({ + op: "logs", + name: "debugger", + lines: 20, + head: false, + follow: false, + timeoutMs: 1_000, + }); + expect(logs.op).toBe("logs"); + if (logs.op !== "logs") throw new Error("unexpected logs result"); + expect(logs.text).toContain("READY"); + expect(logs.text).toContain('INPUT:"run\\r"'); + + const stopped = await first.request({ op: "stop", name: "debugger", timeoutMs: 2_000 }); + expect(stopped.op).toBe("stop"); + if (stopped.op !== "stop") throw new Error("unexpected stop result"); + expect(stopped.daemon.state).toBe("exited"); + } finally { + await shutdown(first); + second.close(); + } + }, 20_000); + + it("stops non-persistent daemons after the last project omp exits", async () => { + const projectDir = await tempDir("omp-daemon-exit-project-"); + const runtimeDir = await tempDir("omp-daemon-exit-runtime-"); + const scriptPath = path.join(projectDir, "service.ts"); + await Bun.write(scriptPath, `process.stdout.write("READY\\n"); setInterval(() => {}, 1000);\n`); + const presence = await registerDaemonProjectPresence(projectDir, runtimeDir); + const first = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 200 }); + const second = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 200 }); + let pid: number | undefined; + try { + const started = await first.request({ + op: "start", + spec: { + name: "server", + application: process.execPath, + args: [scriptPath], + env: {}, + cwd: projectDir, + pty: false, + ready: { log: "READY", timeoutMs: 5_000 }, + restart: "no", + persist: false, + detached: false, + }, + }); + if (started.op !== "start" || started.daemon.pid === undefined) throw new Error("daemon did not start"); + const daemonPid = started.daemon.pid; + pid = daemonPid; + await second.request({ op: "list" }); + + first.close(); + second.close(); + // Cross-process integration: the real broker grace clock cannot be advanced with test fake timers. + await Bun.sleep(500); + expect(processExists(daemonPid)).toBeTrue(); + + await presence.close(); + const stopped = await waitUntil(() => !processExists(daemonPid), 5_000); + const socketRemoved = await waitUntil( + () => + Bun.file(path.join(runtimeDir, "broker.sock")) + .exists() + .then(exists => !exists), + 5_000, + ); + expect(stopped).toBeTrue(); + expect(socketRemoved).toBeTrue(); + } finally { + first.close(); + second.close(); + await presence.close(); + if (pid !== undefined && processExists(pid)) { + const rescue = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 1_000 }); + await shutdown(rescue); + } + } + }, 20_000); + + it("keeps detached daemons alive through broker replacement", async () => { + const projectDir = await tempDir("omp-daemon-detached-project-"); + const runtimeDir = await tempDir("omp-daemon-detached-runtime-"); + const scriptPath = path.join(projectDir, "service.ts"); + await Bun.write(scriptPath, `process.stdout.write("READY\\n"); setInterval(() => {}, 1000);\n`); + const first = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 5_000 }); + let recovered: DaemonBrokerClient | undefined; + let pid: number | undefined; + try { + const started = await first.request({ + op: "start", + spec: { + name: "detached", + application: process.execPath, + args: [scriptPath], + env: {}, + cwd: projectDir, + pty: false, + ready: { log: "READY", timeoutMs: 5_000 }, + restart: "no", + persist: false, + detached: true, + }, + }); + if (started.op !== "start" || started.daemon.pid === undefined) + throw new Error("detached daemon did not start"); + pid = started.daemon.pid; + expect(started.daemon.persist).toBeTrue(); + expect(started.daemon.detached).toBeTrue(); + + await first.request({ op: "shutdown" }); + first.close(); + // Broker shutdown happens in another process, so fake timers cannot observe its lease release. + const brokerStopped = await waitUntil( + () => + Bun.file(path.join(runtimeDir, "broker.pid")) + .exists() + .then(exists => !exists), + 5_000, + ); + expect(brokerStopped).toBeTrue(); + expect(processExists(pid)).toBeTrue(); + + recovered = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 5_000 }); + const described = await recovered.request({ op: "describe", name: "detached" }); + if (described.op !== "describe") throw new Error("detached daemon did not recover"); + expect(described.daemon.pid).toBe(pid); + expect(described.daemon.detached).toBeTrue(); + expect(described.spec.persist).toBeTrue(); + + const stopped = await recovered.request({ op: "stop", name: "detached", timeoutMs: 2_000 }); + if (stopped.op !== "stop") throw new Error("detached daemon did not stop"); + expect(stopped.daemon.state).toBe("exited"); + await shutdown(recovered); + recovered = undefined; + } finally { + first.close(); + recovered?.close(); + if (pid !== undefined && processExists(pid)) { + const rescue = await createDaemonBrokerClient(projectDir, { runtimeDir, idleGraceMs: 1_000 }); + try { + await rescue.request({ op: "stop", name: "detached", timeoutMs: 2_000 }); + } finally { + await shutdown(rescue); + } + } + } + }, 20_000); +});