diff --git a/apps/server/src/terminal/Manager.test.ts b/apps/server/src/terminal/Manager.test.ts index f04161cfbd3c..325233bf59d9 100644 --- a/apps/server/src/terminal/Manager.test.ts +++ b/apps/server/src/terminal/Manager.test.ts @@ -56,6 +56,7 @@ class FakePtyProcess implements PtyAdapter.PtyProcess { private readonly dataListeners = new Set<(data: string) => void>(); private readonly exitListeners = new Set<(event: PtyAdapter.PtyExitEvent) => void>(); killed = false; + exitOnSubscribe: PtyAdapter.PtyExitEvent | undefined; constructor(pid: number) { this.pid = pid; @@ -88,6 +89,7 @@ class FakePtyProcess implements PtyAdapter.PtyProcess { } onExit(callback: (event: PtyAdapter.PtyExitEvent) => void): () => void { + if (this.exitOnSubscribe) callback(this.exitOnSubscribe); this.exitListeners.add(callback); return () => { this.exitListeners.delete(callback); @@ -113,6 +115,7 @@ class FakePtyAdapter { readonly spawnFailures: Error[] = []; private readonly mode: "sync" | "async"; private nextPid = 9000; + exitOnSubscribe: PtyAdapter.PtyExitEvent | undefined; constructor(mode: "sync" | "async" = "sync") { this.mode = mode; @@ -133,6 +136,7 @@ class FakePtyAdapter { ); } const process = new FakePtyProcess(this.nextPid++); + process.exitOnSubscribe = this.exitOnSubscribe; this.processes.push(process); if (this.mode === "async") { return Effect.tryPromise({ @@ -628,6 +632,36 @@ it.layer( }), ); + it.effect("handles an exit replayed during subscription after publishing startup", () => + Effect.gen(function* () { + const ptyAdapter = new FakePtyAdapter(); + ptyAdapter.exitOnSubscribe = { exitCode: 7, signal: null }; + const { manager, getEvents } = yield* createManager(5, { ptyAdapter }); + const exited = yield* Deferred.make(); + const unsubscribe = yield* manager.subscribe((event) => + event.type === "exited" + ? Deferred.succeed(exited, undefined).pipe(Effect.asVoid) + : Effect.void, + ); + yield* Effect.addFinalizer(() => Effect.sync(unsubscribe)); + yield* manager.open(openInput()); + yield* Deferred.await(exited); + const events = yield* getEvents; + expect(events.map((event) => event.type)).toEqual(["started", "exited"]); + expect(events[1]).toMatchObject({ exitCode: 7 }); + const attached: TerminalAttachStreamEvent[] = []; + const stopAttach = yield* manager.attachStream(openInput(), (event) => + Effect.sync(() => { + attached.push(event); + }), + ); + yield* Effect.addFinalizer(() => Effect.sync(stopAttach)); + expect(attached.find((event) => event.type === "snapshot")).toMatchObject({ + snapshot: { status: "exited", exitCode: 7 }, + }); + }), + ); + it.effect("supports asynchronous PTY spawn effects", () => Effect.gen(function* () { const { manager, ptyAdapter } = yield* createManager(5, { diff --git a/apps/server/src/terminal/Manager.ts b/apps/server/src/terminal/Manager.ts index 75d592c00c0d..43d6a0750e65 100644 --- a/apps/server/src/terminal/Manager.ts +++ b/apps/server/src/terminal/Manager.ts @@ -2243,18 +2243,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func startedShell = spawnResult.shellLabel; const processPid = ptyProcess.pid; - const unsubscribeData = ptyProcess.onData((data) => { - if (!enqueueProcessEvent(session, processPid, { type: "output", data })) { - return; - } - runFork(drainProcessEvents(session, processPid)); - }); - const unsubscribeExit = ptyProcess.onExit((event) => { - if (!enqueueProcessEvent(session, processPid, { type: "exit", event })) { - return; - } - runFork(drainProcessEvents(session, processPid)); - }); + let eventsActivated = false; let eventStamp: ReturnType = { updatedAt: session.updatedAt, @@ -2264,8 +2253,19 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func session.process = ptyProcess; session.pid = processPid; session.status = "running"; - session.unsubscribeData = unsubscribeData; - session.unsubscribeExit = unsubscribeExit; + // onExit may replay an exit immediately; accept it before subscribing. + session.unsubscribeData = spawnResult.process.onData((data) => { + if (!enqueueProcessEvent(session, processPid, { type: "output", data })) { + return; + } + if (eventsActivated) runFork(drainProcessEvents(session, processPid)); + }); + session.unsubscribeExit = spawnResult.process.onExit((event) => { + if (!enqueueProcessEvent(session, processPid, { type: "exit", event })) { + return; + } + if (eventsActivated) runFork(drainProcessEvents(session, processPid)); + }); eventStamp = advanceEventSequence(session); return [undefined, state] as const; }); @@ -2277,6 +2277,9 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func sequence: eventStamp.sequence, snapshot: snapshot(session), }); + // Publish startup before draining any events replayed during subscription. + eventsActivated = true; + if (session.processEventDrainRunning) runFork(drainProcessEvents(session, processPid)); }), ), ), diff --git a/apps/server/src/terminal/NodePtyAdapter.test.ts b/apps/server/src/terminal/NodePtyAdapter.test.ts index e6650025f70f..8107e19e6165 100644 --- a/apps/server/src/terminal/NodePtyAdapter.test.ts +++ b/apps/server/src/terminal/NodePtyAdapter.test.ts @@ -1,23 +1,63 @@ +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 Logger from "effect/Logger"; +import * as Scheduler from "effect/Scheduler"; +import { expect, 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 +75,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"), 1); + + 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* () { @@ -153,3 +278,88 @@ it.effect("reports native module load failures as structured startup defects", ( ), ), ); + +for (const budget of [2048, 8]) { + it.effect(`preserves an exit during readiness handoff with scheduler budget ${budget}`, () => + Effect.gen(function* () { + const { nativeProcess, subscribed } = preparePendingProcess(); + const adapter = yield* PtyAdapter.PtyAdapter; + const exits: PtyAdapter.PtyExitEvent[] = []; + const fiber = yield* Effect.gen(function* () { + const process = yield* adapter.spawn(spawnInput); + process.onExit((event) => exits.push(event)); + }).pipe( + Effect.provideService(Scheduler.MaxOpsBeforeYield, budget), + Effect.provideService(Scheduler.PreventSchedulerYield, false), + Effect.forkChild, + ); + yield* subscribed; + nativeProcess.pid = 12345; + nativeProcess._socket.emit("ready_datapipe"); + nativeProcess.events.emit("exit", { exitCode: 0 }); + yield* Fiber.join(fiber); + assert.equal(exits.length, 1); + }).pipe(Effect.provide(testLayer)), + ); +} + +it.effect("replays an exit to late subscribers and respects unsubscription", () => + Effect.gen(function* () { + const adapter = yield* PtyAdapter.PtyAdapter; + const process = yield* adapter.spawn(spawnInput); + const nativeProcess = spawn.mock.results.at(-1)!.value; + const removed = vi.fn(); + process.onExit(removed)(); + nativeProcess.events.emit("exit", { exitCode: 7, signal: 2 }); + const late = vi.fn(); + process.onExit(late); + nativeProcess.events.emit("exit", { exitCode: 9 }); + assert.equal(removed.mock.calls.length, 0); + assert.deepEqual(late.mock.calls, [[{ exitCode: 7, signal: 2 }]]); + assert.equal(nativeProcess.events.listenerCount("exit"), 0); + }).pipe(Effect.provide(testLayer)), +); + +for (const failure of ["spawn", "interrupt"] as const) { + it.effect(`logs cleanup failures without replacing ${failure}`, () => + Effect.gen(function* () { + const { nativeProcess, subscribed } = preparePendingProcess(); + const killError = new Error("native kill failed"); + nativeProcess._agent.kill.mockImplementation(() => { + throw killError; + }); + const messages: unknown[] = []; + const logger = Logger.make(({ message }) => { + messages.push(message); + }); + const adapter = yield* PtyAdapter.PtyAdapter; + const fiber = yield* adapter + .spawn(spawnInput) + .pipe( + Effect.provide(Logger.layer([logger], { mergeWithExisting: false })), + Effect.forkChild, + ); + yield* subscribed; + const spawnError = new Error("pipe failed"); + if (failure === "interrupt") yield* Fiber.interrupt(fiber); + else nativeProcess._socket.emit("error", spawnError); + const exit = yield* Fiber.await(fiber); + assert.isTrue(Exit.isFailure(exit)); + if (Exit.isFailure(exit)) { + if (failure === "interrupt") assert.isTrue(Cause.hasInterrupts(exit.cause)); + else { + const error = Cause.squash(exit.cause); + assert.instanceOf(error, PtyAdapter.PtySpawnError); + assert.equal(error.cause, spawnError); + } + } + assert.equal(messages.length, 1); + expect(messages[0]).toMatchObject([ + "failed to cancel Windows terminal startup", + { terminalPid: 0, cause: { cause: killError } }, + ]); + assert.equal(nativeProcess.events.listenerCount("exit"), 0); + assert.equal(nativeProcess._socket.listenerCount("ready_datapipe"), 0); + }).pipe(Effect.provide(testLayer)), + ); +} diff --git a/apps/server/src/terminal/NodePtyAdapter.ts b/apps/server/src/terminal/NodePtyAdapter.ts index 67cdcecdd53a..8ba78ee4d288 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,13 +83,112 @@ const ensureNodePtySpawnHelperExecutable = Effect.fn(function* () { yield* fs.chmod(helperPath, 0o755).pipe(Effect.orElseSucceed(() => undefined)); }); +/** + * Waits for Windows process creation so the manager receives a valid PID. + * node-pty defers 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. + */ +const waitForWindowsPid = ( + process: import("node-pty").IPty, + trackedProcess: NodePtyProcess, + 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.")); + let stopExit = () => {}; + const cleanup = () => { + socket.off("ready_datapipe", onReady); + socket.off("error", onError); + socket.off("close", onClose); + stopExit(); + }; + socket.once("ready_datapipe", onReady); + socket.once("error", onError); + socket.once("close", onClose); + stopExit = trackedProcess.onExit(({ exitCode }) => + onError( + new Error(`Windows PTY exited before its PID was available (exit code ${exitCode}).`), + ), + ); + return Effect.sync(cleanup); + }); + +/** + * Cancels Windows startup without waiting for the first output, unlike public kill(). + * The private agent can cancel the pending connection before a child exists. + * Cleanup failures are logged without replacing the startup failure. + */ +const killStartingWindowsPty = (process: import("node-pty").IPty) => + Effect.try(() => { + 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.catch((error) => + Effect.logWarning("failed to cancel Windows terminal startup", { + terminalPid: process.pid, + cause: error, + }), + ), + ); + class NodePtyProcess implements PtyAdapter.PtyProcess { private readonly process: import("node-pty").IPty; private readonly platform: NodeJS.Platform; + private exitEvent: PtyAdapter.PtyExitEvent | undefined; + private readonly exitListeners = new Set<(event: PtyAdapter.PtyExitEvent) => void>(); + private readonly exitSubscription: import("node-pty").IDisposable; constructor(process: import("node-pty").IPty, platform: NodeJS.Platform) { this.process = process; this.platform = platform; + // Retain exits while Windows readiness and the manager hand off the process. + this.exitSubscription = process.onExit((event) => { + if (this.exitEvent) return; + this.exitEvent = { exitCode: event.exitCode, signal: event.signal ?? null }; + this.exitSubscription.dispose(); + for (const listener of this.exitListeners) listener(this.exitEvent); + this.exitListeners.clear(); + }); } get pid(): number { @@ -116,16 +216,20 @@ class NodePtyProcess implements PtyAdapter.PtyProcess { } onExit(callback: (event: PtyAdapter.PtyExitEvent) => void): () => void { - const disposable = this.process.onExit((event) => { - callback({ - exitCode: event.exitCode, - signal: event.signal ?? null, - }); - }); + if (this.exitEvent) { + callback(this.exitEvent); + return () => {}; + } + this.exitListeners.add(callback); return () => { - disposable.dispose(); + this.exitListeners.delete(callback); }; } + + disposeExitSubscription(): void { + this.exitSubscription.dispose(); + this.exitListeners.clear(); + } } export const make = Effect.fn("NodePtyAdapter.make")(function* () { @@ -166,14 +270,16 @@ export const make = Effect.fn("NodePtyAdapter.make")(function* () { ? { ...input.env, TERM: "xterm-256color" } : input.env; const ptyProcess = yield* Effect.try({ - try: () => - nodePty.spawn(input.shell, input.args ?? [], { + try: () => { + const nativeProcess = nodePty.spawn(input.shell, input.args ?? [], { cwd: input.cwd, cols: input.cols, rows: input.rows, env, name: "xterm-256color", - }), + }); + return { nativeProcess, process: new NodePtyProcess(nativeProcess, platform) }; + }, catch: (cause) => new PtyAdapter.PtySpawnError({ adapter: "node-pty", @@ -181,7 +287,16 @@ export const make = Effect.fn("NodePtyAdapter.make")(function* () { cause, }), }); - return new NodePtyProcess(ptyProcess, platform); + if (platform === "win32") { + yield* waitForWindowsPid(ptyProcess.nativeProcess, ptyProcess.process, input.shell).pipe( + Effect.onError(() => + Effect.sync(() => ptyProcess.process.disposeExitSubscription()).pipe( + Effect.andThen(killStartingWindowsPty(ptyProcess.nativeProcess)), + ), + ), + ); + } + return ptyProcess.process; }), }); });