Files
kuber/command/logs.ts
T

213 lines
5.4 KiB
TypeScript

import { defineCommand } from "citty";
import https from "node:https";
import { createLogger } from "../lib/logger";
import { core, kc } from "../lib/k8s";
import {
getPodContainerName,
listManagedDeployments,
listPodsForDeployment,
} from "../lib/shared";
function formatError(error: unknown): string {
if (error instanceof Error) return error.message;
return String(error);
}
function delay(ms: number) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
async function streamResponseLines(
stream: NodeJS.ReadableStream,
onLine: (line: string) => void,
) {
let buffer = "";
for await (const chunk of stream) {
buffer += chunk.toString();
let newline = buffer.indexOf("\n");
while (newline !== -1) {
const line = buffer.slice(0, newline).replace(/\r$/, "");
buffer = buffer.slice(newline + 1);
if (line) onLine(line);
newline = buffer.indexOf("\n");
}
}
const line = buffer.replace(/\r$/, "");
if (line) onLine(line);
}
async function followPodLogs(
namespace: string,
podName: string,
containerName: string,
writeLine: (line: string) => void,
signal: AbortSignal,
) {
const cluster = kc.getCurrentCluster();
if (!cluster) throw new Error("No currently active cluster");
const requestURL = new URL(`${cluster.server}/api/v1/namespaces/${namespace}/pods/${podName}/log`);
requestURL.searchParams.set("container", containerName);
requestURL.searchParams.set("follow", "true");
const options: https.RequestOptions = {
method: "GET",
protocol: requestURL.protocol,
hostname: requestURL.hostname,
port: requestURL.port,
path: `${requestURL.pathname}${requestURL.search}`,
signal,
};
await kc.applyToHTTPSOptions(options);
await new Promise<void>((resolve, reject) => {
const request = https.request(options, async (response) => {
if ((response.statusCode ?? 0) < 200 || (response.statusCode ?? 0) > 299) {
const chunks: Buffer[] = [];
response.on("data", (chunk) => {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
});
response.on("end", () => {
reject(Buffer.concat(chunks).toString("utf8") || `HTTP ${response.statusCode}`);
});
return;
}
try {
await streamResponseLines(response, writeLine);
resolve();
} catch (error) {
reject(error);
}
});
request.on("error", reject);
request.end();
});
}
async function followDeploymentLogs(name: string, signal: AbortSignal) {
const logger = createLogger(name);
const activePods = new Set<string>();
while (!signal.aborted) {
const pods = await listPodsForDeployment(name);
for (const pod of pods) {
const podName = pod.metadata?.name;
if (!podName || activePods.has(podName) || pod.status?.phase !== "Running") {
continue;
}
activePods.add(podName);
void followPodLogs(
pod.metadata?.namespace ?? "default",
podName,
getPodContainerName(pod),
(line) => logger`${line}`,
signal,
)
.catch((error) => {
if (!signal.aborted) logger.warn`${podName}: ${formatError(error)}`;
})
.finally(() => {
activePods.delete(podName);
});
}
await delay(2000);
}
}
async function logDeployment(name: string, follow: boolean) {
const logger = createLogger(name);
if (!follow) {
const pods = await listPodsForDeployment(name);
if (pods.length === 0) throw new Error(`No pods found for deployment ${name}`);
for (const pod of pods) {
const podName = pod.metadata?.name;
if (!podName) continue;
try {
const text = await core.readNamespacedPodLog({
namespace: pod.metadata?.namespace ?? "default",
name: podName,
container: getPodContainerName(pod),
});
for (const line of text.split(/\r?\n/)) {
if (line) logger`${line}`;
}
} catch (error) {
logger.warn`${podName}: ${formatError(error)}`;
}
}
return;
}
const controller = new AbortController();
await new Promise<void>((resolve, reject) => {
const cleanup = () => {
process.off("SIGINT", abort);
process.off("SIGTERM", abort);
};
const abort = () => {
controller.abort();
cleanup();
resolve();
};
process.on("SIGINT", abort);
process.on("SIGTERM", abort);
void followDeploymentLogs(name, controller.signal)
.then(() => {
cleanup();
resolve();
})
.catch((error) => {
cleanup();
if (controller.signal.aborted) {
resolve();
return;
}
controller.abort();
reject(error);
});
});
}
export const logs = defineCommand({
meta: {
name: "logs",
description: "Show deployment logs",
},
args: {
follow: {
type: "boolean",
alias: "f",
description: "Follow log output",
},
},
async run({ args }) {
const [deployment] = args._;
const names = deployment
? [deployment]
: (await listManagedDeployments())
.map((item) => item.metadata?.name)
.filter((item): item is string => Boolean(item));
if (args.follow) {
await Promise.all(names.map((name) => logDeployment(name, true)));
return;
}
for (const name of names) await logDeployment(name, false);
},
});