import { AppsV1Api, CoreV1Api, KubeConfig, type V1LabelSelector, } from "@kubernetes/client-node"; import http, { type IncomingMessage, type RequestOptions } from "node:http"; import https from "node:https"; import { logKubernetesRequest } from "../lib/request-log"; import { KubernetesLogError, type ContainerLogRequest, type KubernetesLogsBackend, } from "./log-service"; function selector(labels: Readonly>): string { return Object.entries(labels) .sort(([left], [right]) => left.localeCompare(right)) .map(([key, value]) => `${key}=${value}`) .join(","); } function matchLabels( value: V1LabelSelector | undefined, ): Record { if (value?.matchExpressions?.length) return {}; return value?.matchLabels ?? {}; } function statusCode(error: unknown): number | undefined { if (!error || typeof error !== "object") return; if ("code" in error && typeof error.code === "number") return error.code; if ("statusCode" in error && typeof error.statusCode === "number") return error.statusCode; } function logError(error: unknown): KubernetesLogError { const status = statusCode(error); const code = error && typeof error === "object" && "code" in error ? error.code : undefined; const transientNetworkCodes = new Set([ "ECONNRESET", "ETIMEDOUT", "EPIPE", "ECONNREFUSED", "EHOSTUNREACH", "ENETUNREACH", "ECONNABORTED", "EAI_AGAIN", "UND_ERR_SOCKET", "UND_ERR_CONNECT_TIMEOUT", ]); const message = error instanceof Error ? error.message : ""; return new KubernetesLogError( error instanceof Error ? error.message : "Kubernetes log request failed", typeof code === "string" ? transientNetworkCodes.has(code) : status === undefined && /socket hang up|socket closed/i.test(message) ? true : status === undefined || status === 408 || status === 429 || status >= 500, ); } function requestLogs( url: URL, options: RequestOptions, signal: AbortSignal, ): Promise { return new Promise((resolve, reject) => { const transport = url.protocol === "https:" ? https : http; const request = transport.request( url, { ...options, method: "GET", signal }, resolve, ); request.once("error", reject); request.end(); }); } async function responseBody(response: IncomingMessage): Promise { const chunks: Buffer[] = []; for await (const chunk of response) chunks.push(Buffer.from(chunk)); return Buffer.concat(chunks).toString("utf8"); } export class KubernetesLogs implements KubernetesLogsBackend { constructor( private readonly config: KubeConfig, private readonly apps = config.makeApiClient(AppsV1Api), private readonly core = config.makeApiClient(CoreV1Api), ) {} async listDeployments( namespace: string, labels: Readonly>, ) { try { const deployments = await this.apps.listNamespacedDeployment({ namespace, labelSelector: selector(labels), }); return deployments.items.flatMap((deployment) => { const name = deployment.metadata?.name; const labels = matchLabels(deployment.spec?.selector); return name && Object.keys(labels).length ? [{ name, selector: labels }] : []; }); } catch (error) { throw logError(error); } } async getService(namespace: string, name: string) { try { const service = await this.core.readNamespacedService({ namespace, name, }); return { name, selector: service.spec?.selector ?? {} }; } catch (error) { if (statusCode(error) === 404) return; throw logError(error); } } async listPods(namespace: string, labels: Readonly>) { try { const pods = await this.core.listNamespacedPod({ namespace, labelSelector: selector(labels), }); return pods.items.map((pod) => ({ name: pod.metadata?.name ?? "", uid: pod.metadata?.uid ?? "", phase: pod.status?.phase, containers: (pod.spec?.containers ?? []).map( (container) => container.name, ), })); } catch (error) { throw logError(error); } } async readContainerLogs(request: ContainerLogRequest): Promise { try { return await this.core.readNamespacedPodLog({ namespace: request.namespace, name: request.pod, container: request.container, tailLines: request.tailLines, sinceSeconds: request.sinceSeconds, timestamps: request.timestamps, }); } catch (error) { throw logError(error); } } async *streamContainerLogs( request: ContainerLogRequest, signal: AbortSignal, ): AsyncIterable { try { if (signal.aborted) return; const cluster = this.config.getCurrentCluster(); if (!cluster) throw new Error("No currently active cluster"); const url = new URL( `${cluster.server}/api/v1/namespaces/${encodeURIComponent(request.namespace)}/pods/${encodeURIComponent(request.pod)}/log`, ); url.searchParams.set("container", request.container); url.searchParams.set("follow", "true"); url.searchParams.set("timestamps", String(request.timestamps)); if (request.tailLines !== undefined) url.searchParams.set("tailLines", String(request.tailLines)); if (request.sinceSeconds !== undefined) url.searchParams.set("sinceSeconds", String(request.sinceSeconds)); if (request.sinceTime !== undefined) url.searchParams.set("sinceTime", request.sinceTime); const options: RequestOptions = {}; await this.config.applyToHTTPSOptions(options); const response = await logKubernetesRequest( { method: "GET", pathname: url.pathname, body: undefined, attempt: 1, }, async () => { const response = await requestLogs(url, options, signal); const status = response.statusCode ?? 0; if (status < 200 || status >= 300) { const body = await responseBody(response); let detail = body.trim(); try { const status = JSON.parse(body) as { message?: unknown }; if (typeof status.message === "string") detail = status.message; } catch { // Proxies can return plain text rather than a Kubernetes Status. } throw new KubernetesLogError( `Kubernetes log request failed (HTTP ${status})${detail ? `: ${detail.slice(0, 512)}` : ""}`, status === 408 || status === 429 || status >= 500, ); } return response; }, ); for await (const chunk of response) yield new Uint8Array(chunk); } catch (error) { if (!signal.aborted) throw error instanceof KubernetesLogError ? error : logError(error); } } }