Files
kuber/tests/server/log-service.test.ts

770 lines
26 KiB
TypeScript

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<string, string | Error>();
streams = new Map<
string,
(signal: AbortSignal) => AsyncIterable<string | Uint8Array>
>();
deploymentCalls: Array<Readonly<Record<string, string>>> = [];
podSelectors: Array<Readonly<Record<string, string>>> = [];
requests: ContainerLogRequest[] = [];
async listDeployments(
_namespace: string,
labels: Readonly<Record<string, string>>,
) {
this.deploymentCalls.push(labels);
return this.deployments;
}
async getService(_namespace: string, _name: string) {
return this.service;
}
async listPods(
_namespace: string,
selector: Readonly<Record<string, string>>,
) {
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<string> {
await untilAbort(signal);
yield* [] as string[];
}
function untilAbort(signal: AbortSignal): Promise<void> {
if (signal.aborted) return Promise.resolve();
return new Promise((resolve) =>
signal.addEventListener("abort", () => resolve(), { once: true }),
);
}
async function* chunks(...values: Array<string | Uint8Array>) {
for (const value of values) yield value;
}
async function nextOfType<T extends LogEvent["type"]>(
iterator: AsyncIterator<LogEvent>,
type: T,
limit = 30,
): Promise<Extract<LogEvent, { type: T }>> {
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<LogEvent, { type: T }>;
}
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<http.Server | https.Server> = [];
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<KubernetesLogs> {
const server = https.createServer(tls, handler);
servers.push(server);
await new Promise<void>((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<KubernetesLogs> {
const server = http.createServer(handler);
servers.push(server);
await new Promise<void>((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<void>((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<void>((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);
});
});