Skip to content
34 changes: 34 additions & 0 deletions apps/server/src/terminal/Manager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand All @@ -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;
Expand All @@ -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({
Expand Down Expand Up @@ -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<void>();
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, {
Expand Down
31 changes: 17 additions & 14 deletions apps/server/src/terminal/Manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof advanceEventSequence> = {
updatedAt: session.updatedAt,
Expand All @@ -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;
});
Expand All @@ -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));
}),
),
),
Expand Down
228 changes: 219 additions & 9 deletions apps/server/src/terminal/NodePtyAdapter.test.ts
Original file line number Diff line number Diff line change
@@ -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<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 +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* () {
Expand Down Expand Up @@ -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)),
);
}
Loading
Loading