Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
139 changes: 131 additions & 8 deletions apps/server/src/terminal/NodePtyAdapter.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>();
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");

Expand All @@ -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* () {
Expand Down
76 changes: 76 additions & 0 deletions apps/server/src/terminal/NodePtyAdapter.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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<void, PtyAdapter.PtySpawnError>((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;
Expand Down Expand Up @@ -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);
}),
});
Expand Down