153 lines
4.3 KiB
TypeScript
153 lines
4.3 KiB
TypeScript
import {
|
|
AppsV1Api,
|
|
CoreV1Api,
|
|
KubeConfig,
|
|
Log,
|
|
type V1LabelSelector,
|
|
} from "@kubernetes/client-node";
|
|
import { PassThrough } from "node:stream";
|
|
import {
|
|
KubernetesLogError,
|
|
type ContainerLogRequest,
|
|
type KubernetesLogsBackend,
|
|
} from "./log-service";
|
|
|
|
function selector(labels: Readonly<Record<string, string>>): string {
|
|
return Object.entries(labels)
|
|
.sort(([left], [right]) => left.localeCompare(right))
|
|
.map(([key, value]) => `${key}=${value}`)
|
|
.join(",");
|
|
}
|
|
|
|
function matchLabels(value: V1LabelSelector | undefined): Record<string, string> {
|
|
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);
|
|
return new KubernetesLogError(
|
|
error instanceof Error ? error.message : "Kubernetes log request failed",
|
|
status === undefined || status === 408 || status === 429 || status >= 500,
|
|
);
|
|
}
|
|
|
|
export class KubernetesLogs implements KubernetesLogsBackend {
|
|
private readonly logger: Log;
|
|
|
|
constructor(
|
|
config: KubeConfig,
|
|
private readonly apps = config.makeApiClient(AppsV1Api),
|
|
private readonly core = config.makeApiClient(CoreV1Api),
|
|
) {
|
|
this.logger = new Log(config);
|
|
}
|
|
|
|
async listDeployments(
|
|
namespace: string,
|
|
labels: Readonly<Record<string, string>>,
|
|
) {
|
|
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<Record<string, string>>,
|
|
) {
|
|
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<string> {
|
|
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<Uint8Array> {
|
|
const output = new PassThrough();
|
|
let upstream: AbortController | undefined;
|
|
const abort = () => {
|
|
upstream?.abort();
|
|
output.destroy();
|
|
};
|
|
signal.addEventListener("abort", abort, { once: true });
|
|
try {
|
|
upstream = await this.logger.log(
|
|
request.namespace,
|
|
request.pod,
|
|
request.container,
|
|
output,
|
|
{
|
|
follow: true,
|
|
tailLines: request.tailLines,
|
|
sinceSeconds: request.sinceSeconds,
|
|
sinceTime: request.sinceTime,
|
|
timestamps: request.timestamps,
|
|
},
|
|
);
|
|
if (signal.aborted) abort();
|
|
for await (const chunk of output) yield new Uint8Array(chunk);
|
|
} catch (error) {
|
|
if (!signal.aborted) throw logError(error);
|
|
} finally {
|
|
signal.removeEventListener("abort", abort);
|
|
upstream?.abort();
|
|
output.destroy();
|
|
}
|
|
}
|
|
}
|