368 lines
11 KiB
TypeScript
368 lines
11 KiB
TypeScript
import { afterEach, describe, expect, mock, test } from "bun:test";
|
|
import {
|
|
DEFAULT_API_TIMEOUT_MS,
|
|
KuberApiError,
|
|
apiRequest,
|
|
apiStreamNdjson,
|
|
apiUpload,
|
|
} from "../../lib/api";
|
|
import type { KuberSession } from "../../lib/session";
|
|
import { createVersionObserver } from "../../lib/version-update";
|
|
|
|
const originalFetch = globalThis.fetch;
|
|
const session: KuberSession = {
|
|
token: "secret-token",
|
|
expiresAt: "2099-01-01T00:00:00.000Z",
|
|
user: { username: "operator", roles: ["operator"] },
|
|
};
|
|
|
|
afterEach(() => {
|
|
globalThis.fetch = originalFetch;
|
|
});
|
|
|
|
test("observes the response header on JSON, errors, uploads and streaming without delaying payloads", async () => {
|
|
const seen: Array<string | null> = [];
|
|
const onServerVersion = (version: string | null) => {
|
|
seen.push(version);
|
|
};
|
|
const options = {
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
onServerVersion,
|
|
};
|
|
const encoder = new TextEncoder();
|
|
globalThis.fetch = mock(async (input) => {
|
|
const path = new URL(String(input)).pathname;
|
|
const headers = { "x-kuber-version": "2.7.0" };
|
|
if (path === "/error")
|
|
return Response.json({ detail: "failure" }, { status: 503, headers });
|
|
if (path === "/upload") return new Response(null, { status: 204, headers });
|
|
if (path === "/stream")
|
|
return new Response(
|
|
new ReadableStream({
|
|
start(controller) {
|
|
controller.enqueue(encoder.encode('{"n":1}\n'));
|
|
},
|
|
}),
|
|
{ headers },
|
|
);
|
|
return Response.json({ ok: true }, { headers });
|
|
}) as unknown as typeof fetch;
|
|
|
|
expect(await apiRequest<{ ok: boolean }>("/json", {}, options)).toEqual({
|
|
ok: true,
|
|
});
|
|
await expect(apiRequest("/error", {}, options)).rejects.toBeInstanceOf(
|
|
KuberApiError,
|
|
);
|
|
await apiUpload("/upload", new Uint8Array([1]), { ...options, offset: 0 });
|
|
|
|
let complete!: (success: boolean) => void;
|
|
const pending = new Promise<boolean>((resolve) => {
|
|
complete = resolve;
|
|
});
|
|
const lines: string[] = [];
|
|
const observer = createVersionObserver({
|
|
currentVersion: "2.6.1",
|
|
runner: () => pending,
|
|
report: (line) => lines.push(line),
|
|
});
|
|
const records = apiStreamNdjson<{ n: number }>(
|
|
"/stream",
|
|
{},
|
|
{
|
|
...options,
|
|
onServerVersion: (version) => {
|
|
onServerVersion(version);
|
|
observer(version);
|
|
},
|
|
},
|
|
);
|
|
expect((await records.next()).value).toEqual({ n: 1 });
|
|
expect(seen).toEqual(["2.7.0", "2.7.0", "2.7.0", "2.7.0"]);
|
|
expect(lines).toEqual([]);
|
|
await records.return();
|
|
complete(false);
|
|
await new Promise<void>((resolve) => setTimeout(resolve, 0));
|
|
expect(lines).toEqual(["\x1b[90m+ New version available: 2.7.0\x1b[0m\n"]);
|
|
});
|
|
|
|
describe("authenticated API transport", () => {
|
|
test("sends authentication and serializes JSON", async () => {
|
|
const fetchMock = mock(
|
|
async (input: Parameters<typeof fetch>[0], init?: RequestInit) => {
|
|
expect(input).toBe("https://api.test/workspaces");
|
|
expect(new Headers(init?.headers).get("authorization")).toBe(
|
|
"Bearer secret-token",
|
|
);
|
|
expect(new Headers(init?.headers).get("content-type")).toBe(
|
|
"application/json",
|
|
);
|
|
expect(new Headers(init?.headers).has("user-agent")).toBe(false);
|
|
expect(init?.body).toBe('{"name":"demo"}');
|
|
return Response.json({ id: "workspace-1" });
|
|
},
|
|
);
|
|
globalThis.fetch = fetchMock as unknown as typeof fetch;
|
|
|
|
await expect(
|
|
apiRequest<{ id: string }>(
|
|
"/workspaces",
|
|
{ method: "POST", json: { name: "demo" } },
|
|
{ baseUrl: "https://api.test", session },
|
|
),
|
|
).resolves.toEqual({ id: "workspace-1" });
|
|
expect(fetchMock).toHaveBeenCalledTimes(1);
|
|
});
|
|
|
|
test("preserves unauthenticated callers and empty responses", async () => {
|
|
globalThis.fetch = mock(async (_input, init) => {
|
|
expect(new Headers(init?.headers).has("authorization")).toBe(false);
|
|
return new Response(null, { status: 204 });
|
|
}) as unknown as typeof fetch;
|
|
|
|
await expect(
|
|
apiRequest<void>(
|
|
"/logout",
|
|
{ method: "POST" },
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
),
|
|
).resolves.toBeUndefined();
|
|
});
|
|
|
|
test("preserves a caller-provided User-Agent without adding one by default", async () => {
|
|
const sent: Array<string | null> = [];
|
|
globalThis.fetch = mock(async (_input, init) => {
|
|
sent.push(new Headers(init?.headers).get("user-agent"));
|
|
return new Response(null, { status: 204 });
|
|
}) as unknown as typeof fetch;
|
|
|
|
const options = { authenticated: false, baseUrl: "https://api.test" };
|
|
await apiRequest("/plain", {}, options);
|
|
await apiRequest("/custom", { headers: { "user-agent": "my-app" } }, options);
|
|
expect(sent).toEqual([null, "my-app"]);
|
|
});
|
|
|
|
test("treats an empty successful JSON response as no content", async () => {
|
|
globalThis.fetch = mock(
|
|
async () => new Response("", { status: 200 }),
|
|
) as unknown as typeof fetch;
|
|
|
|
await expect(
|
|
apiRequest<void>(
|
|
"/empty",
|
|
{},
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
),
|
|
).resolves.toBeUndefined();
|
|
});
|
|
|
|
test("exposes typed problem details and correlation headers", async () => {
|
|
globalThis.fetch = mock(async () =>
|
|
Response.json(
|
|
{
|
|
title: "Conflict",
|
|
status: 409,
|
|
detail: "Workspace already exists",
|
|
code: "WORKSPACE_EXISTS",
|
|
operationId: "op-7",
|
|
},
|
|
{
|
|
status: 409,
|
|
headers: { "x-request-id": "req-4" },
|
|
},
|
|
),
|
|
) as unknown as typeof fetch;
|
|
|
|
try {
|
|
await apiRequest(
|
|
"/workspaces",
|
|
{},
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
);
|
|
throw new Error("expected request to fail");
|
|
} catch (error) {
|
|
expect(error).toBeInstanceOf(KuberApiError);
|
|
expect(error).toMatchObject({
|
|
message: "Workspace already exists",
|
|
status: 409,
|
|
code: "WORKSPACE_EXISTS",
|
|
requestId: "req-4",
|
|
operationId: "op-7",
|
|
});
|
|
expect((error as KuberApiError).problem.title).toBe("Conflict");
|
|
}
|
|
});
|
|
|
|
test("provides a stable fallback problem for non-JSON errors", async () => {
|
|
globalThis.fetch = mock(
|
|
async () => new Response("upstream failure", { status: 502 }),
|
|
) as unknown as typeof fetch;
|
|
|
|
await expect(
|
|
apiRequest(
|
|
"/status",
|
|
{},
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
),
|
|
).rejects.toMatchObject({ status: 502, code: "HTTP_502" });
|
|
});
|
|
|
|
test("applies a finite timeout by default", async () => {
|
|
globalThis.fetch = mock(async (_input, init) => {
|
|
expect(DEFAULT_API_TIMEOUT_MS).toBeGreaterThan(0);
|
|
return await new Promise<Response>((_resolve, reject) => {
|
|
init?.signal?.addEventListener(
|
|
"abort",
|
|
() => reject(init.signal?.reason),
|
|
{
|
|
once: true,
|
|
},
|
|
);
|
|
});
|
|
}) as unknown as typeof fetch;
|
|
|
|
await expect(
|
|
apiRequest(
|
|
"/slow",
|
|
{},
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
timeoutMs: 5,
|
|
},
|
|
),
|
|
).rejects.toMatchObject({ name: "TimeoutError" });
|
|
});
|
|
|
|
test("propagates caller aborts to fetch", async () => {
|
|
const controller = new AbortController();
|
|
const reason = new DOMException("cancelled", "AbortError");
|
|
globalThis.fetch = mock(async (_input, init) => {
|
|
if (init?.signal?.aborted) throw init.signal.reason;
|
|
return await new Promise<Response>((_resolve, reject) => {
|
|
init?.signal?.addEventListener(
|
|
"abort",
|
|
() => reject(init.signal?.reason),
|
|
{
|
|
once: true,
|
|
},
|
|
);
|
|
});
|
|
}) as unknown as typeof fetch;
|
|
|
|
const request = apiRequest(
|
|
"/operations/1",
|
|
{ signal: controller.signal },
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
);
|
|
controller.abort(reason);
|
|
await expect(request).rejects.toBe(reason);
|
|
});
|
|
|
|
test("uploads binary chunks without assigning a JSON content type", async () => {
|
|
const bytes = new Uint8Array([1, 2, 3]);
|
|
globalThis.fetch = mock(async (_input, init) => {
|
|
const headers = new Headers(init?.headers);
|
|
expect(init?.method).toBe("PATCH");
|
|
expect(headers.get("content-type")).toBe("application/octet-stream");
|
|
expect(headers.get("upload-offset")).toBe("12");
|
|
expect(init?.body).toBe(bytes);
|
|
return Response.json({
|
|
uploadId: "upload-1",
|
|
offset: 15,
|
|
complete: false,
|
|
});
|
|
}) as unknown as typeof fetch;
|
|
|
|
await expect(
|
|
apiUpload<{ offset: number }>("/uploads/upload-1", bytes, {
|
|
offset: 12,
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
}),
|
|
).resolves.toMatchObject({ offset: 15 });
|
|
});
|
|
|
|
test("consumes split NDJSON records incrementally", async () => {
|
|
let source: ReadableStreamDefaultController<Uint8Array> | undefined;
|
|
const stream = new ReadableStream<Uint8Array>({
|
|
start(controller) {
|
|
source = controller;
|
|
controller.enqueue(new TextEncoder().encode('{"sequence":1}\n{"seq'));
|
|
},
|
|
});
|
|
globalThis.fetch = mock(async (_input, init) => {
|
|
expect(new Headers(init?.headers).get("accept")).toBe(
|
|
"application/x-ndjson",
|
|
);
|
|
return new Response(stream);
|
|
}) as unknown as typeof fetch;
|
|
|
|
const records = apiStreamNdjson<{ sequence: number }>(
|
|
"/operations/1/events",
|
|
{},
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
);
|
|
const first = await records.next();
|
|
expect(first.value).toEqual({ sequence: 1 });
|
|
source?.enqueue(new TextEncoder().encode('uence":2}\r\n\n'));
|
|
source?.close();
|
|
expect((await records.next()).value).toEqual({ sequence: 2 });
|
|
expect((await records.next()).done).toBe(true);
|
|
});
|
|
|
|
test("rejects malformed NDJSON records", async () => {
|
|
globalThis.fetch = mock(
|
|
async () => new Response('{"ok":true}\nnot-json\n'),
|
|
) as unknown as typeof fetch;
|
|
|
|
const consume = async () => {
|
|
for await (const _record of apiStreamNdjson(
|
|
"/events",
|
|
{},
|
|
{
|
|
authenticated: false,
|
|
baseUrl: "https://api.test",
|
|
},
|
|
)) {
|
|
// Consume the complete stream.
|
|
}
|
|
};
|
|
await expect(consume()).rejects.toBeInstanceOf(SyntaxError);
|
|
});
|
|
|
|
test("parses many records packed into one NDJSON chunk", async () => {
|
|
const expected = Array.from({ length: 2_000 }, (_, sequence) => ({ sequence }));
|
|
globalThis.fetch = mock(async () =>
|
|
new Response(`${expected.map((record) => JSON.stringify(record)).join("\n")}\n`),
|
|
) as unknown as typeof fetch;
|
|
|
|
const actual: Array<{ sequence: number }> = [];
|
|
for await (const record of apiStreamNdjson<{ sequence: number }>(
|
|
"/events",
|
|
{},
|
|
{ authenticated: false, baseUrl: "https://api.test" },
|
|
)) actual.push(record);
|
|
expect(actual).toEqual(expected);
|
|
});
|
|
});
|