import { afterEach, describe, expect, test } from "bun:test"; import { HttpMethod, RequestContext, } from "@kubernetes/client-node/dist/gen/http/http.js"; import http from "node:http"; import { createKubernetesHttpLibrary, delay, parseRetryAfter, } from "../../lib/k8s-http"; const servers: http.Server[] = []; async function serve( handler: http.RequestListener, ): Promise<{ server: http.Server; url: string }> { const server = http.createServer(handler); servers.push(server); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing port"); return { server, url: `http://127.0.0.1:${address.port}` }; } afterEach(async () => { await Promise.all( servers.splice(0).map( (server) => new Promise((resolve, reject) => { server.close((error) => (error ? reject(error) : resolve())); }), ), ); }); describe("Kubernetes HTTP transport", () => { test("parses Retry-After seconds and HTTP dates", () => { expect(parseRetryAfter("2", 0)).toBe(2000); expect(parseRetryAfter("Thu, 01 Jan 1970 00:00:03 GMT", 1000)).toBe(2000); expect(parseRetryAfter("invalid", 0)).toBeUndefined(); }); test("retries 429 responses and honors the retry limit", async () => { let requests = 0; const { url } = await serve((_request, response) => { requests += 1; response.statusCode = requests < 3 ? 429 : 200; response.setHeader("Retry-After", "0"); response.end(requests < 3 ? "slow down" : "ok"); }); const client = createKubernetesHttpLibrary({ minIntervalMs: 0, baseRetryMs: 0, maxRetries: 3, }); const response = await client .send(new RequestContext(url, HttpMethod.GET)) .toPromise(); expect(response.httpStatusCode).toBe(200); expect(requests).toBe(3); requests = 0; const limitedClient = createKubernetesHttpLibrary({ minIntervalMs: 0, baseRetryMs: 0, maxRetries: 1, }); const limitedResponse = await limitedClient .send(new RequestContext(url, HttpMethod.GET)) .toPromise(); expect(limitedResponse.httpStatusCode).toBe(429); expect(requests).toBe(2); }); test("logs Kubernetes request starts, completions, and failures", async () => { const logs: Record[] = []; const logger = { log: (entry: Record) => logs.push(entry), }; const { url } = await serve((_request, response) => response.end("ok")); const client = createKubernetesHttpLibrary({ minIntervalMs: 0, logger }); await client.send(new RequestContext(url, HttpMethod.GET)).toPromise(); expect(logs).toEqual( expect.arrayContaining([ expect.objectContaining({ event: "kuber.k8s.request.start", method: "GET", attempt: 1, }), expect.objectContaining({ event: "kuber.k8s.request.end", status: 200, responseBody: "ok", }), ]), ); const failedLogs: Record[] = []; const { url: failedUrl } = await serve((request) => request.socket.destroy(), ); const failedClient = createKubernetesHttpLibrary({ minIntervalMs: 0, logger: { log: (entry) => failedLogs.push(entry) }, }); await expect( failedClient .send(new RequestContext(failedUrl, HttpMethod.GET)) .toPromise(), ).rejects.toBeDefined(); expect(failedLogs).toEqual( expect.arrayContaining([ expect.objectContaining({ event: "kuber.k8s.request.start" }), expect.objectContaining({ event: "kuber.k8s.request.failed" }), ]), ); }); test("paces concurrent request starts", async () => { const starts: number[] = []; const { url } = await serve((_request, response) => { starts.push(Date.now()); response.end("ok"); }); const client = createKubernetesHttpLibrary({ maxConcurrent: 3, minIntervalMs: 25, }); await Promise.all( Array.from({ length: 3 }, () => client.send(new RequestContext(url, HttpMethod.GET)).toPromise(), ), ); expect(starts).toHaveLength(3); expect(starts[2]! - starts[0]!).toBeGreaterThanOrEqual(35); }); test("cancels retry waits", async () => { const controller = new AbortController(); const waiting = delay(10_000, controller.signal); controller.abort(); await expect(waiting).rejects.toMatchObject({ name: "AbortError" }); }); });