import { afterEach, describe, expect, test } from "bun:test"; import { KubeConfig } from "@kubernetes/client-node"; import http from "node:http"; import https from "node:https"; import { KubernetesLogs } from "../../server/kubernetes-logs"; import { createLogService, encodeLogEvent, KubernetesLogError, MANAGED_BY_SELECTOR, type ContainerLogRequest, type KubernetesDeployment, type KubernetesLogsBackend, type KubernetesPod, type KubernetesService, type LogEvent, } from "../../server/log-service"; class FakeBackend implements KubernetesLogsBackend { deployments: KubernetesDeployment[] = []; service: KubernetesService | undefined; pods: KubernetesPod[] = []; reads = new Map(); streams = new Map< string, (signal: AbortSignal) => AsyncIterable >(); deploymentCalls: Array>> = []; podSelectors: Array>> = []; requests: ContainerLogRequest[] = []; async listDeployments( _namespace: string, labels: Readonly>, ) { this.deploymentCalls.push(labels); return this.deployments; } async getService(_namespace: string, _name: string) { return this.service; } async listPods( _namespace: string, selector: Readonly>, ) { this.podSelectors.push(selector); return this.pods; } async readContainerLogs(request: ContainerLogRequest) { this.requests.push(request); const value = this.reads.get(`${request.pod}/${request.container}`) ?? ""; if (value instanceof Error) throw value; return value; } streamContainerLogs(request: ContainerLogRequest, signal: AbortSignal) { this.requests.push(request); const stream = this.streams.get(`${request.pod}/${request.container}`); if (stream) return stream(signal); return openStream(signal); } } async function* openStream(signal: AbortSignal): AsyncGenerator { await untilAbort(signal); yield* [] as string[]; } function untilAbort(signal: AbortSignal): Promise { if (signal.aborted) return Promise.resolve(); return new Promise((resolve) => signal.addEventListener("abort", () => resolve(), { once: true }), ); } async function* chunks(...values: Array) { for (const value of values) yield value; } async function nextOfType( iterator: AsyncIterator, type: T, limit = 30, ): Promise> { for (let index = 0; index < limit; index += 1) { const result = await iterator.next(); if (result.done) throw new Error(`Stream ended before ${type}`); if (result.value.type === type) return result.value as Extract; } throw new Error(`No ${type} event received`); } function setup() { const backend = new FakeBackend(); backend.deployments = [{ name: "web", selector: { app: "web" } }]; backend.pods = [ { name: "web-abc", uid: "uid-1", phase: "Running", containers: ["web", "sidecar"], }, ]; const service = createLogService( backend, () => new Date("2026-09-02T12:00:00.000Z"), ); return { backend, service }; } const servers: Array = []; const tlsFixture = { key: Buffer.from(`-----BEGIN PRIVATE KEY----- MIIEvAIBADANBgkqhkiG9w0BAQEFAASCBKYwggSiAgEAAoIBAQCUiNil5x0/YDDX VWBT9hQQOldUzaRdCooc+Nx3CaZaKtjIlWDJQ24GXz+bJPi9kDhNTpoGAKy7EiMj C+8pXS3ggnMEGIHF89eYzgsMlKUZeZ6x+q9blw04tdKRHRpA96Hps5MZ0FQFZ6p5 WzPqlyymvP8lAbeYhugekQjzqWfpmecNE3in7WpJH4XkM9w30u9mvjzV+5iyyIuD HI27Hqujbqm2zkEN07EU1RQMTqTXGzlz7oX5bmfhG04w/47Jz1rLJ80qISWCmVD9 8hbRsg2j67tHWMOoBhOa7qh1Tme67LCNEFdbFuc5UD1mP+22K9D/zCC2RtJQu5Sq dEJW3iovAgMBAAECggEAD2K3ckPk0y4/EOcOkdPhEyc/6ZBdkKepU8PxbkEpIpji mLBkdKSP7ogKOiNTwqsAMf3M1YdXXQ9NZXF0hgfZWzKYAFobgyo1cGYTXeu9yExB RHVPmcClRXUMCS0HDai49FC+EYPzWBX7YhOw5oFfRiw4j5hEcL+0ponmb/rhwSAg AGGIM5nGBT3wpXcbOubdxe8/Ex7vg+refUrDGs4T7sbYijoWm5utz3bdc/JxLDf0 KwMesf7RZ+eeO2C6qOB4gKNjZr5JxApcJZ2WoG9MdWkaYnU8Kyquf2vhJdDZmqtx kV0y8KSWiPTbVqXZOYcg5AKcymdF7P04zVe4rkEuoQKBgQDRDz9jem5r9U25J6VT Oi7oWZmgsDcchfeerRBhQmEPiASkOFu2VgiM3CeUeGnNmrRlRsKUEoS00Bph8AMk z0SfnBTFjIIiB9gERnqrH1vSnCbLhvnhkd7XVtoIG3JjXnJgALoIkmp9DM82VTFu kBaPpOx1i5vQyI68xpAUNhjVjwKBgQC14p5qxrwn5y0VQg/9MkgLUfbsROqRX8SP O1RRhR8teA8lL+jnTzq7HauPMnP8kEjfMHXjV124eDLKUHLrOoNwTg7S0MmgYcRU DWaOCA5u5MpHxCfz2h+OGDVp0FP/QCyLwu6vjkvK+dMhvzK92QHzTqLsia2Tpspx 3HFLEm1RYQKBgG44E7tmuQDB+5A6jrcqXcCyPISzYtru5nYJ2DDuxi1iENBjxjaD dU6OY2+rbFyxy5n5jGx0tvJ9JOutlnq5q/xaVbkxMwquB/15CwNdLRQEr49uQh/i wBHYAGt1zQEGslZbC7mpN+tl7Xk/wSgBX2OsF96BFE0m79om9Z8yRjWRAoGAAeBI iglqv26fBG0eBRqTq6o4xc8gLEe0m1WdVQnufGWUommQGXKzxGJV9rAqihxi5Ap3 7NRl3xU+UN/rj4mW+X2UoZANxF29zLAmsqhancI2Y+8eCmHhmXGee2zusN9Ulkx4 cc8h8QIKr3ptZ4/peT0CaTYyWCeMRwhjEscp4YECgYBlQyBkjrhBLswIK+Avul0C heCCZkDZstdkYTlAh6+tnf0ycMEKXPF9T1O67fSRmzQnTuvb6VWgtMh3zvOMy+Ce UnVMOhKLS0wPPTZc75kKju1C2hHZ2eSgOFbNFrQvWJMXGNjEvdscXjiCoCBbhBt+ Pc3eI1dwQpPwp363CGPLrg== -----END PRIVATE KEY-----\n`), cert: Buffer.from(`-----BEGIN CERTIFICATE----- MIIC3DCCAcSgAwIBAgIIMGH60RwCdvswDQYJKoZIhvcNAQEMBQAwFDESMBAGA1UE AxMJMTI3LjAuMC4xMB4XDTI2MDkyNzEwMjUyN1oXDTM2MDkyNDEwMjUyN1owFDES MBAGA1UEAxMJMTI3LjAuMC4xMIIBIjANBgkqhkiG9w0BAQEFAAOCAQ8AMIIBCgKC AQEAlIjYpecdP2Aw11VgU/YUEDpXVM2kXQqKHPjcdwmmWirYyJVgyUNuBl8/myT4 vZA4TU6aBgCsuxIjIwvvKV0t4IJzBBiBxfPXmM4LDJSlGXmesfqvW5cNOLXSkR0a QPeh6bOTGdBUBWeqeVsz6pcsprz/JQG3mIboHpEI86ln6ZnnDRN4p+1qSR+F5DPc N9LvZr481fuYssiLgxyNux6ro26pts5BDdOxFNUUDE6k1xs5c+6F+W5n4RtOMP+O yc9ayyfNKiElgplQ/fIW0bINo+u7R1jDqAYTmu6odU5nuuywjRBXWxbnOVA9Zj/t tivQ/8wgtkbSULuUqnRCVt4qLwIDAQABozIwMDAdBgNVHQ4EFgQUeDmjpd05dSFB yENQKuBjSFVItWQwDwYDVR0RBAgwBocEfwAAATANBgkqhkiG9w0BAQwFAAOCAQEA SZLDBqpXGPrmKfz3qvPdda3wJUp5MAvy1ZEA2IAMHMutUGemIcEXoU+4qd8aNhzz EVS2gKsaFZ59MXqC0xhRnNreqP3lP5lPVUV+EAOXvAF+vrg07KBRrxYuE2hzgKxl FP6swGVsMGyWF266PFe4p3QhcxSUvLBtCuYyJTEvL5FDhByE5M+F7zh5RFiiHvx5 txljZfJcP0wDZVUalNGbJR1nXJ7GRVPO92ec38ly6PIvlR/BAvNCoLdIF69C12qC c+5Sp92YbZ6+Lq1p75M6mY44A8FSXE40x3HXM8HcpeatXOxZuM3YAkrO337AgEx0 9CZrJPz0ekLRl4UW5LG8mA== -----END CERTIFICATE-----\n`), }; const unrelatedCa = Buffer.from(`-----BEGIN CERTIFICATE----- MIICwzCCAaugAwIBAgIISfrhnhbLiWQwDQYJKoZIhvcNAQEMBQAwEDEOMAwGA1UE AxMFd3JvbmcwHhcNMjYwOTI3MTAyNzAyWhcNMzYwOTI0MTAyNzAyWjAQMQ4wDAYD VQQDEwV3cm9uZzCCASIwDQYJKoZIhvcNAQEBBQADggEPADCCAQoCggEBAK4MZhty U3AIPjb92r7F7tsFX3UUo9fdZUeLjli2a97DZTiO/hs7a/oce27MlXq0BcpP6Rm5 d0kCIb6IU7SgP3sH8dovaLnC1Xv9KdyFnkRp0eXD6xExmZ9Y6ftDHRkEI+la41ai GdRp/ebLnTaRLxRAQUpd2xCvaa/weenGDxKFNNunVncF9WzAWYhUZqFR1wXBtGo0 0AAuzuxYoenyFLHoozeoBJ7Cr23aR/MP0zWQK3blz6TZ7fTcZtDkgsThga+6F8O+ OJGM7di6FRx7O37DRK6yVphgI1nClbEmw/sAtqc3MYQAh6CzoRFbceGQ6QEdGjq3 zrv7tJhjQ2x/CNECAwEAAaMhMB8wHQYDVR0OBBYEFKIQJVhaqJhy9cwoA1lCUPff O//lMA0GCSqGSIb3DQEBDAUAA4IBAQCQHcUDcgLie6wFWpCUFgcT0UMBcZ7FL3me 3vmvOv7T/ydOJGbgaKqz+c3NdrApsL9Jk0wmh0YelYK7TuhZnsNJ0pbehe509vJ9 10B8bKNmZqgVgdTppsaw9UDFcEQn1UMhkhtYs51+acVntNMXz8/qEhrto8mGdPKO 39ydfekUA6fAw9Q8F/sRtvG7bB6Cmj/8Tieq/dJjO7ZXgqY6HaYrX9qQxlapRu2l LgdidPbFADDoUsSXBVqX3F4Il+qRQrrGmB/TPPbGCsW/A0TA2jOPHGoh3mczjEm7 UDgBJ/wg7flQRRYcz278Hm+71JrSOEdpC4F1ssbttKE5Jhg+91uU -----END CERTIFICATE-----\n`); async function kubernetesHttpsServer( tls: typeof tlsFixture, options: { ca?: Buffer; skipTLSVerify?: boolean }, handler: http.RequestListener, ): Promise { const server = https.createServer(tls, handler); servers.push(server); await new Promise((resolve, reject) => { server.once("error", reject); server.listen(0, "127.0.0.1", resolve); }); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing port"); const config = new KubeConfig(); config.loadFromOptions({ clusters: [{ name: "test", server: `https://127.0.0.1:${address.port}`, ...(options.ca && { caData: options.ca.toString("base64") }), ...(options.skipTLSVerify && { skipTLSVerify: true }), }], users: [{ name: "test" }], contexts: [{ name: "test", cluster: "test", user: "test" }], currentContext: "test", }); return new KubernetesLogs(config); } async function kubernetesServer( handler: http.RequestListener, ): Promise { 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"); const config = new KubeConfig(); config.loadFromOptions({ clusters: [ { name: "test", server: `http://127.0.0.1:${address.port}`, skipTLSVerify: true, }, ], users: [{ name: "test" }], contexts: [{ name: "test", cluster: "test", user: "test" }], currentContext: "test", }); return new KubernetesLogs(config); } afterEach(async () => { await Promise.all( servers.splice(0).map( (server) => new Promise((resolve, reject) => server.listening ? server.close((error) => (error ? reject(error) : resolve())) : resolve(), ), ), ); }); describe("LogService", () => { test("enumerates every container in managed deployment pods", async () => { const { backend, service } = setup(); expect( await service.listContainers({ namespace: "demo", target: { kind: "managed-deployments" }, }), ).toEqual([ { namespace: "demo", targetKind: "deployment", targetName: "web", pod: "web-abc", podUid: "uid-1", container: "web", }, { namespace: "demo", targetKind: "deployment", targetName: "web", pod: "web-abc", podUid: "uid-1", container: "sidecar", }, ]); expect(backend.deploymentCalls).toEqual([MANAGED_BY_SELECTOR]); expect(backend.podSelectors).toEqual([{ app: "web" }]); }); test("uses a named service selector without listing deployments", async () => { const { backend, service } = setup(); backend.service = { name: "frontend", selector: { role: "frontend" } }; const result = await service.listContainers({ namespace: "demo", target: { kind: "service", name: "frontend" }, }); expect(result[0]).toMatchObject({ targetKind: "service", targetName: "frontend", }); expect(backend.deploymentCalls).toHaveLength(0); expect(backend.podSelectors).toEqual([{ role: "frontend" }]); }); test("rejects missing and selectorless services", async () => { const { backend, service } = setup(); const options = { namespace: "demo", target: { kind: "service" as const, name: "missing" }, }; await expect(service.listContainers(options)).rejects.toMatchObject({ retryable: false, }); backend.service = { name: "missing", selector: {} }; await expect(service.listContainers(options)).rejects.toThrow( "no pod selector", ); }); test("collects timestamped structured lines and passes log options", async () => { const { backend, service } = setup(); backend.reads.set( "web-abc/web", "2026-09-02T11:59:00.123456Z first\nsecond\n", ); backend.reads.set("web-abc/sidecar", "side\n"); const events = await service.collect({ namespace: "demo", target: { kind: "managed-deployments" }, tailLines: 25, sinceTime: new Date("2026-09-02T11:00:00Z"), timestamps: true, }); expect(events).toHaveLength(3); expect(events[0]).toMatchObject({ type: "log", timestamp: "2026-09-02T12:00:00.000Z", logTimestamp: "2026-09-02T11:59:00.123456Z", message: "first", pod: "web-abc", container: "web", }); expect(events[1]).toMatchObject({ message: "second" }); expect(backend.requests[0]).toEqual({ namespace: "demo", pod: "web-abc", container: "web", tailLines: 25, sinceSeconds: undefined, sinceTime: "2026-09-02T11:00:00.000Z", timestamps: true, }); expect(JSON.parse(encodeLogEvent(events[0]!))).toEqual(events[0]); }); test("passes a line count through a followed Kubernetes log request", async () => { let received: URL | undefined; const kubernetes = await kubernetesServer((request, response) => { received = new URL(request.url!, "http://localhost"); response.writeHead(200); response.end("latest\n"); }); const controller = new AbortController(); const iterator = kubernetes.streamContainerLogs( { namespace: "demo", pod: "web-abc", container: "web", tailLines: 12, timestamps: false, }, controller.signal, ); const chunks: Uint8Array[] = []; for await (const chunk of iterator) chunks.push(chunk); expect(new TextDecoder().decode(Buffer.concat(chunks))).toBe("latest\n"); expect(received?.searchParams.get("tailLines")).toBe("12"); expect(received?.searchParams.get("follow")).toBe("true"); }); test("streams logs from HTTPS trusted by the configured cluster CA", async () => { const tls = tlsFixture; const kubernetes = await kubernetesHttpsServer( tls, { ca: tls.cert }, (_request, response) => response.end("trusted\n"), ); const controller = new AbortController(); const chunks: Uint8Array[] = []; for await (const chunk of kubernetes.streamContainerLogs({ namespace: "demo", pod: "web-abc", container: "web", timestamps: false, }, controller.signal)) chunks.push(chunk); expect(Buffer.concat(chunks).toString()).toBe("trusted\n"); }); test("rejects HTTPS logs with an untrusted cluster CA", async () => { const tls = tlsFixture; const kubernetes = await kubernetesHttpsServer( tls, { ca: unrelatedCa }, (_request, response) => response.end("should not be reached\n"), ); const controller = new AbortController(); await expect(async () => { for await (const _chunk of kubernetes.streamContainerLogs({ namespace: "demo", pod: "web-abc", container: "web", timestamps: false, }, controller.signal)) { /* consume until TLS rejects */ } }).toThrow(); }); test("allows an untrusted HTTPS log certificate only with explicit opt-in", async () => { const tls = tlsFixture; const kubernetes = await kubernetesHttpsServer( tls, { skipTLSVerify: true }, (_request, response) => response.end("opted in\n"), ); const controller = new AbortController(); const chunks: Uint8Array[] = []; for await (const chunk of kubernetes.streamContainerLogs({ namespace: "demo", pod: "web-abc", container: "web", timestamps: false, }, controller.signal)) chunks.push(chunk); expect(Buffer.concat(chunks).toString()).toBe("opted in\n"); }); test("returns per-container safe errors and continues collection", async () => { const { backend, service } = setup(); backend.reads.set( "web-abc/web", new KubernetesLogError("container unavailable", false), ); backend.reads.set("web-abc/sidecar", "healthy"); const events = await service.collect({ namespace: "demo", target: { kind: "managed-deployments" }, }); expect(events[0]).toEqual({ type: "error", timestamp: "2026-09-02T12:00:00.000Z", namespace: "demo", pod: "web-abc", container: "web", message: "container unavailable", retryable: false, }); expect(events[1]).toMatchObject({ type: "log", message: "healthy" }); expect(JSON.stringify(events)).not.toContain("env"); expect(JSON.stringify(events)).not.toContain("Secret"); }); test("validates tail and since options before Kubernetes calls", async () => { const { backend, service } = setup(); await expect( service.collect({ namespace: "demo", target: { kind: "managed-deployments" }, tailLines: -1, }), ).rejects.toThrow("tailLines"); await expect( service.collect({ namespace: "demo", target: { kind: "managed-deployments" }, sinceSeconds: 5, sinceTime: "2026-09-02T00:00:00Z", }), ).rejects.toThrow("mutually exclusive"); expect(backend.deploymentCalls).toHaveLength(0); }); test("follows chunked lines and stops active streams on abort", async () => { const { backend, service } = setup(); let cancelled = false; backend.streams.set("web-abc/web", (signal) => (async function* () { try { yield new TextEncoder().encode("hel"); yield "lo\npartial"; await untilAbort(signal); } finally { cancelled = true; } })(), ); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, discoveryIntervalMs: 5, heartbeatIntervalMs: 1_000, }); expect(await nextOfType(iterator, "log")).toMatchObject({ message: "hello", }); controller.abort(); expect((await iterator.next()).done).toBe(true); expect(cancelled).toBe(true); }); test("discards queued lines immediately on external abort", async () => { const { backend, service } = setup(); backend.pods[0] = { ...backend.pods[0]!, containers: ["web"] }; backend.streams.set("web-abc/web", () => chunks("one\ntwo\nthree\n")); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, queueCapacity: 5, discoveryIntervalMs: 1_000, heartbeatIntervalMs: 1_000, }); expect(await nextOfType(iterator, "log")).toMatchObject({ message: "one" }); await Bun.sleep(1); controller.abort(); expect((await iterator.next()).done).toBe(true); }); test("discovers replacement pod UIDs while following", async () => { const { backend, service } = setup(); backend.pods[0] = { ...backend.pods[0]!, containers: ["web"] }; backend.streams.set("web-abc/web", () => chunks("old\n")); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, discoveryIntervalMs: 5, heartbeatIntervalMs: 1_000, }); expect(await nextOfType(iterator, "log")).toMatchObject({ message: "old", podUid: "uid-1", }); backend.pods = [ { name: "web-new", uid: "uid-2", containers: ["web"], phase: "Running" }, ]; backend.streams.set("web-new/web", () => chunks("new\n")); expect(await nextOfType(iterator, "log")).toMatchObject({ message: "new", pod: "web-new", podUid: "uid-2", }); controller.abort(); await iterator.next(); }); test("emits heartbeats and retryable stream errors before retrying", async () => { const { backend, service } = setup(); backend.pods[0] = { ...backend.pods[0]!, containers: ["web"] }; let attempts = 0; backend.streams.set("web-abc/web", () => (async function* () { attempts += 1; if (attempts === 1) throw new Error("temporary disconnect"); yield "recovered\n"; })(), ); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, discoveryIntervalMs: 5, heartbeatIntervalMs: 2, retryIntervalMs: 7, }); expect(await nextOfType(iterator, "error")).toMatchObject({ message: "temporary disconnect", retryable: true, retryAfterMs: 7, pod: "web-abc", }); expect(await nextOfType(iterator, "heartbeat")).toMatchObject({ timestamp: "2026-09-02T12:00:00.000Z", }); expect(await nextOfType(iterator, "log")).toMatchObject({ message: "recovered", }); controller.abort(); await iterator.next(); expect(attempts).toBe(2); }); test("emits a coded TLS failure once without retrying", async () => { const { backend, service } = setup(); backend.pods[0] = { ...backend.pods[0]!, containers: ["web"] }; let attempts = 0; backend.streams.set("web-abc/web", () => (async function* () { attempts += 1; yield* [] as string[]; const error = new Error("self signed certificate in certificate chain") as Error & { code: string; }; error.code = "SELF_SIGNED_CERT_IN_CHAIN"; throw error; })(), ); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, discoveryIntervalMs: 5, heartbeatIntervalMs: 1_000, }); expect(await nextOfType(iterator, "error")).toMatchObject({ message: "self signed certificate in certificate chain", retryable: false, }); expect(await nextOfType(iterator, "heartbeat")).toMatchObject({ type: "heartbeat", }); expect(attempts).toBe(1); controller.abort(); expect((await iterator.next()).done).toBe(true); expect(attempts).toBe(1); }); test("reports a broken upstream log connection per container without losing other streams", async () => { const urls: URL[] = []; const kubernetes = await kubernetesServer((request, response) => { const url = new URL(request.url!, "http://localhost"); urls.push(url); if (url.searchParams.get("container") === "web") { request.socket.destroy(); } else { response.writeHead(200); response.write("sidecar healthy\n"); } }); const { backend, service } = setup(); backend.streams.set("web-abc/web", (signal) => kubernetes.streamContainerLogs( { namespace: "demo", pod: "web-abc", container: "web", timestamps: false }, signal, ), ); backend.streams.set("web-abc/sidecar", (signal) => kubernetes.streamContainerLogs( { namespace: "demo", pod: "web-abc", container: "sidecar", timestamps: false }, signal, ), ); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, discoveryIntervalMs: 1_000, heartbeatIntervalMs: 1_000, retryIntervalMs: 1_000, }); try { const events = await Promise.all([iterator.next(), iterator.next()]); expect(events.map(({ value }) => value?.type).sort()).toEqual([ "error", "log", ]); const error = events.find(({ value }) => value?.type === "error")!.value; expect(error).toMatchObject({ pod: "web-abc", container: "web", retryable: true, }); expect(error?.message).toMatch( /socket.*connection.*closed|socket hang up|fetch failed/i, ); expect(error?.message).not.toContain("Error occurred in log request"); expect( events.find(({ value }) => value?.type === "log")!.value, ).toMatchObject({ container: "sidecar", message: "sidecar healthy" }); expect(urls).toHaveLength(2); expect(urls[0]!.searchParams.get("follow")).toBe("true"); } finally { controller.abort(); await iterator.return(undefined); } }); test.each([ [ 403, JSON.stringify({ kind: "Status", message: "logs forbidden" }), "logs forbidden", false, ], [500, "proxy unavailable", "proxy unavailable", true], ])( "keeps upstream HTTP %i status and body for follow errors", async (status, body, detail, retryable) => { const kubernetes = await kubernetesServer((_request, response) => { response.writeHead(status); response.end(body); }); const read = async () => { for await (const _chunk of kubernetes.streamContainerLogs( { namespace: "demo", pod: "web-abc", container: "web", timestamps: false, }, new AbortController().signal, )) { // The upstream response must fail before yielding log data. } }; await expect(read()).rejects.toMatchObject({ message: `Kubernetes log request failed (HTTP ${status}): ${detail}`, retryable, }); }, ); test("aborting a follow request closes the upstream stream", async () => { let disconnected!: () => void; const closed = new Promise((resolve) => (disconnected = resolve)); const kubernetes = await kubernetesServer((_request, response) => { response.on("close", disconnected); response.writeHead(200); response.write("first\n"); }); const controller = new AbortController(); const iterator = kubernetes.streamContainerLogs( { namespace: "demo", pod: "web-abc", container: "web", timestamps: false }, controller.signal, )[Symbol.asyncIterator](); expect(new TextDecoder().decode((await iterator.next()).value)).toBe( "first\n", ); controller.abort(); expect((await iterator.next()).done).toBe(true); await closed; }); test("emits one terminal target error and closes follow mode", async () => { const { service } = setup(); const iterator = service.follow({ namespace: "demo", target: { kind: "service", name: "missing" }, signal: new AbortController().signal, discoveryIntervalMs: 5, heartbeatIntervalMs: 5, }); expect(await nextOfType(iterator, "error")).toMatchObject({ message: "Service missing was not found", retryable: false, }); expect((await iterator.next()).done).toBe(true); }); test("applies bounded-queue backpressure to fast producers", async () => { const { backend, service } = setup(); backend.pods[0] = { ...backend.pods[0]!, containers: ["web"] }; let produced = 0; backend.streams.set("web-abc/web", () => (async function* () { for (let index = 0; index < 20; index += 1) { produced += 1; yield `${index}\n`; } })(), ); const controller = new AbortController(); const iterator = service.follow({ namespace: "demo", target: { kind: "managed-deployments" }, signal: controller.signal, queueCapacity: 1, discoveryIntervalMs: 1_000, heartbeatIntervalMs: 1_000, }); expect((await iterator.next()).value).toMatchObject({ message: "0" }); await Bun.sleep(10); expect(produced).toBeLessThanOrEqual(3); controller.abort(); await iterator.return?.(undefined); }); });