diff --git a/packages/server/package.json b/packages/server/package.json index 9e035101c..4ddc8a72b 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -31,7 +31,7 @@ "dev": "cross-env PASEO_NODE_ENV=development node --import tsx scripts/dev-runner.ts", "dev:tsx": "cross-env PASEO_NODE_ENV=development tsx watch --ignore '**/*.timestamp-*' src/server/index.ts", "build": "node -e \"require('node:fs').rmSync('dist',{ recursive: true, force: true })\" && npm run build:lib && npm run build:scripts", - "build:lib": "tsc -p tsconfig.server.json --incremental false && node -e \"const fs=require('node:fs'); fs.mkdirSync('dist/server/server/speech/providers/local/sherpa/assets',{recursive:true}); fs.copyFileSync('src/server/speech/providers/local/sherpa/assets/silero_vad.onnx','dist/server/server/speech/providers/local/sherpa/assets/silero_vad.onnx'); fs.cpSync('src/terminal/shell-integration','dist/server/terminal/shell-integration',{recursive:true}); fs.cpSync('src/terminal/shell-integration','dist/src/terminal/shell-integration',{recursive:true});\"", + "build:lib": "tsc -p tsconfig.server.json --incremental false && node -e \"const fs=require('node:fs'); fs.mkdirSync('dist/server/server/speech/providers/local/sherpa/assets',{recursive:true}); fs.copyFileSync('src/server/speech/providers/local/sherpa/assets/silero_vad.onnx','dist/server/server/speech/providers/local/sherpa/assets/silero_vad.onnx'); fs.cpSync('src/terminal/shell-integration','dist/server/terminal/shell-integration',{recursive:true}); fs.cpSync('src/terminal/shell-integration','dist/src/terminal/shell-integration',{recursive:true}); fs.copyFileSync('src/terminal/terminal-ts-loader.mjs','dist/server/terminal/terminal-ts-loader.mjs');\"", "build:scripts": "tsc -p tsconfig.scripts.json --incremental false && node -e \"const fs=require('node:fs'); fs.mkdirSync('dist/scripts',{recursive:true}); fs.copyFileSync('scripts/mcp-stdio-socket-bridge-cli.mjs','dist/scripts/mcp-stdio-socket-bridge-cli.mjs');\"", "prepack": "npm run build", "start": "node dist/server/server/index.js", diff --git a/packages/server/src/server/bootstrap.ts b/packages/server/src/server/bootstrap.ts index 6e7161525..3f8e9c8be 100644 --- a/packages/server/src/server/bootstrap.ts +++ b/packages/server/src/server/bootstrap.ts @@ -115,7 +115,8 @@ import { DaemonConfigStore } from "./daemon-config-store.js"; import { WorkspaceGitServiceImpl } from "./workspace-git-service.js"; import { archivePersistedWorkspaceRecord } from "./workspace-archive-service.js"; import { wrapSessionMessage, type SessionOutboundMessage } from "./messages.js"; -import { createTerminalManager, type TerminalManager } from "../terminal/terminal-manager.js"; +import type { TerminalManager } from "../terminal/terminal-manager.js"; +import { createConfiguredTerminalManager } from "../terminal/terminal-manager-factory.js"; import { createConnectionOfferV2, encodeOfferToFragmentUrl } from "./connection-offer.js"; import { loadOrCreateDaemonKeyPair } from "./daemon-keypair.js"; import { startRelayTransport, type RelayTransportController } from "./relay-transport.js"; @@ -434,7 +435,7 @@ export async function createPaseoDaemon( paseoHome: config.paseoHome, logger, }); - const terminalManager = createTerminalManager(); + const terminalManager = createConfiguredTerminalManager(); const github = createGitHubService(); const workspaceGitService = new WorkspaceGitServiceImpl({ logger, diff --git a/packages/server/src/terminal/terminal-manager-factory.test.ts b/packages/server/src/terminal/terminal-manager-factory.test.ts new file mode 100644 index 000000000..583718cd8 --- /dev/null +++ b/packages/server/src/terminal/terminal-manager-factory.test.ts @@ -0,0 +1,10 @@ +import { expect, it } from "vitest"; +import { resolveTerminalBackend } from "./terminal-manager-factory.js"; + +it("uses the worker terminal backend by default", () => { + expect(resolveTerminalBackend({})).toBe("worker"); +}); + +it("allows explicitly opting back into the in-process terminal backend", () => { + expect(resolveTerminalBackend({ PASEO_TERMINAL_BACKEND: "in-process" })).toBe("in-process"); +}); diff --git a/packages/server/src/terminal/terminal-manager-factory.ts b/packages/server/src/terminal/terminal-manager-factory.ts new file mode 100644 index 000000000..a2d6a5d26 --- /dev/null +++ b/packages/server/src/terminal/terminal-manager-factory.ts @@ -0,0 +1,18 @@ +import { createTerminalManager, type TerminalManager } from "./terminal-manager.js"; +import { createWorkerTerminalManager } from "./worker-terminal-manager.js"; + +export type TerminalBackend = "in-process" | "worker"; + +export function resolveTerminalBackend(env: NodeJS.ProcessEnv = process.env): TerminalBackend { + return env.PASEO_TERMINAL_BACKEND === "in-process" ? "in-process" : "worker"; +} + +export function createConfiguredTerminalManager(options?: { + backend?: TerminalBackend; +}): TerminalManager { + const backend = options?.backend ?? resolveTerminalBackend(); + if (backend === "worker") { + return createWorkerTerminalManager(); + } + return createTerminalManager(); +} diff --git a/packages/server/src/terminal/terminal-ts-loader.mjs b/packages/server/src/terminal/terminal-ts-loader.mjs new file mode 100644 index 000000000..16b4537f1 --- /dev/null +++ b/packages/server/src/terminal/terminal-ts-loader.mjs @@ -0,0 +1,20 @@ +import { existsSync } from "node:fs"; +import { fileURLToPath } from "node:url"; + +export async function resolve(specifier, context, nextResolve) { + if ( + context.parentURL?.startsWith("file:") && + specifier.startsWith(".") && + specifier.endsWith(".js") + ) { + const candidateUrl = new URL(specifier.replace(/\.js$/, ".ts"), context.parentURL); + if (existsSync(fileURLToPath(candidateUrl))) { + return { + url: candidateUrl.href, + shortCircuit: true, + }; + } + } + + return nextResolve(specifier, context); +} diff --git a/packages/server/src/terminal/terminal-worker-process.ts b/packages/server/src/terminal/terminal-worker-process.ts new file mode 100644 index 000000000..672da646e --- /dev/null +++ b/packages/server/src/terminal/terminal-worker-process.ts @@ -0,0 +1,213 @@ +import { createTerminalManager } from "./terminal-manager.js"; +import type { TerminalSession } from "./terminal.js"; +import type { + TerminalWorkerRequest, + TerminalWorkerToParentMessage, + WorkerTerminalInfo, +} from "./terminal-worker-protocol.js"; + +const manager = createTerminalManager(); +const unsubscribeByTerminalId = new Map void>>(); +let ipcClosing = false; + +function sendToParent(message: TerminalWorkerToParentMessage): void { + if (ipcClosing || !process.connected || !process.send) { + return; + } + try { + process.send(message, (error) => { + if (error) { + ipcClosing = true; + } + }); + } catch { + ipcClosing = true; + } +} + +function toTerminalInfo(session: TerminalSession): WorkerTerminalInfo { + return { + id: session.id, + name: session.name, + cwd: session.cwd, + ...(session.getTitle() ? { title: session.getTitle() } : {}), + }; +} + +function clearTerminalSubscriptions(terminalId: string): void { + const subscriptions = unsubscribeByTerminalId.get(terminalId); + if (subscriptions) { + for (const unsubscribe of subscriptions) { + try { + unsubscribe(); + } catch { + // no-op + } + } + } + unsubscribeByTerminalId.delete(terminalId); +} + +function watchTerminal(session: TerminalSession): void { + clearTerminalSubscriptions(session.id); + + const unsubscribeMessage = session.subscribe((message) => { + sendToParent({ + type: "terminalMessage", + terminalId: session.id, + message, + }); + }); + const unsubscribeExit = session.onExit((info) => { + clearTerminalSubscriptions(session.id); + sendToParent({ + type: "terminalExit", + terminalId: session.id, + info, + }); + }); + const unsubscribeTitle = session.onTitleChange((title) => { + sendToParent({ + type: "terminalTitleChange", + terminalId: session.id, + title, + }); + }); + const unsubscribeCommandFinished = session.onCommandFinished((info) => { + sendToParent({ + type: "terminalCommandFinished", + terminalId: session.id, + info, + }); + }); + + unsubscribeByTerminalId.set(session.id, [ + unsubscribeMessage, + unsubscribeExit, + unsubscribeTitle, + unsubscribeCommandFinished, + ]); +} + +manager.subscribeTerminalsChanged((event) => { + sendToParent({ + type: "terminalsChanged", + cwd: event.cwd, + terminals: event.terminals, + }); +}); + +async function handleRequest(message: TerminalWorkerRequest): Promise { + switch (message.type) { + case "getTerminals": { + const terminals = await manager.getTerminals(message.cwd); + sendToParent({ + type: "response", + requestId: message.requestId, + ok: true, + result: terminals.map(toTerminalInfo), + }); + return; + } + + case "createTerminal": { + const session = await manager.createTerminal(message.options); + watchTerminal(session); + sendToParent({ + type: "terminalCreated", + terminal: toTerminalInfo(session), + state: session.getState(), + }); + sendToParent({ + type: "response", + requestId: message.requestId, + ok: true, + result: { + terminal: toTerminalInfo(session), + state: session.getState(), + }, + }); + return; + } + + case "registerCwdEnv": { + manager.registerCwdEnv({ cwd: message.cwd, env: message.env }); + sendToParent({ type: "response", requestId: message.requestId, ok: true }); + return; + } + + case "killTerminal": { + const session = manager.getTerminal(message.terminalId); + const cwd = session?.cwd; + manager.killTerminal(message.terminalId); + clearTerminalSubscriptions(message.terminalId); + if (cwd) { + sendToParent({ + type: "terminalRemoved", + terminalId: message.terminalId, + cwd, + }); + } + sendToParent({ type: "response", requestId: message.requestId, ok: true }); + return; + } + + case "killTerminalAndWait": { + const session = manager.getTerminal(message.terminalId); + const cwd = session?.cwd; + await manager.killTerminalAndWait(message.terminalId, message.options); + clearTerminalSubscriptions(message.terminalId); + if (cwd) { + sendToParent({ + type: "terminalRemoved", + terminalId: message.terminalId, + cwd, + }); + } + sendToParent({ type: "response", requestId: message.requestId, ok: true }); + return; + } + + case "listDirectories": { + sendToParent({ + type: "response", + requestId: message.requestId, + ok: true, + result: manager.listDirectories(), + }); + return; + } + + case "killAll": { + manager.killAll(); + for (const terminalId of Array.from(unsubscribeByTerminalId.keys())) { + clearTerminalSubscriptions(terminalId); + } + sendToParent({ type: "response", requestId: message.requestId, ok: true }); + return; + } + + case "send": { + const session = manager.getTerminal(message.terminalId); + session?.send(message.message); + sendToParent({ type: "response", requestId: message.requestId, ok: true }); + return; + } + } +} + +process.on("message", (message: TerminalWorkerRequest) => { + void handleRequest(message).catch((error: unknown) => { + sendToParent({ + type: "response", + requestId: message.requestId, + ok: false, + error: error instanceof Error ? error.message : "Terminal worker request failed", + }); + }); +}); + +process.once("disconnect", () => { + ipcClosing = true; + manager.killAll(); +}); diff --git a/packages/server/src/terminal/terminal-worker-protocol.ts b/packages/server/src/terminal/terminal-worker-protocol.ts new file mode 100644 index 000000000..71f0b3c55 --- /dev/null +++ b/packages/server/src/terminal/terminal-worker-protocol.ts @@ -0,0 +1,122 @@ +import type { TerminalExitInfo, ServerMessage, ClientMessage } from "./terminal.js"; +import type { TerminalState } from "../shared/messages.js"; + +export interface WorkerTerminalInfo { + id: string; + name: string; + cwd: string; + title?: string; +} + +export interface WorkerCreateTerminalOptions { + id?: string; + cwd: string; + name?: string; + title?: string; + env?: Record; + command?: string; + args?: string[]; +} + +export interface WorkerKillAndWaitOptions { + gracefulTimeoutMs?: number; + forceTimeoutMs?: number; +} + +export type TerminalWorkerRequest = + | { + type: "getTerminals"; + requestId: string; + cwd: string; + } + | { + type: "createTerminal"; + requestId: string; + options: WorkerCreateTerminalOptions; + } + | { + type: "registerCwdEnv"; + requestId: string; + cwd: string; + env: Record; + } + | { + type: "killTerminal"; + requestId: string; + terminalId: string; + } + | { + type: "killTerminalAndWait"; + requestId: string; + terminalId: string; + options?: WorkerKillAndWaitOptions; + } + | { + type: "listDirectories"; + requestId: string; + } + | { + type: "killAll"; + requestId: string; + } + | { + type: "send"; + requestId: string; + terminalId: string; + message: ClientMessage; + }; + +export type TerminalWorkerResponse = + | { + type: "response"; + requestId: string; + ok: true; + result?: unknown; + } + | { + type: "response"; + requestId: string; + ok: false; + error: string; + }; + +export type TerminalWorkerEvent = + | { + type: "terminalCreated"; + terminal: WorkerTerminalInfo; + state: TerminalState; + } + | { + type: "terminalRemoved"; + terminalId: string; + cwd: string; + } + | { + type: "terminalMessage"; + terminalId: string; + message: ServerMessage; + } + | { + type: "terminalExit"; + terminalId: string; + info: TerminalExitInfo; + } + | { + type: "terminalTitleChange"; + terminalId: string; + title?: string; + } + | { + type: "terminalCommandFinished"; + terminalId: string; + info: { + exitCode: number | null; + }; + } + | { + type: "terminalsChanged"; + cwd: string; + terminals: WorkerTerminalInfo[]; + }; + +export type TerminalWorkerToParentMessage = TerminalWorkerResponse | TerminalWorkerEvent; diff --git a/packages/server/src/terminal/worker-terminal-manager.test.ts b/packages/server/src/terminal/worker-terminal-manager.test.ts new file mode 100644 index 000000000..2466e31d6 --- /dev/null +++ b/packages/server/src/terminal/worker-terminal-manager.test.ts @@ -0,0 +1,211 @@ +import { afterEach, expect, it } from "vitest"; +import { EventEmitter } from "node:events"; +import { existsSync, mkdtempSync, readFileSync, rmSync } from "node:fs"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { createWorkerTerminalManager } from "./worker-terminal-manager.js"; +import type { TerminalManager } from "./terminal-manager.js"; +import type { TerminalSession } from "./terminal.js"; +import type { TerminalState } from "../shared/messages.js"; +import type { + TerminalWorkerRequest, + TerminalWorkerToParentMessage, +} from "./terminal-worker-protocol.js"; + +async function waitForCondition( + predicate: () => boolean, + timeoutMs: number, + intervalMs = 25, +): Promise { + const start = Date.now(); + while (Date.now() - start < timeoutMs) { + if (predicate()) { + return; + } + await new Promise((resolve) => setTimeout(resolve, intervalMs)); + } + throw new Error(`Timed out after ${timeoutMs}ms waiting for condition`); +} + +async function withShell(shell: string, run: () => Promise): Promise { + const originalShell = process.env.SHELL; + process.env.SHELL = shell; + try { + return await run(); + } finally { + if (originalShell === undefined) { + delete process.env.SHELL; + } else { + process.env.SHELL = originalShell; + } + } +} + +function getVisibleText(session: TerminalSession): string { + return session + .getState() + .grid.map((row) => + row + .map((cell) => cell.char) + .join("") + .trimEnd(), + ) + .join("\n"); +} + +function createTerminalState(): TerminalState { + const blankCell = { char: " " }; + return { + rows: 1, + cols: 1, + grid: [[blankCell]], + scrollback: [], + cursor: { row: 0, col: 0 }, + }; +} + +class FakeTerminalWorker extends EventEmitter { + connected = true; + killed = false; + readonly sentMessages: TerminalWorkerRequest[] = []; + + send(message: TerminalWorkerRequest, callback: (error: Error | null) => void): boolean { + this.sentMessages.push(message); + callback(null); + return true; + } + + disconnect(): void { + this.connected = false; + this.emit("exit", 0, null); + } + + kill(): boolean { + this.killed = true; + this.connected = false; + this.emit("exit", 0, null); + return true; + } + + emitWorkerMessage(message: TerminalWorkerToParentMessage): void { + this.emit("message", message); + } +} + +let manager: TerminalManager | null = null; +const temporaryDirs: string[] = []; + +afterEach(() => { + manager?.killAll(); + manager = null; + while (temporaryDirs.length > 0) { + const dir = temporaryDirs.pop(); + if (dir) { + rmSync(dir, { recursive: true, force: true }); + } + } +}); + +it("creates a terminal through the worker and streams output", async () => { + await withShell("/bin/sh", async () => { + manager = createWorkerTerminalManager(); + const session = await manager.createTerminal({ cwd: "/tmp", env: { PS1: "$ " } }); + const messages: string[] = []; + let snapshots = 0; + const unsubscribe = session.subscribe((message) => { + if (message.type === "output") { + messages.push(message.data); + } + if (message.type === "snapshot") { + snapshots += 1; + } + }); + await new Promise((resolve) => setTimeout(resolve, 100)); + const snapshotsBeforeOutput = snapshots; + + session.send({ type: "input", data: "printf worker-output\\r" }); + + await waitForCondition( + () => + messages.join("").includes("worker-output") || + getVisibleText(session).includes("worker-output"), + 10000, + ); + await new Promise((resolve) => setTimeout(resolve, 100)); + unsubscribe(); + + expect(messages.join("") + getVisibleText(session)).toContain("worker-output"); + expect(snapshots).toBe(snapshotsBeforeOutput); + }); +}); + +it("does not surface fire-and-forget send timeouts as unhandled rejections", async () => { + const worker = new FakeTerminalWorker(); + manager = createWorkerTerminalManager({ + requestTimeoutMs: 5, + forkWorker: () => worker, + }); + + worker.emitWorkerMessage({ + type: "terminalCreated", + terminal: { id: "terminal-1", name: "Terminal", cwd: "/tmp" }, + state: createTerminalState(), + }); + const session = manager.getTerminal("terminal-1"); + expect(session).toBeDefined(); + + const unhandledRejections: unknown[] = []; + const onUnhandledRejection = (reason: unknown) => { + unhandledRejections.push(reason); + }; + process.on("unhandledRejection", onUnhandledRejection); + try { + session?.send({ type: "input", data: "x" }); + await new Promise((resolve) => setTimeout(resolve, 25)); + } finally { + process.off("unhandledRejection", onUnhandledRejection); + } + + expect(worker.sentMessages.some((message) => message.type === "send")).toBe(true); + expect(unhandledRejections).toEqual([]); +}); + +it("keeps registered cwd env inheritance behind the worker manager interface", async () => { + await withShell("/bin/sh", async () => { + manager = createWorkerTerminalManager(); + const cwd = mkdtempSync(join(tmpdir(), "worker-terminal-manager-env-")); + temporaryDirs.push(cwd); + const markerPath = join(cwd, "env.txt"); + + manager.registerCwdEnv({ + cwd, + env: { PASEO_WORKER_TERMINAL_TEST: "worker-env" }, + }); + const session = await manager.createTerminal({ cwd, env: { PS1: "$ " } }); + session.send({ + type: "input", + data: `printf '%s' "$PASEO_WORKER_TERMINAL_TEST" > ${JSON.stringify(markerPath)}\r`, + }); + + await waitForCondition(() => existsSync(markerPath), 10000); + + expect(readFileSync(markerPath, "utf8")).toBe("worker-env"); + }); +}); + +it("removes worker terminals after killAndWait", async () => { + await withShell("/bin/sh", async () => { + manager = createWorkerTerminalManager(); + const session = await manager.createTerminal({ cwd: "/tmp", env: { PS1: "$ " } }); + + await manager.killTerminalAndWait(session.id, { + gracefulTimeoutMs: 1000, + forceTimeoutMs: 500, + }); + + await waitForCondition(() => manager?.getTerminal(session.id) === undefined, 5000); + + expect(manager.getTerminal(session.id)).toBeUndefined(); + expect(manager.listDirectories()).not.toContain("/tmp"); + }); +}); diff --git a/packages/server/src/terminal/worker-terminal-manager.ts b/packages/server/src/terminal/worker-terminal-manager.ts new file mode 100644 index 000000000..8d771cf26 --- /dev/null +++ b/packages/server/src/terminal/worker-terminal-manager.ts @@ -0,0 +1,527 @@ +import { fork, type ChildProcess } from "node:child_process"; +import { fileURLToPath } from "node:url"; +import { randomUUID } from "node:crypto"; +import type { TerminalState } from "../shared/messages.js"; +import type { + ClientMessage, + ServerMessage, + TerminalCommandFinishedInfo, + TerminalExitInfo, + TerminalSession, +} from "./terminal.js"; +import type { + TerminalListItem, + TerminalManager, + TerminalsChangedEvent, + TerminalsChangedListener, +} from "./terminal-manager.js"; +import type { + TerminalWorkerRequest, + TerminalWorkerResponse, + TerminalWorkerToParentMessage, + WorkerCreateTerminalOptions, + WorkerTerminalInfo, +} from "./terminal-worker-protocol.js"; + +const REQUEST_TIMEOUT_MS = 10000; + +type TerminalWorkerRequestInput = TerminalWorkerRequest extends infer Request + ? Request extends TerminalWorkerRequest + ? Omit + : never + : never; + +interface PendingRequest { + resolve: (value: unknown) => void; + reject: (error: Error) => void; + timeout: ReturnType; +} + +interface WorkerTerminalRecord { + info: WorkerTerminalInfo; + state: TerminalState; + exitInfo: TerminalExitInfo | null; + messageListeners: Set<(msg: ServerMessage) => void>; + exitListeners: Set<(info: TerminalExitInfo) => void>; + commandFinishedListeners: Set<(info: TerminalCommandFinishedInfo) => void>; + titleChangeListeners: Set<(title?: string) => void>; + session: TerminalSession; +} + +interface TerminalWorkerProcess { + connected: boolean; + killed: boolean; + send(message: TerminalWorkerRequest, callback: (error: Error | null) => void): boolean; + disconnect(): void; + kill(): boolean; + on(event: "message", listener: (message: TerminalWorkerToParentMessage) => void): this; + on(event: "exit", listener: (code: number | null, signal: NodeJS.Signals | null) => void): this; +} + +interface WorkerTerminalManagerOptions { + requestTimeoutMs?: number; + forkWorker?: () => TerminalWorkerProcess; +} + +function resolveWorkerUrl(): URL { + const currentUrl = import.meta.url; + if (currentUrl.endsWith(".ts")) { + return new URL("./terminal-worker-process.ts", currentUrl); + } + return new URL("./terminal-worker-process.js", currentUrl); +} + +function resolveWorkerExecArgv(): string[] { + if (!import.meta.url.endsWith(".ts")) { + return []; + } + const loaderUrl = new URL("./terminal-ts-loader.mjs", import.meta.url).href; + const importSource = [ + 'import { register } from "node:module";', + 'import { pathToFileURL } from "node:url";', + `register(${JSON.stringify(loaderUrl)}, pathToFileURL("./"));`, + ].join(" "); + return [ + "--experimental-strip-types", + "--import", + `data:text/javascript,${encodeURIComponent(importSource)}`, + ]; +} + +function isResponse(message: TerminalWorkerToParentMessage): message is TerminalWorkerResponse { + return message.type === "response"; +} + +function cloneTerminalInfo(info: WorkerTerminalInfo): WorkerTerminalInfo { + return { + id: info.id, + name: info.name, + cwd: info.cwd, + ...(info.title ? { title: info.title } : {}), + }; +} + +function forkTerminalWorker(): TerminalWorkerProcess { + return fork(fileURLToPath(resolveWorkerUrl()), [], { + execArgv: resolveWorkerExecArgv(), + serialization: "advanced", + stdio: ["ignore", "ignore", "inherit", "ipc"], + }) as ChildProcess as TerminalWorkerProcess; +} + +export function createWorkerTerminalManager( + managerOptions: WorkerTerminalManagerOptions = {}, +): TerminalManager { + const worker = managerOptions.forkWorker ? managerOptions.forkWorker() : forkTerminalWorker(); + const requestTimeoutMs = managerOptions.requestTimeoutMs ?? REQUEST_TIMEOUT_MS; + const pendingRequests = new Map(); + const recordsById = new Map(); + const terminalIdsByCwd = new Map>(); + const terminalsChangedListeners = new Set(); + let workerExited = false; + let workerShutdownTimer: ReturnType | null = null; + + function emitTerminalsChanged(event: TerminalsChangedEvent): void { + for (const listener of Array.from(terminalsChangedListeners)) { + try { + listener(event); + } catch { + // no-op + } + } + } + + function listTerminalItemsForCwd(cwd: string): TerminalListItem[] { + const terminalIds = terminalIdsByCwd.get(cwd); + if (!terminalIds) { + return []; + } + const terminals: TerminalListItem[] = []; + for (const terminalId of terminalIds) { + const record = recordsById.get(terminalId); + if (!record) { + continue; + } + terminals.push({ + id: record.info.id, + name: record.info.name, + cwd: record.info.cwd, + ...(record.info.title ? { title: record.info.title } : {}), + }); + } + return terminals; + } + + function registerRecord(input: { + info: WorkerTerminalInfo; + state: TerminalState; + }): TerminalSession { + const existing = recordsById.get(input.info.id); + if (existing) { + existing.info = cloneTerminalInfo(input.info); + existing.state = input.state; + return existing.session; + } + + const record: WorkerTerminalRecord = { + info: cloneTerminalInfo(input.info), + state: input.state, + exitInfo: null, + messageListeners: new Set(), + exitListeners: new Set(), + commandFinishedListeners: new Set(), + titleChangeListeners: new Set(), + session: undefined as unknown as TerminalSession, + }; + + const session: TerminalSession = { + get id() { + return record.info.id; + }, + get name() { + return record.info.name; + }, + get cwd() { + return record.info.cwd; + }, + send(message: ClientMessage): void { + sendBestEffortRequest({ type: "send", terminalId: record.info.id, message }); + }, + subscribe(listener: (msg: ServerMessage) => void): () => void { + record.messageListeners.add(listener); + queueMicrotask(() => { + if (record.messageListeners.has(listener)) { + listener({ type: "snapshot", state: record.state }); + } + }); + return () => { + record.messageListeners.delete(listener); + }; + }, + onExit(listener: (info: TerminalExitInfo) => void): () => void { + if (record.exitInfo) { + queueMicrotask(() => listener(record.exitInfo!)); + return () => {}; + } + record.exitListeners.add(listener); + return () => { + record.exitListeners.delete(listener); + }; + }, + onCommandFinished(listener: (info: TerminalCommandFinishedInfo) => void): () => void { + record.commandFinishedListeners.add(listener); + return () => { + record.commandFinishedListeners.delete(listener); + }; + }, + onTitleChange(listener: (title?: string) => void): () => void { + record.titleChangeListeners.add(listener); + if (record.info.title !== undefined) { + queueMicrotask(() => { + if (record.titleChangeListeners.has(listener)) { + listener(record.info.title); + } + }); + } + return () => { + record.titleChangeListeners.delete(listener); + }; + }, + getSize(): { rows: number; cols: number } { + return { + rows: record.state.rows, + cols: record.state.cols, + }; + }, + getState(): TerminalState { + return record.state; + }, + getTitle(): string | undefined { + return record.info.title; + }, + getExitInfo(): TerminalExitInfo | null { + return record.exitInfo; + }, + kill(): void { + sendBestEffortRequest({ type: "killTerminal", terminalId: record.info.id }); + }, + killAndWait(options?: { + gracefulTimeoutMs?: number; + forceTimeoutMs?: number; + }): Promise { + return sendRequest({ + type: "killTerminalAndWait", + terminalId: record.info.id, + ...(options ? { options } : {}), + }).then(() => undefined); + }, + }; + + record.session = session; + recordsById.set(record.info.id, record); + const terminalIds = terminalIdsByCwd.get(record.info.cwd) ?? new Set(); + terminalIds.add(record.info.id); + terminalIdsByCwd.set(record.info.cwd, terminalIds); + return session; + } + + function removeRecord(terminalId: string): WorkerTerminalRecord | undefined { + const record = recordsById.get(terminalId); + if (!record) { + return undefined; + } + recordsById.delete(terminalId); + const terminalIds = terminalIdsByCwd.get(record.info.cwd); + if (terminalIds) { + terminalIds.delete(terminalId); + if (terminalIds.size === 0) { + terminalIdsByCwd.delete(record.info.cwd); + } + } + return record; + } + + function handleWorkerEvent(message: TerminalWorkerToParentMessage): void { + switch (message.type) { + case "terminalCreated": { + registerRecord({ info: message.terminal, state: message.state }); + return; + } + + case "terminalRemoved": { + removeRecord(message.terminalId); + emitTerminalsChanged({ + cwd: message.cwd, + terminals: listTerminalItemsForCwd(message.cwd), + }); + return; + } + + case "terminalMessage": { + const record = recordsById.get(message.terminalId); + if (!record) { + return; + } + if (message.message.type === "snapshot") { + record.state = message.message.state; + } + for (const listener of Array.from(record.messageListeners)) { + listener(message.message); + } + return; + } + + case "terminalExit": { + const record = recordsById.get(message.terminalId); + if (!record) { + return; + } + record.exitInfo = message.info; + for (const listener of Array.from(record.exitListeners)) { + listener(message.info); + } + record.exitListeners.clear(); + removeRecord(message.terminalId); + emitTerminalsChanged({ + cwd: record.info.cwd, + terminals: listTerminalItemsForCwd(record.info.cwd), + }); + return; + } + + case "terminalTitleChange": { + const record = recordsById.get(message.terminalId); + if (!record) { + return; + } + record.info = { + ...record.info, + ...(message.title ? { title: message.title } : { title: undefined }), + }; + for (const listener of Array.from(record.titleChangeListeners)) { + listener(message.title); + } + emitTerminalsChanged({ + cwd: record.info.cwd, + terminals: listTerminalItemsForCwd(record.info.cwd), + }); + return; + } + + case "terminalCommandFinished": { + const record = recordsById.get(message.terminalId); + if (!record) { + return; + } + for (const listener of Array.from(record.commandFinishedListeners)) { + listener(message.info); + } + return; + } + + case "terminalsChanged": { + emitTerminalsChanged({ + cwd: message.cwd, + terminals: message.terminals.map((terminal) => ({ + id: terminal.id, + name: terminal.name, + cwd: terminal.cwd, + ...(terminal.title ? { title: terminal.title } : {}), + })), + }); + return; + } + } + } + + function rejectPendingRequests(error: Error): void { + for (const [requestId, pending] of pendingRequests) { + clearTimeout(pending.timeout); + pending.reject(error); + pendingRequests.delete(requestId); + } + } + + worker.on("message", (message: TerminalWorkerToParentMessage) => { + if (isResponse(message)) { + const pending = pendingRequests.get(message.requestId); + if (!pending) { + return; + } + clearTimeout(pending.timeout); + pendingRequests.delete(message.requestId); + if (message.ok) { + pending.resolve(message.result); + } else { + pending.reject(new Error(message.error)); + } + return; + } + handleWorkerEvent(message); + }); + + worker.on("exit", (code, signal) => { + workerExited = true; + if (workerShutdownTimer) { + clearTimeout(workerShutdownTimer); + workerShutdownTimer = null; + } + rejectPendingRequests(new Error(`Terminal worker exited (${signal ?? code ?? "unknown"})`)); + }); + + function sendRequest(input: TerminalWorkerRequestInput): Promise { + if (workerExited || !worker.connected) { + return Promise.reject(new Error("Terminal worker is not running")); + } + const requestId = randomUUID(); + const message = { ...input, requestId } as TerminalWorkerRequest; + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + pendingRequests.delete(requestId); + reject(new Error(`Terminal worker request timed out: ${input.type}`)); + }, requestTimeoutMs); + pendingRequests.set(requestId, { resolve, reject, timeout }); + worker.send(message, (error) => { + if (!error) { + return; + } + clearTimeout(timeout); + pendingRequests.delete(requestId); + reject(error); + }); + }); + } + + function sendBestEffortRequest(input: TerminalWorkerRequestInput): void { + void sendRequest(input).catch(() => { + // The public terminal methods that call this are intentionally synchronous. + // Worker failures are surfaced through awaitable manager methods and worker + // lifecycle state; do not let fire-and-forget sends crash the daemon. + }); + } + + function toSessions(terminals: WorkerTerminalInfo[]): TerminalSession[] { + return terminals + .map((terminal) => recordsById.get(terminal.id)?.session) + .filter((session): session is TerminalSession => Boolean(session)); + } + + return { + async getTerminals(cwd: string): Promise { + const result = (await sendRequest({ type: "getTerminals", cwd })) as WorkerTerminalInfo[]; + return toSessions(result); + }, + + async createTerminal(options: WorkerCreateTerminalOptions): Promise { + const result = (await sendRequest({ type: "createTerminal", options })) as { + terminal: WorkerTerminalInfo; + state: TerminalState; + }; + return registerRecord({ info: result.terminal, state: result.state }); + }, + + registerCwdEnv(options: { cwd: string; env: Record }): void { + sendBestEffortRequest({ + type: "registerCwdEnv", + cwd: options.cwd, + env: options.env, + }); + }, + + getTerminal(id: string): TerminalSession | undefined { + return recordsById.get(id)?.session; + }, + + killTerminal(id: string): void { + void sendRequest({ type: "killTerminal", terminalId: id }).catch(() => { + // no-op; kill is intentionally best-effort and synchronous in the public interface. + }); + }, + + async killTerminalAndWait( + id: string, + options?: { gracefulTimeoutMs?: number; forceTimeoutMs?: number }, + ): Promise { + await sendRequest({ + type: "killTerminalAndWait", + terminalId: id, + ...(options ? { options } : {}), + }); + }, + + listDirectories(): string[] { + return Array.from(terminalIdsByCwd.keys()); + }, + + killAll(): void { + void sendRequest({ type: "killAll" }) + .catch(() => { + // no-op + }) + .finally(() => { + if (worker.connected) { + worker.disconnect(); + } + if (!worker.killed && !workerShutdownTimer) { + workerShutdownTimer = setTimeout(() => { + worker.kill(); + }, 1000); + } + }); + for (const terminalId of Array.from(recordsById.keys())) { + removeRecord(terminalId); + } + }, + + subscribeTerminalsChanged(listener: TerminalsChangedListener): () => void { + terminalsChangedListeners.add(listener); + return () => { + terminalsChangedListeners.delete(listener); + }; + }, + }; +} + +export function terminateWorkerTerminalManager(manager: TerminalManager): void { + manager.killAll(); +} diff --git a/packages/server/tsconfig.server.json b/packages/server/tsconfig.server.json index 8b6228b86..aee1b6367 100644 --- a/packages/server/tsconfig.server.json +++ b/packages/server/tsconfig.server.json @@ -22,7 +22,7 @@ "declarationMap": true, "sourceMap": true }, - "include": ["src/server/**/*", "src/services/**/*"], + "include": ["src/server/**/*", "src/services/**/*", "src/terminal/**/*"], "exclude": [ "node_modules", "dist",