Files

228 lines
7.3 KiB
TypeScript

import { KUBER_API_BASE_URL } from "../const";
import { readSession, type KuberSession } from "./session";
export type ExecStartFrame = {
type: "start";
version: 1;
deployment: string;
command: string[];
tty: boolean;
container?: string;
columns?: number;
rows?: number;
};
export type ExecClientWireFrame =
| ExecStartFrame
| { type: "stdin"; data: string; encoding: "base64"; eof?: boolean }
| { type: "resize"; columns: number; rows: number }
| { type: "close" };
export type ExecServerWireFrame =
| { type: "stdout" | "stderr"; data: string; encoding: "base64" }
| { type: "exit"; exitCode: number; reason?: string; message?: string }
| { type: "error"; code: string; message: string };
export type ExecOutputFrame =
| { type: "stdout"; data: Uint8Array }
| { type: "stderr"; data: Uint8Array }
| Exclude<ExecServerWireFrame, { type: "stdout" | "stderr" }>;
export interface ExecWebSocket {
binaryType: "arraybuffer" | "blob";
readyState: number;
send(data: string): void;
close(code?: number, reason?: string): void;
addEventListener(type: string, listener: (event: any) => void): void;
removeEventListener(type: string, listener: (event: any) => void): void;
}
export type ExecWebSocketFactory = (
url: string,
headers: Readonly<Record<string, string>>,
) => ExecWebSocket;
export type OpenExecOptions = {
baseUrl?: string;
session?: KuberSession;
socketFactory?: ExecWebSocketFactory;
};
export interface ExecApiSession extends AsyncIterable<ExecOutputFrame> {
sendStdin(data: Uint8Array, eof?: boolean): void;
resize(columns: number, rows: number): void;
close(): void;
}
function websocketUrl(baseUrl: string, path: string): string {
const url = new URL(`${baseUrl}${path}`);
url.protocol = url.protocol === "https:" ? "wss:" : "ws:";
return url.toString();
}
const defaultSocketFactory: ExecWebSocketFactory = (url, headers) =>
new WebSocket(url, {
headers,
} as unknown as string[]) as unknown as ExecWebSocket;
function parseFrame(value: unknown): ExecOutputFrame {
if (typeof value !== "string")
throw new Error("Exec frame must be JSON text");
const frame: unknown = JSON.parse(value);
if (!frame || typeof frame !== "object" || !("type" in frame))
throw new Error("Invalid exec frame");
const wire = frame as ExecServerWireFrame;
if (wire.type === "stdout" || wire.type === "stderr") {
if (wire.encoding !== "base64" || typeof wire.data !== "string")
throw new Error("Invalid exec output frame");
return {
type: wire.type,
data: Uint8Array.from(Buffer.from(wire.data, "base64")),
};
}
if (wire.type === "exit" && Number.isSafeInteger(wire.exitCode)) return wire;
if (wire.type === "error" && typeof wire.message === "string") return wire;
throw new Error("Invalid exec frame");
}
export async function openExecSession(
project: string,
start: Omit<ExecStartFrame, "type" | "version">,
signal: AbortSignal,
options: OpenExecOptions = {},
): Promise<ExecApiSession> {
if (signal.aborted)
throw signal.reason ?? new DOMException("Aborted", "AbortError");
const session = options.session ?? (await readSession());
if (!session) throw new Error("Not logged in. Run kuber login first.");
const path = `/workspaces/${encodeURIComponent(project)}/exec`;
const socket = (options.socketFactory ?? defaultSocketFactory)(
websocketUrl(options.baseUrl ?? KUBER_API_BASE_URL, path),
{
authorization: `Bearer ${session.token}`,
},
);
socket.binaryType = "arraybuffer";
await new Promise<void>((resolve, reject) => {
const onOpen = () => {
cleanup();
socket.send(JSON.stringify({ type: "start", version: 1, ...start }));
resolve();
};
const onError = (event: { error?: unknown }) => {
cleanup();
reject(event.error ?? new Error("Exec WebSocket connection failed"));
};
const onClose = () => {
cleanup();
reject(new Error("Exec WebSocket closed before opening"));
};
const onAbort = () => {
cleanup();
socket.close(1000, "aborted");
reject(signal.reason ?? new DOMException("Aborted", "AbortError"));
};
const cleanup = () => {
socket.removeEventListener("open", onOpen);
socket.removeEventListener("error", onError);
socket.removeEventListener("close", onClose);
signal.removeEventListener("abort", onAbort);
};
socket.addEventListener("open", onOpen);
socket.addEventListener("error", onError);
socket.addEventListener("close", onClose);
signal.addEventListener("abort", onAbort, { once: true });
});
const frames: ExecOutputFrame[] = [];
const readers: Array<{
resolve: (result: IteratorResult<ExecOutputFrame>) => void;
reject: (error: unknown) => void;
}> = [];
let closed = false;
let failure: unknown;
let live = false;
const pending: ExecClientWireFrame[] = [];
const finish = (error?: unknown) => {
if (closed) return;
closed = true;
failure = error;
socket.removeEventListener("message", onMessage);
socket.removeEventListener("error", onError);
socket.removeEventListener("close", onClose);
signal.removeEventListener("abort", onAbort);
for (const reader of readers.splice(0)) {
if (error) reader.reject(error);
else reader.resolve({ done: true, value: undefined });
}
};
const onMessage = (event: { data: unknown }) => {
try {
if (!live) {
live = true;
for (const frame of pending.splice(0))
socket.send(JSON.stringify(frame));
}
const frame = parseFrame(event.data);
const reader = readers.shift();
if (reader) reader.resolve({ done: false, value: frame });
else frames.push(frame);
} catch (error) {
socket.close(1002, "invalid frame");
finish(error);
}
};
const onError = (event: { error?: unknown }) =>
finish(event.error ?? new Error("Exec WebSocket failed"));
const onClose = () => finish();
const onAbort = () => {
socket.close(1000, "aborted");
finish();
};
socket.addEventListener("message", onMessage);
socket.addEventListener("error", onError);
socket.addEventListener("close", onClose);
signal.addEventListener("abort", onAbort, { once: true });
const send = (frame: ExecClientWireFrame) => {
if (closed) throw failure ?? new Error("Exec session is closed");
if (!live) {
pending.push(frame);
return;
}
socket.send(JSON.stringify(frame));
};
return {
sendStdin(data, eof) {
send({
type: "stdin",
data: Buffer.from(data).toString("base64"),
encoding: "base64",
...(eof ? { eof: true } : {}),
});
},
resize(columns, rows) {
send({ type: "resize", columns, rows });
},
close() {
if (!closed && socket.readyState === 1) send({ type: "close" });
socket.close(1000, "client closed");
finish();
},
[Symbol.asyncIterator]() {
return {
next(): Promise<IteratorResult<ExecOutputFrame>> {
const frame = frames.shift();
if (frame) return Promise.resolve({ done: false, value: frame });
if (failure) return Promise.reject(failure);
if (closed) return Promise.resolve({ done: true, value: undefined });
return new Promise((resolve, reject) =>
readers.push({ resolve, reject }),
);
},
};
},
};
}