diff --git a/apps/server/src/terminal/NodePtyAdapter.test.ts b/apps/server/src/terminal/NodePtyAdapter.test.ts index e6650025f70f..49eca9b6955b 100644 --- a/apps/server/src/terminal/NodePtyAdapter.test.ts +++ b/apps/server/src/terminal/NodePtyAdapter.test.ts @@ -1,23 +1,61 @@ +import * as NodeEvents from "node:events"; +import * as NodeNet from "node:net"; + import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import { vi } from "vite-plus/test"; import * as NodePtyAdapter from "./NodePtyAdapter.ts"; import * as PtyAdapter from "./PtyAdapter.ts"; -const spawn = vi.fn(() => ({ - pid: 42, - write: vi.fn(), - resize: vi.fn(), - kill: vi.fn(), - onData: vi.fn(() => ({ dispose: vi.fn() })), - onExit: vi.fn(() => ({ dispose: vi.fn() })), -})); +function makeNativeProcess(pid = 42) { + const events = new NodeEvents.EventEmitter(); + return { + pid, + _socket: new NodeNet.Socket(), + _agent: { kill: vi.fn() }, + write: vi.fn(), + resize: vi.fn(), + kill: vi.fn(), + onData: vi.fn((callback: (data: string) => void) => { + events.on("data", callback); + return { + dispose: () => { + events.off("data", callback); + }, + }; + }), + onExit: vi.fn((callback: (event: { exitCode: number; signal?: number }) => void) => { + events.on("exit", callback); + return { + dispose: () => { + events.off("exit", callback); + }, + }; + }), + events, + }; +} + +const spawn = vi.fn(() => makeNativeProcess()); + +function preparePendingProcess() { + const nativeProcess = makeNativeProcess(0); + const subscribed = Promise.withResolvers(); + nativeProcess._socket.on("newListener", (event) => { + if (event === "ready_datapipe") queueMicrotask(() => subscribed.resolve()); + }); + spawn.mockReturnValueOnce(nativeProcess); + return { nativeProcess, subscribed: Effect.promise(() => subscribed.promise) }; +} + +const spawnInput = { shell: "powershell.exe", cwd: ".", cols: 80, rows: 24, env: {} }; const fakeNodePty = { spawn } as unknown as typeof import("node-pty"); @@ -35,6 +73,91 @@ const makeTestLayer = (platform: NodeJS.Platform = "win32") => const testLayer = makeTestLayer(); +it.effect("waits for the Windows PID without requiring output", () => + Effect.gen(function* () { + const { nativeProcess, subscribed } = preparePendingProcess(); + const adapter = yield* PtyAdapter.PtyAdapter; + let completed = false; + const fiber = yield* adapter.spawn(spawnInput).pipe( + Effect.tap(() => + Effect.sync(() => { + completed = true; + }), + ), + Effect.forkChild, + ); + yield* subscribed; + assert.isFalse(completed); + nativeProcess.pid = 12345; + nativeProcess._socket.emit("ready_datapipe"); + const process = yield* Fiber.join(fiber); + assert.equal(process.pid, 12345); + assert.equal(nativeProcess._socket.listenerCount("ready_datapipe"), 0); + assert.equal(nativeProcess.events.listenerCount("exit"), 0); + + const output: string[] = []; + const exits: PtyAdapter.PtyExitEvent[] = []; + const stopData = process.onData((data) => output.push(data)); + const stopExit = process.onExit((event) => exits.push(event)); + nativeProcess.events.emit("data", "first output"); + nativeProcess.events.emit("exit", { exitCode: 0 }); + assert.deepEqual(output, ["first output"]); + assert.deepEqual(exits, [{ exitCode: 0, signal: null }]); + stopData(); + stopExit(); + }).pipe(Effect.provide(testLayer)), +); + +for (const failure of ["exit", "close", "error", "invalid-pid"] as const) { + it.effect(`fails Windows startup on ${failure} and cleans up`, () => + Effect.gen(function* () { + const { nativeProcess, subscribed } = preparePendingProcess(); + const adapter = yield* PtyAdapter.PtyAdapter; + const fiber = yield* adapter.spawn(spawnInput).pipe(Effect.result, Effect.forkChild); + yield* subscribed; + if (failure === "exit") nativeProcess.events.emit("exit", { exitCode: 1 }); + else if (failure === "error") nativeProcess._socket.emit("error", new Error("pipe failed")); + else nativeProcess._socket.emit(failure === "close" ? "close" : "ready_datapipe"); + const result = yield* Fiber.join(fiber); + assert.equal(result._tag, "Failure"); + if (result._tag === "Failure") assert.instanceOf(result.failure, PtyAdapter.PtySpawnError); + assert.equal(nativeProcess._socket.listenerCount("ready_datapipe"), 0); + assert.equal(nativeProcess._socket.listenerCount("error"), 0); + assert.equal(nativeProcess._socket.listenerCount("close"), 0); + assert.equal(nativeProcess.events.listenerCount("exit"), 0); + assert.equal(nativeProcess._agent.kill.mock.calls.length, 1); + }).pipe(Effect.provide(testLayer)), + ); +} + +it.effect("cancels the Windows connection without waiting for output", () => + Effect.gen(function* () { + const { nativeProcess, subscribed } = preparePendingProcess(); + const adapter = yield* PtyAdapter.PtyAdapter; + const fiber = yield* adapter.spawn(spawnInput).pipe(Effect.forkChild); + yield* subscribed; + yield* Fiber.interrupt(fiber); + assert.equal(nativeProcess._agent.kill.mock.calls.length, 1); + assert.equal(nativeProcess.kill.mock.calls.length, 0); + assert.equal(nativeProcess._socket.listenerCount("ready_datapipe"), 0); + assert.equal(nativeProcess.events.listenerCount("exit"), 0); + }).pipe(Effect.provide(testLayer)), +); + +it.effect("reports an incompatible Windows readiness API instead of hanging", () => + Effect.gen(function* () { + const nativeProcess = makeNativeProcess(0); + Reflect.deleteProperty(nativeProcess, "_socket"); + spawn.mockReturnValueOnce(nativeProcess); + const adapter = yield* PtyAdapter.PtyAdapter; + const error = yield* adapter.spawn(spawnInput).pipe(Effect.flip); + assert.instanceOf(error, PtyAdapter.PtySpawnError); + assert.instanceOf(error.cause, Error); + assert.equal(error.cause.message, "Windows PTY readiness socket is unavailable."); + assert.equal(nativeProcess._agent.kill.mock.calls.length, 1); + }).pipe(Effect.provide(testLayer)), +); + for (const platform of ["win32", "linux", "darwin"] as const) { it.effect(`terminates through node-pty using ${platform} semantics`, () => Effect.gen(function* () { diff --git a/apps/server/src/terminal/NodePtyAdapter.ts b/apps/server/src/terminal/NodePtyAdapter.ts index 67cdcecdd53a..0ab308c537f2 100644 --- a/apps/server/src/terminal/NodePtyAdapter.ts +++ b/apps/server/src/terminal/NodePtyAdapter.ts @@ -1,4 +1,5 @@ import * as NodeModule from "node:module"; +import * as NodeNet from "node:net"; import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; @@ -82,6 +83,76 @@ const ensureNodePtySpawnHelperExecutable = Effect.fn(function* () { yield* fs.chmod(helperPath, 0o755).pipe(Effect.orElseSucceed(() => undefined)); }); +// node-pty now defers Windows process creation to avoid blocking on named pipes: +// https://github.com/microsoft/node-pty/pull/885 +// T3 adopted that behavior when upgrading from 1.1.0 to 1.2.0-beta.15: +// https://github.com/pingdotgg/t3code/pull/13748 +// Its public API has no readiness event. The private ready_datapipe handler sets +// pid before our listener runs; wait here so the manager always receives a real PID. +const waitForWindowsPid = (process: import("node-pty").IPty, shell: string) => + Effect.callback((resume) => { + const hasPid = () => Number.isInteger(process.pid) && process.pid > 0; + const failure = (cause: unknown) => + Effect.fail(new PtyAdapter.PtySpawnError({ adapter: "node-pty", shell, cause })); + + if (hasPid()) { + resume(Effect.void); + return; + } + + if (!("_socket" in process) || !(process._socket instanceof NodeNet.Socket)) { + resume(failure(new Error("Windows PTY readiness socket is unavailable."))); + return; + } + + const socket = process._socket; + const onReady = () => { + cleanup(); + resume( + hasPid() + ? Effect.void + : failure(new Error("Windows PTY became ready without a valid PID.")), + ); + }; + const onError = (cause: Error) => { + cleanup(); + resume(failure(cause)); + }; + const onClose = () => onError(new Error("Windows PTY closed before its PID was available.")); + const exitListener = process.onExit(({ exitCode }) => + onError( + new Error(`Windows PTY exited before its PID was available (exit code ${exitCode}).`), + ), + ); + const cleanup = () => { + socket.off("ready_datapipe", onReady); + socket.off("error", onError); + socket.off("close", onClose); + exitListener.dispose(); + }; + socket.once("ready_datapipe", onReady); + socket.once("error", onError); + socket.once("close", onClose); + return Effect.sync(cleanup); + }); + +const killStartingWindowsPty = (process: import("node-pty").IPty) => + Effect.try(() => { + // Public kill() waits for the first output on Windows. The agent can cancel + // the pending connection even when no child or output exists yet. + if ( + "_agent" in process && + typeof process._agent === "object" && + process._agent !== null && + "kill" in process._agent && + typeof process._agent.kill === "function" + ) { + process._agent.kill(); + } else { + process.kill(); + } + }).pipe(Effect.ignore); + class NodePtyProcess implements PtyAdapter.PtyProcess { private readonly process: import("node-pty").IPty; private readonly platform: NodeJS.Platform; @@ -181,6 +252,11 @@ export const make = Effect.fn("NodePtyAdapter.make")(function* () { cause, }), }); + if (platform === "win32") { + yield* waitForWindowsPid(ptyProcess, input.shell).pipe( + Effect.onError(() => killStartingWindowsPty(ptyProcess)), + ); + } return new NodePtyProcess(ptyProcess, platform); }), });