feat: add restart command for managed deployments and enhance logging functionality
This commit is contained in:
+138
-68
@@ -1,55 +1,134 @@
|
||||
import { defineCommand } from "citty";
|
||||
import { finished } from "node:stream/promises";
|
||||
import { Writable } from "node:stream";
|
||||
import https from "node:https";
|
||||
import { createLogger } from "../lib/logger";
|
||||
import { core, logClient } from "../lib/k8s";
|
||||
import { core, kc } from "../lib/k8s";
|
||||
import {
|
||||
getPodContainerName,
|
||||
listManagedDeployments,
|
||||
listPodsForDeployment,
|
||||
} from "../lib/shared";
|
||||
|
||||
class LoggerStream extends Writable {
|
||||
#buffer = "";
|
||||
|
||||
constructor(private readonly writeLine: (line: string) => void) {
|
||||
super();
|
||||
}
|
||||
|
||||
override _write(chunk: Buffer | string, _encoding: BufferEncoding, callback: (error?: Error | null) => void) {
|
||||
this.#buffer += chunk.toString();
|
||||
|
||||
let newline = this.#buffer.indexOf("\n");
|
||||
while (newline !== -1) {
|
||||
const line = this.#buffer.slice(0, newline).replace(/\r$/, "");
|
||||
if (line) this.writeLine(line);
|
||||
this.#buffer = this.#buffer.slice(newline + 1);
|
||||
newline = this.#buffer.indexOf("\n");
|
||||
}
|
||||
|
||||
callback();
|
||||
}
|
||||
|
||||
override _final(callback: (error?: Error | null) => void) {
|
||||
const line = this.#buffer.replace(/\r$/, "");
|
||||
if (line) this.writeLine(line);
|
||||
this.#buffer = "";
|
||||
callback();
|
||||
}
|
||||
}
|
||||
|
||||
function formatError(error: unknown): string {
|
||||
if (error instanceof Error) return error.message;
|
||||
return String(error);
|
||||
}
|
||||
|
||||
async function logDeployment(name: string, follow: boolean) {
|
||||
const pods = await listPodsForDeployment(name);
|
||||
if (pods.length === 0) throw new Error(`No pods found for deployment ${name}`);
|
||||
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;
|
||||
@@ -71,45 +150,36 @@ async function logDeployment(name: string, follow: boolean) {
|
||||
return;
|
||||
}
|
||||
|
||||
const controllers = await Promise.all(
|
||||
pods.map(async (pod) => {
|
||||
const podName = pod.metadata?.name;
|
||||
if (!podName) return;
|
||||
const controller = new AbortController();
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const cleanup = () => {
|
||||
process.off("SIGINT", abort);
|
||||
process.off("SIGTERM", abort);
|
||||
};
|
||||
|
||||
const stream = new LoggerStream((line) => logger`${line}`);
|
||||
try {
|
||||
const controller = await logClient.log(
|
||||
pod.metadata?.namespace ?? "default",
|
||||
podName,
|
||||
getPodContainerName(pod),
|
||||
stream,
|
||||
{ follow: true },
|
||||
);
|
||||
|
||||
return { controller, stream };
|
||||
} catch (error) {
|
||||
logger.warn`${podName}: ${formatError(error)}`;
|
||||
stream.end();
|
||||
}
|
||||
}),
|
||||
);
|
||||
|
||||
const activeControllers = controllers.filter(
|
||||
(item): item is NonNullable<(typeof controllers)[number]> => Boolean(item),
|
||||
);
|
||||
if (activeControllers.length === 0) return;
|
||||
|
||||
await new Promise<void>((resolve) => {
|
||||
const abort = () => {
|
||||
for (const item of activeControllers) item.controller.abort();
|
||||
controller.abort();
|
||||
cleanup();
|
||||
resolve();
|
||||
};
|
||||
|
||||
process.once("SIGINT", abort);
|
||||
process.once("SIGTERM", abort);
|
||||
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);
|
||||
});
|
||||
});
|
||||
|
||||
await Promise.all(activeControllers.map((item) => finished(item.stream)));
|
||||
}
|
||||
|
||||
export const logs = defineCommand({
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
import { defineCommand } from "citty";
|
||||
import { Listr } from "listr2";
|
||||
import { assertManagedNamespace } from "../lib/apply";
|
||||
import {
|
||||
getProject,
|
||||
listManagedDeployments,
|
||||
restartDeployment,
|
||||
} from "../lib/shared";
|
||||
|
||||
export const restart = defineCommand({
|
||||
meta: {
|
||||
name: "restart",
|
||||
description: "Roll out a restart for managed deployments",
|
||||
},
|
||||
async run() {
|
||||
await assertManagedNamespace(getProject());
|
||||
const deployments = await listManagedDeployments();
|
||||
|
||||
await new Listr([
|
||||
{
|
||||
title: "Restart deployments",
|
||||
skip: () =>
|
||||
deployments.length > 0 ? false : "No managed deployments found",
|
||||
task: (_ctx, task) => {
|
||||
task.output = `${deployments.length} deployments queued`;
|
||||
return task.newListr(
|
||||
deployments.map((deployment) => ({
|
||||
title: `Deployment ${deployment.metadata?.name}`,
|
||||
task: () => restartDeployment(deployment.metadata!.name!),
|
||||
})),
|
||||
{ concurrent: false, exitOnError: true },
|
||||
);
|
||||
},
|
||||
},
|
||||
]).run();
|
||||
},
|
||||
});
|
||||
Reference in New Issue
Block a user