mirror of
https://github.com/getpaseo/paseo.git
synced 2026-07-29 12:01:31 +00:00
Move terminal sessions to worker process
This commit is contained in:
@@ -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",
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
18
packages/server/src/terminal/terminal-manager-factory.ts
Normal file
18
packages/server/src/terminal/terminal-manager-factory.ts
Normal file
@@ -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();
|
||||
}
|
||||
20
packages/server/src/terminal/terminal-ts-loader.mjs
Normal file
20
packages/server/src/terminal/terminal-ts-loader.mjs
Normal file
@@ -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);
|
||||
}
|
||||
213
packages/server/src/terminal/terminal-worker-process.ts
Normal file
213
packages/server/src/terminal/terminal-worker-process.ts
Normal file
@@ -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<string, Array<() => 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<void> {
|
||||
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();
|
||||
});
|
||||
122
packages/server/src/terminal/terminal-worker-protocol.ts
Normal file
122
packages/server/src/terminal/terminal-worker-protocol.ts
Normal file
@@ -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<string, string>;
|
||||
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<string, string>;
|
||||
}
|
||||
| {
|
||||
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;
|
||||
211
packages/server/src/terminal/worker-terminal-manager.test.ts
Normal file
211
packages/server/src/terminal/worker-terminal-manager.test.ts
Normal file
@@ -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<void> {
|
||||
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<T>(shell: string, run: () => Promise<T>): Promise<T> {
|
||||
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");
|
||||
});
|
||||
});
|
||||
527
packages/server/src/terminal/worker-terminal-manager.ts
Normal file
527
packages/server/src/terminal/worker-terminal-manager.ts
Normal file
@@ -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<Request, "requestId">
|
||||
: never
|
||||
: never;
|
||||
|
||||
interface PendingRequest {
|
||||
resolve: (value: unknown) => void;
|
||||
reject: (error: Error) => void;
|
||||
timeout: ReturnType<typeof setTimeout>;
|
||||
}
|
||||
|
||||
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<string, PendingRequest>();
|
||||
const recordsById = new Map<string, WorkerTerminalRecord>();
|
||||
const terminalIdsByCwd = new Map<string, Set<string>>();
|
||||
const terminalsChangedListeners = new Set<TerminalsChangedListener>();
|
||||
let workerExited = false;
|
||||
let workerShutdownTimer: ReturnType<typeof setTimeout> | 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<void> {
|
||||
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<string>();
|
||||
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<unknown> {
|
||||
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<TerminalSession[]> {
|
||||
const result = (await sendRequest({ type: "getTerminals", cwd })) as WorkerTerminalInfo[];
|
||||
return toSessions(result);
|
||||
},
|
||||
|
||||
async createTerminal(options: WorkerCreateTerminalOptions): Promise<TerminalSession> {
|
||||
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<string, string> }): 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<void> {
|
||||
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();
|
||||
}
|
||||
@@ -22,7 +22,7 @@
|
||||
"declarationMap": true,
|
||||
"sourceMap": true
|
||||
},
|
||||
"include": ["src/server/**/*", "src/services/**/*"],
|
||||
"include": ["src/server/**/*", "src/services/**/*", "src/terminal/**/*"],
|
||||
"exclude": [
|
||||
"node_modules",
|
||||
"dist",
|
||||
|
||||
Reference in New Issue
Block a user