Files
pikasTech-unidesk/scripts/src/ssh-output-flush.test.ts
T
2026-07-20 05:24:14 +02:00

129 lines
4.4 KiB
TypeScript

import { mkdtempSync, rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { describe, expect, test } from "bun:test";
import {
flushWebSocketSendQueue,
sshProxySendDrainTimeoutMs,
waitForWebSocketSendDrain,
} from "../../src/components/frontend/src/ssh-proxy-flush";
import { flushWritableStreams } from "./ssh-output-flush";
import { sshBrokerSource } from "./ssh-runtime";
describe("ssh output completion", () => {
test("waits for delayed writable callbacks", async () => {
let flushed = false;
const stream = {
write(_chunk: string, callback: () => void) {
setTimeout(() => {
flushed = true;
callback();
}, 5);
return false;
},
} as unknown as NodeJS.WritableStream;
await flushWritableStreams([stream]);
expect(flushed).toBe(true);
});
test("waits for websocket buffered data before close", async () => {
let amount = 4096;
setTimeout(() => {
amount = 0;
}, 5);
await waitForWebSocketSendDrain(() => amount, { timeoutMs: 100, pollMs: 1 });
expect(amount).toBe(0);
});
test("bounds websocket drain wait", async () => {
const startedAt = Date.now();
await waitForWebSocketSendDrain(() => 1, { timeoutMs: 10, pollMs: 1 });
expect(Date.now() - startedAt).toBeGreaterThanOrEqual(9);
});
test("supports waiting without a transfer-size-dependent close deadline", async () => {
let amount = 1;
setTimeout(() => {
amount = 0;
}, 15);
await waitForWebSocketSendDrain(() => amount, { timeoutMs: null, pollMs: 1 });
expect(amount).toBe(0);
});
test("allows sustained high-pressure websocket drain before the hard timeout", async () => {
expect(sshProxySendDrainTimeoutMs).toBe(60_000);
let amount = 64 * 1024;
setTimeout(() => {
amount = 0;
}, 1100);
await waitForWebSocketSendDrain(() => amount, { pollMs: 5 });
expect(amount).toBe(0);
});
test("retains a dropped websocket frame and resumes in order after drain", () => {
const queue = ["one", "two", "exit"];
const firstStatuses = [3, -1];
const firstSent: string[] = [];
expect(flushWebSocketSendQueue(queue, (message) => {
firstSent.push(String(message));
return firstStatuses.shift()!;
})).toEqual({ backpressured: true, dropped: false });
expect(firstSent).toEqual(["one", "two"]);
expect(queue).toEqual(["exit"]);
expect(flushWebSocketSendQueue(queue, () => 0)).toEqual({ backpressured: true, dropped: true });
expect(queue).toEqual(["exit"]);
const resumed: string[] = [];
expect(flushWebSocketSendQueue(queue, (message) => {
resumed.push(String(message));
return 4;
})).toEqual({ backpressured: false, dropped: false });
expect(resumed).toEqual(["exit"]);
expect(queue).toEqual([]);
});
test("broker preserves the final stdout chunk before exit", () => {
const dir = mkdtempSync(join(tmpdir(), "unidesk-ssh-broker-test-"));
const scriptPath = join(dir, "broker.mjs");
const bytes = 256 * 1024;
const mockWebSocket = `
class MockWebSocket {
static OPEN = 1;
readyState = MockWebSocket.OPEN;
listeners = new Map();
constructor() { queueMicrotask(() => this.emit("open", {})); }
addEventListener(type, listener) { this.listeners.set(type, listener); }
emit(type, event) { this.listeners.get(type)?.(event); }
send(text) {
const message = JSON.parse(text);
if (message.type !== "ssh.open") return;
queueMicrotask(() => {
this.emit("message", { data: JSON.stringify({ type: "ssh.opened" }) });
this.emit("message", { data: JSON.stringify({ type: "ssh.data", stream: "stdout", data: Buffer.alloc(${bytes}, "x").toString("base64") }) });
this.emit("message", { data: JSON.stringify({ type: "ssh.exit", exitCode: 0 }) });
});
}
close() { this.readyState = 3; this.emit("close", {}); }
}
globalThis.WebSocket = MockWebSocket;
`;
writeFileSync(scriptPath, `${mockWebSocket}\n${sshBrokerSource()}`, "utf8");
try {
const result = Bun.spawnSync({
cmd: [process.execPath, scriptPath, JSON.stringify({ providerId: "test", runtimeTimeoutMs: 1000 })],
stdout: "pipe",
stderr: "pipe",
});
expect(result.exitCode).toBe(0);
expect(result.stdout.byteLength).toBe(bytes);
expect(result.stderr.toString()).toBe("");
} finally {
rmSync(dir, { recursive: true, force: true });
}
});
});