From cfc3f384d77fef104f269db5ca089f9b3f6e95ee Mon Sep 17 00:00:00 2001 From: SegFaultZero Date: Sun, 27 Sep 2026 12:44:03 +0800 Subject: [PATCH 1/2] fix(server): handle asynchronous Windows terminal process IDs --- apps/server/src/terminal/Manager.test.ts | 82 +++++++++++++++++++++++- apps/server/src/terminal/Manager.ts | 43 ++++++++----- 2 files changed, 107 insertions(+), 18 deletions(-) diff --git a/apps/server/src/terminal/Manager.test.ts b/apps/server/src/terminal/Manager.test.ts index f04161cfbd3c..022577fe8132 100644 --- a/apps/server/src/terminal/Manager.test.ts +++ b/apps/server/src/terminal/Manager.test.ts @@ -11,6 +11,8 @@ import { ProviderInstanceId, ServerSettingsError, TerminalProviderInstanceNotFoundError, + TerminalSessionSnapshot, + TerminalSummary, } from "@t3tools/contracts"; import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as Data from "effect/Data"; @@ -28,6 +30,7 @@ import * as PlatformError from "effect/PlatformError"; import * as Path from "effect/Path"; import * as Ref from "effect/Ref"; import * as Schedule from "effect/Schedule"; +import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; @@ -42,6 +45,9 @@ import * as ServerSettings from "../serverSettings.ts"; import * as TerminalManager from "./Manager.ts"; import * as PtyAdapter from "./PtyAdapter.ts"; +const encodeTerminalSnapshot = Schema.encodeEffect(TerminalSessionSnapshot); +const encodeTerminalSummary = Schema.encodeEffect(TerminalSummary); + class WaitForConditionError extends Data.TaggedError("WaitForConditionError")<{ readonly message: string; }> {} @@ -50,7 +56,7 @@ class FakePtyProcess implements PtyAdapter.PtyProcess { readonly writes: string[] = []; readonly resizeCalls: Array<{ cols: number; rows: number }> = []; readonly killSignals: Array = []; - readonly pid: number; + pid: number; writeFailure: unknown | undefined; resizeFailure: unknown | undefined; private readonly dataListeners = new Set<(data: string) => void>(); @@ -112,7 +118,7 @@ class FakePtyAdapter { readonly processes: FakePtyProcess[] = []; readonly spawnFailures: Error[] = []; private readonly mode: "sync" | "async"; - private nextPid = 9000; + nextPid = 9000; constructor(mode: "sync" | "async" = "sync") { this.mode = mode; @@ -642,6 +648,78 @@ it.layer( }), ); + it.effect("keeps Windows terminals usable while ConPTY assigns the PID", () => + Effect.gen(function* () { + const adapter = new FakePtyAdapter(); + adapter.nextPid = 0; + const inspectedPids: number[] = []; + const { manager } = yield* createManager(5, { + ptyAdapter: adapter, + subprocessInspector: (pid) => + Effect.sync(() => { + inspectedPids.push(pid); + return { hasRunningSubprocess: true, childCommand: "node", processIds: [pid] }; + }), + }); + const readMetadata = Effect.gen(function* () { + let terminals: ReadonlyArray = []; + const unsubscribe = yield* manager.subscribeMetadata((event) => + Effect.sync(() => { + if (event.type === "snapshot") terminals = event.terminals; + }), + ); + unsubscribe(); + return terminals; + }); + const initial = yield* manager.open(openInput()); + assert.equal(initial.pid, null); + yield* encodeTerminalSnapshot(initial); + + const metadata = yield* readMetadata; + assert.equal(metadata[0]?.pid, null); + yield* Effect.forEach(metadata, (item) => encodeTerminalSummary(item)); + yield* manager.closeIdle({ threadId: "thread-1" }); + expect(inspectedPids).toEqual([]); + + const attached = yield* Ref.make>([]); + const unsubscribeAttach = yield* manager.attachStream(openInput(), (event) => + Ref.update(attached, (events) => [...events, event]), + ); + yield* Effect.addFinalizer(() => Effect.sync(unsubscribeAttach)); + const output = yield* Deferred.make(); + const exited = yield* Deferred.make(); + const unsubscribe = yield* manager.subscribe((event) => + event.type === "output" + ? Deferred.succeed(output, undefined).pipe(Effect.asVoid) + : event.type === "exited" + ? Deferred.succeed(exited, undefined).pipe(Effect.asVoid) + : Effect.void, + ); + yield* Effect.addFinalizer(() => Effect.sync(unsubscribe)); + + const process = adapter.processes[0]!; + process.pid = 12345; + process.emitData("Windows prompt> "); + yield* Deferred.await(output); + const ready = yield* manager.open(openInput()); + assert.equal(ready.pid, 12345); + assert.equal(ready.history, "Windows prompt> "); + assert.equal((yield* readMetadata)[0]?.pid, 12345); + yield* manager.closeIdle({ threadId: "thread-1" }); + expect(inspectedPids).toContain(12345); + expect(inspectedPids).not.toContain(0); + expect(yield* Ref.get(attached)).toContainEqual( + expect.objectContaining({ type: "output", data: "Windows prompt> " }), + ); + + process.emitExit({ exitCode: 0, signal: null }); + yield* Deferred.await(exited); + const finalMetadata = yield* readMetadata; + assert.equal(finalMetadata[0]?.status, "exited"); + assert.equal(finalMetadata[0]?.pid, null); + }), + ); + it.effect("forwards write and resize to active pty process", () => Effect.gen(function* () { const { manager, ptyAdapter } = yield* createManager(); diff --git a/apps/server/src/terminal/Manager.ts b/apps/server/src/terminal/Manager.ts index 75d592c00c0d..d9f0696a02c4 100644 --- a/apps/server/src/terminal/Manager.ts +++ b/apps/server/src/terminal/Manager.ts @@ -276,7 +276,7 @@ interface TerminalSessionState { cwd: string; worktreePath: string | null; status: TerminalSessionStatus; - pid: number | null; + readonly pid: number | null; history: BoundedTerminalHistory; pendingHistoryControlSequence: string; pendingProcessEvents: Array; @@ -366,6 +366,13 @@ function terminalWireLabel(session: TerminalSessionState): string { return truncateTerminalWireLabel(getTerminalLabel(session.terminalId)); } +// ConPTY assigns the PID asynchronously. Read it from the process instead of +// freezing its startup sentinel into snapshots and subprocess checks. +function terminalProcessPid(process: PtyAdapter.PtyProcess | null): number | null { + const pid = process?.pid; + return pid !== undefined && Number.isInteger(pid) && pid > 0 ? pid : null; +} + function snapshot(session: TerminalSessionState): TerminalSessionSnapshot { return { threadId: session.threadId, @@ -463,10 +470,10 @@ function cleanupProcessHandles(session: TerminalSessionState): void { function enqueueProcessEvent( session: TerminalSessionState, - expectedPid: number, + expectedProcess: PtyAdapter.PtyProcess, event: PendingProcessEvent, ): boolean { - if (!session.process || session.status !== "running" || session.pid !== expectedPid) { + if (!session.process || session.status !== "running" || session.process !== expectedProcess) { return false; } @@ -2001,11 +2008,15 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const drainProcessEvents = Effect.fn("terminal.drainProcessEvents")(function* ( session: TerminalSessionState, - expectedPid: number, + expectedProcess: PtyAdapter.PtyProcess, ) { while (true) { const action: DrainProcessEventAction = yield* Effect.sync(() => { - if (session.pid !== expectedPid || !session.process || session.status !== "running") { + if ( + session.process !== expectedProcess || + !session.process || + session.status !== "running" + ) { session.pendingProcessEvents = []; session.pendingProcessEventIndex = 0; session.processEventDrainRunning = false; @@ -2050,7 +2061,6 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const process = session.process; cleanupProcessHandles(session); session.process = null; - session.pid = null; session.hasRunningSubprocess = false; session.childCommandLabel = null; session.status = "exited"; @@ -2122,7 +2132,6 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func yield* modifyManagerState((state) => { cleanupProcessHandles(session); session.process = null; - session.pid = null; session.hasRunningSubprocess = false; session.childCommandLabel = null; session.status = "exited"; @@ -2242,18 +2251,18 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func ptyProcess = spawnResult.process; startedShell = spawnResult.shellLabel; - const processPid = ptyProcess.pid; + const process = ptyProcess; const unsubscribeData = ptyProcess.onData((data) => { - if (!enqueueProcessEvent(session, processPid, { type: "output", data })) { + if (!enqueueProcessEvent(session, process, { type: "output", data })) { return; } - runFork(drainProcessEvents(session, processPid)); + runFork(drainProcessEvents(session, process)); }); const unsubscribeExit = ptyProcess.onExit((event) => { - if (!enqueueProcessEvent(session, processPid, { type: "exit", event })) { + if (!enqueueProcessEvent(session, process, { type: "exit", event })) { return; } - runFork(drainProcessEvents(session, processPid)); + runFork(drainProcessEvents(session, process)); }); let eventStamp: ReturnType = { @@ -2262,7 +2271,6 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func }; yield* modifyManagerState((state) => { session.process = ptyProcess; - session.pid = processPid; session.status = "running"; session.unsubscribeData = unsubscribeData; session.unsubscribeExit = unsubscribeExit; @@ -2295,7 +2303,6 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func yield* modifyManagerState((state) => { cleanupProcessHandles(session); session.status = "error"; - session.pid = null; session.process = null; session.hasRunningSubprocess = false; session.childCommandLabel = null; @@ -2550,7 +2557,9 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func cwd: input.cwd, worktreePath: input.worktreePath ?? null, status: "starting", - pid: null, + get pid() { + return terminalProcessPid(this.process); + }, history, pendingHistoryControlSequence: "", pendingProcessEvents: [], @@ -2973,7 +2982,9 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func cwd: input.cwd, worktreePath: input.worktreePath ?? null, status: "starting", - pid: null, + get pid() { + return terminalProcessPid(this.process); + }, history: new BoundedTerminalHistory(historyLineLimit, "", historyByteLimit), pendingHistoryControlSequence: "", pendingProcessEvents: [], From 6c680ef141a124002a968d17cd2dd15ad5b21868 Mon Sep 17 00:00:00 2001 From: SegFaultZero Date: Sun, 27 Sep 2026 14:05:03 +0800 Subject: [PATCH 2/2] fix(server): wait for Windows PTY readiness in the adapter --- apps/server/src/terminal/Manager.test.ts | 82 +---------- apps/server/src/terminal/Manager.ts | 43 ++---- .../src/terminal/NodePtyAdapter.test.ts | 139 +++++++++++++++++- apps/server/src/terminal/NodePtyAdapter.ts | 76 ++++++++++ 4 files changed, 225 insertions(+), 115 deletions(-) diff --git a/apps/server/src/terminal/Manager.test.ts b/apps/server/src/terminal/Manager.test.ts index 022577fe8132..f04161cfbd3c 100644 --- a/apps/server/src/terminal/Manager.test.ts +++ b/apps/server/src/terminal/Manager.test.ts @@ -11,8 +11,6 @@ import { ProviderInstanceId, ServerSettingsError, TerminalProviderInstanceNotFoundError, - TerminalSessionSnapshot, - TerminalSummary, } from "@t3tools/contracts"; import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as Data from "effect/Data"; @@ -30,7 +28,6 @@ import * as PlatformError from "effect/PlatformError"; import * as Path from "effect/Path"; import * as Ref from "effect/Ref"; import * as Schedule from "effect/Schedule"; -import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; @@ -45,9 +42,6 @@ import * as ServerSettings from "../serverSettings.ts"; import * as TerminalManager from "./Manager.ts"; import * as PtyAdapter from "./PtyAdapter.ts"; -const encodeTerminalSnapshot = Schema.encodeEffect(TerminalSessionSnapshot); -const encodeTerminalSummary = Schema.encodeEffect(TerminalSummary); - class WaitForConditionError extends Data.TaggedError("WaitForConditionError")<{ readonly message: string; }> {} @@ -56,7 +50,7 @@ class FakePtyProcess implements PtyAdapter.PtyProcess { readonly writes: string[] = []; readonly resizeCalls: Array<{ cols: number; rows: number }> = []; readonly killSignals: Array = []; - pid: number; + readonly pid: number; writeFailure: unknown | undefined; resizeFailure: unknown | undefined; private readonly dataListeners = new Set<(data: string) => void>(); @@ -118,7 +112,7 @@ class FakePtyAdapter { readonly processes: FakePtyProcess[] = []; readonly spawnFailures: Error[] = []; private readonly mode: "sync" | "async"; - nextPid = 9000; + private nextPid = 9000; constructor(mode: "sync" | "async" = "sync") { this.mode = mode; @@ -648,78 +642,6 @@ it.layer( }), ); - it.effect("keeps Windows terminals usable while ConPTY assigns the PID", () => - Effect.gen(function* () { - const adapter = new FakePtyAdapter(); - adapter.nextPid = 0; - const inspectedPids: number[] = []; - const { manager } = yield* createManager(5, { - ptyAdapter: adapter, - subprocessInspector: (pid) => - Effect.sync(() => { - inspectedPids.push(pid); - return { hasRunningSubprocess: true, childCommand: "node", processIds: [pid] }; - }), - }); - const readMetadata = Effect.gen(function* () { - let terminals: ReadonlyArray = []; - const unsubscribe = yield* manager.subscribeMetadata((event) => - Effect.sync(() => { - if (event.type === "snapshot") terminals = event.terminals; - }), - ); - unsubscribe(); - return terminals; - }); - const initial = yield* manager.open(openInput()); - assert.equal(initial.pid, null); - yield* encodeTerminalSnapshot(initial); - - const metadata = yield* readMetadata; - assert.equal(metadata[0]?.pid, null); - yield* Effect.forEach(metadata, (item) => encodeTerminalSummary(item)); - yield* manager.closeIdle({ threadId: "thread-1" }); - expect(inspectedPids).toEqual([]); - - const attached = yield* Ref.make>([]); - const unsubscribeAttach = yield* manager.attachStream(openInput(), (event) => - Ref.update(attached, (events) => [...events, event]), - ); - yield* Effect.addFinalizer(() => Effect.sync(unsubscribeAttach)); - const output = yield* Deferred.make(); - const exited = yield* Deferred.make(); - const unsubscribe = yield* manager.subscribe((event) => - event.type === "output" - ? Deferred.succeed(output, undefined).pipe(Effect.asVoid) - : event.type === "exited" - ? Deferred.succeed(exited, undefined).pipe(Effect.asVoid) - : Effect.void, - ); - yield* Effect.addFinalizer(() => Effect.sync(unsubscribe)); - - const process = adapter.processes[0]!; - process.pid = 12345; - process.emitData("Windows prompt> "); - yield* Deferred.await(output); - const ready = yield* manager.open(openInput()); - assert.equal(ready.pid, 12345); - assert.equal(ready.history, "Windows prompt> "); - assert.equal((yield* readMetadata)[0]?.pid, 12345); - yield* manager.closeIdle({ threadId: "thread-1" }); - expect(inspectedPids).toContain(12345); - expect(inspectedPids).not.toContain(0); - expect(yield* Ref.get(attached)).toContainEqual( - expect.objectContaining({ type: "output", data: "Windows prompt> " }), - ); - - process.emitExit({ exitCode: 0, signal: null }); - yield* Deferred.await(exited); - const finalMetadata = yield* readMetadata; - assert.equal(finalMetadata[0]?.status, "exited"); - assert.equal(finalMetadata[0]?.pid, null); - }), - ); - it.effect("forwards write and resize to active pty process", () => Effect.gen(function* () { const { manager, ptyAdapter } = yield* createManager(); diff --git a/apps/server/src/terminal/Manager.ts b/apps/server/src/terminal/Manager.ts index d9f0696a02c4..75d592c00c0d 100644 --- a/apps/server/src/terminal/Manager.ts +++ b/apps/server/src/terminal/Manager.ts @@ -276,7 +276,7 @@ interface TerminalSessionState { cwd: string; worktreePath: string | null; status: TerminalSessionStatus; - readonly pid: number | null; + pid: number | null; history: BoundedTerminalHistory; pendingHistoryControlSequence: string; pendingProcessEvents: Array; @@ -366,13 +366,6 @@ function terminalWireLabel(session: TerminalSessionState): string { return truncateTerminalWireLabel(getTerminalLabel(session.terminalId)); } -// ConPTY assigns the PID asynchronously. Read it from the process instead of -// freezing its startup sentinel into snapshots and subprocess checks. -function terminalProcessPid(process: PtyAdapter.PtyProcess | null): number | null { - const pid = process?.pid; - return pid !== undefined && Number.isInteger(pid) && pid > 0 ? pid : null; -} - function snapshot(session: TerminalSessionState): TerminalSessionSnapshot { return { threadId: session.threadId, @@ -470,10 +463,10 @@ function cleanupProcessHandles(session: TerminalSessionState): void { function enqueueProcessEvent( session: TerminalSessionState, - expectedProcess: PtyAdapter.PtyProcess, + expectedPid: number, event: PendingProcessEvent, ): boolean { - if (!session.process || session.status !== "running" || session.process !== expectedProcess) { + if (!session.process || session.status !== "running" || session.pid !== expectedPid) { return false; } @@ -2008,15 +2001,11 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const drainProcessEvents = Effect.fn("terminal.drainProcessEvents")(function* ( session: TerminalSessionState, - expectedProcess: PtyAdapter.PtyProcess, + expectedPid: number, ) { while (true) { const action: DrainProcessEventAction = yield* Effect.sync(() => { - if ( - session.process !== expectedProcess || - !session.process || - session.status !== "running" - ) { + if (session.pid !== expectedPid || !session.process || session.status !== "running") { session.pendingProcessEvents = []; session.pendingProcessEventIndex = 0; session.processEventDrainRunning = false; @@ -2061,6 +2050,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const process = session.process; cleanupProcessHandles(session); session.process = null; + session.pid = null; session.hasRunningSubprocess = false; session.childCommandLabel = null; session.status = "exited"; @@ -2132,6 +2122,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func yield* modifyManagerState((state) => { cleanupProcessHandles(session); session.process = null; + session.pid = null; session.hasRunningSubprocess = false; session.childCommandLabel = null; session.status = "exited"; @@ -2251,18 +2242,18 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func ptyProcess = spawnResult.process; startedShell = spawnResult.shellLabel; - const process = ptyProcess; + const processPid = ptyProcess.pid; const unsubscribeData = ptyProcess.onData((data) => { - if (!enqueueProcessEvent(session, process, { type: "output", data })) { + if (!enqueueProcessEvent(session, processPid, { type: "output", data })) { return; } - runFork(drainProcessEvents(session, process)); + runFork(drainProcessEvents(session, processPid)); }); const unsubscribeExit = ptyProcess.onExit((event) => { - if (!enqueueProcessEvent(session, process, { type: "exit", event })) { + if (!enqueueProcessEvent(session, processPid, { type: "exit", event })) { return; } - runFork(drainProcessEvents(session, process)); + runFork(drainProcessEvents(session, processPid)); }); let eventStamp: ReturnType = { @@ -2271,6 +2262,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func }; yield* modifyManagerState((state) => { session.process = ptyProcess; + session.pid = processPid; session.status = "running"; session.unsubscribeData = unsubscribeData; session.unsubscribeExit = unsubscribeExit; @@ -2303,6 +2295,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func yield* modifyManagerState((state) => { cleanupProcessHandles(session); session.status = "error"; + session.pid = null; session.process = null; session.hasRunningSubprocess = false; session.childCommandLabel = null; @@ -2557,9 +2550,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func cwd: input.cwd, worktreePath: input.worktreePath ?? null, status: "starting", - get pid() { - return terminalProcessPid(this.process); - }, + pid: null, history, pendingHistoryControlSequence: "", pendingProcessEvents: [], @@ -2982,9 +2973,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func cwd: input.cwd, worktreePath: input.worktreePath ?? null, status: "starting", - get pid() { - return terminalProcessPid(this.process); - }, + pid: null, history: new BoundedTerminalHistory(historyLineLimit, "", historyByteLimit), pendingHistoryControlSequence: "", pendingProcessEvents: [], 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); }), });