Files
kuber/tests/server/exec-websocket.test.ts

354 lines
9.3 KiB
TypeScript

import { describe, expect, test } from "bun:test";
import {
authorizeExecConnection,
execProblem,
WireExecSession,
type ExecConnection,
type ExecWebSocketLink,
} from "../../server/app";
import { hashToken, MemoryAuthStore } from "../../server/auth";
import {
createExecService,
EXEC_MANAGED_BY_LABEL,
EXEC_MANAGED_BY_VALUE,
EXEC_WORKSPACE_UID_LABEL,
type ExecChunk,
type ExecDeployment,
type ExecPod,
type KubernetesExecBackend,
type KubernetesExecExit,
type KubernetesExecProcess,
type KubernetesExecRequest,
} from "../../server/exec-service";
import { MemoryWorkspaceStore } from "../../server/workspace-store";
const now = () => Date.parse("2026-09-02T00:00:00.000Z");
const workspace = { project: "shop", uid: "workspace-1" };
async function* chunks(...values: ExecChunk[]): AsyncGenerator<ExecChunk> {
for (const value of values) yield value;
}
function deferred<T>() {
let resolve!: (value: T) => void;
let reject!: (reason?: unknown) => void;
const promise = new Promise<T>((resolvePromise, rejectPromise) => {
resolve = resolvePromise;
reject = rejectPromise;
});
return { promise, resolve, reject };
}
class FakeProcess implements KubernetesExecProcess {
stdout: AsyncIterable<ExecChunk> = chunks();
stderr: AsyncIterable<ExecChunk> = chunks();
status: Promise<KubernetesExecExit> = Promise.resolve({ exitCode: 0 });
stdin: Uint8Array[] = [];
resizes: Array<[number, number]> = [];
stdinClosed = false;
closeCalls = 0;
writeStdin(data: Uint8Array) {
this.stdin.push(data);
}
closeStdin() {
this.stdinClosed = true;
}
resize(columns: number, rows: number) {
this.resizes.push([columns, rows]);
}
wait() {
return this.status;
}
close() {
this.closeCalls += 1;
}
}
class FakeBackend implements KubernetesExecBackend {
deployment: ExecDeployment | undefined = {
name: "api",
uid: "deployment-1",
labels: {
[EXEC_MANAGED_BY_LABEL]: EXEC_MANAGED_BY_VALUE,
[EXEC_WORKSPACE_UID_LABEL]: workspace.uid,
},
selector: { app: "api" },
containers: ["api"],
};
pods: ExecPod[] = [
{
name: "api-0",
uid: "pod-0",
deploymentUid: "deployment-1",
phase: "Running",
containers: [{ name: "api", running: true, ready: true }],
},
];
process = new FakeProcess();
requests: KubernetesExecRequest[] = [];
async getDeployment() {
return this.deployment;
}
async listPods() {
return this.pods;
}
async exec(request: KubernetesExecRequest) {
this.requests.push(request);
return this.process;
}
}
class FakeSocket implements ExecWebSocketLink {
sent: string[] = [];
closes: Array<[number | undefined, string | undefined]> = [];
sendText(data: string) {
this.sent.push(data);
}
close(code?: number, reason?: string) {
this.closes.push([code, reason]);
}
}
function workspaceRecord(uid = "workspace-1") {
return {
apiVersion: "kuber.astrxl.dev/v2",
kind: "Workspace",
metadata: {
name: "shop",
uid,
resourceVersion: "1",
creationTimestamp: "2026-09-02T00:00:00.000Z",
},
spec: {
source: { uri: "oci://example/demo", digest: "sha256:abc" },
},
status: { latestRevision: 1 },
} as const;
}
function connection(role: "viewer" | "operator" | "admin" = "operator") {
return {
identity: {
user: {
username: role,
passwordHash: "hash",
roles: [role],
authVersion: 1,
},
session: {
tokenHash: hashToken("token"),
username: role,
authVersion: 1,
expiresAt: "2030-01-01T00:00:00.000Z",
},
},
workspace: workspaceRecord(),
} as ExecConnection;
}
async function authenticatedStore(
role: "viewer" | "operator" | "admin" = "operator",
) {
const store = new MemoryAuthStore();
await store.putUser({ username: role, passwordHash: "hash", roles: [role] });
await store.putSession({
tokenHash: hashToken("token"),
username: role,
authVersion: 1,
expiresAt: "2030-01-01T00:00:00.000Z",
});
return store;
}
function request(path = "/api/v2/workspaces/shop/exec") {
return new Request(`https://kuber.astrxl.dev${path}`, {
headers: { authorization: "Bearer token" },
});
}
describe("authorizeExecConnection", () => {
test("resolves the workspace UID server-side for authorized operators", async () => {
const store = await authenticatedStore("operator");
const workspaceStore = new MemoryWorkspaceStore({
uid: () => "workspace-1",
});
await workspaceStore.create({
id: "shop",
source: { uri: "oci://example/demo", digest: "sha256:abc" },
});
const result = await authorizeExecConnection(
{ store, workspaceStore, now },
request(),
"shop",
);
expect(result.identity.user.roles).toContain("operator");
expect(result.workspace.metadata.uid).toBe("workspace-1");
});
test("rejects unauthenticated upgrades with a bearer challenge", async () => {
const result = await authorizeExecConnection(
{ store: new MemoryAuthStore(), now },
new Request("https://kuber.astrxl.dev/api/v2/workspaces/shop/exec"),
"shop",
).then(
() => null,
(error) => error,
);
expect((result as { status: number }).status).toBe(401);
});
test("rejects operators without the exec capability", async () => {
const store = await authenticatedStore("viewer");
const error = await authorizeExecConnection(
{ store, now },
request(),
"shop",
).then(
() => null,
(error) => error,
);
expect((error as { status: number }).status).toBe(403);
});
test("does not resolve unauthenticated unknown workspaces", async () => {
const store = await authenticatedStore("operator");
const workspaceStore = new MemoryWorkspaceStore();
const error = await authorizeExecConnection(
{ store, workspaceStore, now },
request(),
"missing",
).then(
() => null,
(error) => error,
);
expect((error as { status: number }).status).toBe(404);
});
});
describe("WireExecSession", () => {
test("opens on start and maps stdin/resize/output/exit wire frames", async () => {
const backend = new FakeBackend();
backend.process.stdout = chunks("out");
backend.process.stderr = chunks(new TextEncoder().encode("err"));
backend.process.status = Promise.resolve({ exitCode: 3 });
const socket = new FakeSocket();
const session = new WireExecSession(
socket,
createExecService(backend),
connection(),
);
await session.receive(
JSON.stringify({
type: "start",
version: 1,
deployment: "api",
command: ["sh", "-c", "echo hi"],
tty: true,
}),
);
await session.receive(
JSON.stringify({
type: "stdin",
data: Buffer.from("hello").toString("base64"),
encoding: "base64",
}),
);
await session.receive(
JSON.stringify({ type: "resize", columns: 100, rows: 40 }),
);
await Bun.sleep(1);
expect(new TextDecoder().decode(backend.process.stdin[0])).toBe("hello");
expect(backend.process.resizes).toEqual([[100, 40]]);
const frames = socket.sent.map((text) => JSON.parse(text));
expect(frames).toContainEqual({
type: "stdout",
data: Buffer.from("out").toString("base64"),
encoding: "base64",
});
expect(frames).toContainEqual({ type: "exit", exitCode: 3 });
expect(socket.closes[0]?.[0]).toBe(1000);
});
test("rejects an unknown first frame", async () => {
const socket = new FakeSocket();
const session = new WireExecSession(
socket,
createExecService(new FakeBackend()),
connection(),
);
await session.receive(
JSON.stringify({ type: "resize", columns: 1, rows: 1 }),
);
expect(JSON.parse(socket.sent[0]!)).toMatchObject({
type: "error",
code: "EXEC_INVALID",
});
});
test("rejects oversize frames", async () => {
const socket = new FakeSocket();
const session = new WireExecSession(
socket,
createExecService(new FakeBackend()),
connection(),
{ maxFrameBytes: 4 },
);
await session.receive(JSON.stringify({ type: "close" }));
expect(JSON.parse(socket.sent[0]!)).toMatchObject({ type: "error" });
expect(socket.closes[0]?.[0]).toBe(1009);
});
test("aborts the process when the socket closes", async () => {
const backend = new FakeBackend();
backend.process.stdout = (async function* () {
yield "data";
})();
const socket = new FakeSocket();
const session = new WireExecSession(
socket,
createExecService(backend),
connection(),
);
await session.receive(
JSON.stringify({
type: "start",
version: 1,
deployment: "api",
command: ["sh"],
tty: false,
}),
);
await Bun.sleep(1);
session.close();
});
});
describe("execProblem", () => {
test("formats HTTP errors as problem+json without leaking internals", async () => {
const response = execProblem(
new (class extends Error {
readonly status = 401;
readonly title = "Unauthorized";
readonly code = "UNAUTHORIZED";
})("boom"),
);
expect(response.status).toBe(500);
expect((await response.json()) as { code: string }).toMatchObject({
code: "INTERNAL_ERROR",
});
});
});