feat: release 2.5.1
This commit is contained in:
+76
-39
@@ -2,10 +2,10 @@ import {
|
||||
AppsV1Api,
|
||||
CoreV1Api,
|
||||
KubeConfig,
|
||||
Log,
|
||||
type V1LabelSelector,
|
||||
} from "@kubernetes/client-node";
|
||||
import { PassThrough } from "node:stream";
|
||||
import http, { type IncomingMessage, type RequestOptions } from "node:http";
|
||||
import https from "node:https";
|
||||
import { logKubernetesRequest } from "../lib/request-log";
|
||||
import {
|
||||
KubernetesLogError,
|
||||
@@ -36,22 +36,46 @@ function statusCode(error: unknown): number | undefined {
|
||||
|
||||
function logError(error: unknown): KubernetesLogError {
|
||||
const status = statusCode(error);
|
||||
const code =
|
||||
error && typeof error === "object" && "code" in error
|
||||
? error.code
|
||||
: undefined;
|
||||
return new KubernetesLogError(
|
||||
error instanceof Error ? error.message : "Kubernetes log request failed",
|
||||
status === undefined || status === 408 || status === 429 || status >= 500,
|
||||
code === undefined &&
|
||||
(status === undefined || status === 408 || status === 429 || status >= 500),
|
||||
);
|
||||
}
|
||||
|
||||
export class KubernetesLogs implements KubernetesLogsBackend {
|
||||
private readonly logger: Log;
|
||||
function requestLogs(
|
||||
url: URL,
|
||||
options: RequestOptions,
|
||||
signal: AbortSignal,
|
||||
): Promise<IncomingMessage> {
|
||||
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<string> {
|
||||
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(
|
||||
config: KubeConfig,
|
||||
private readonly config: KubeConfig,
|
||||
private readonly apps = config.makeApiClient(AppsV1Api),
|
||||
private readonly core = config.makeApiClient(CoreV1Api),
|
||||
) {
|
||||
this.logger = new Log(config);
|
||||
}
|
||||
) {}
|
||||
|
||||
async listDeployments(
|
||||
namespace: string,
|
||||
@@ -125,44 +149,57 @@ export class KubernetesLogs implements KubernetesLogsBackend {
|
||||
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 logKubernetesRequest(
|
||||
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: `/api/v1/namespaces/${request.namespace}/pods/${request.pod}/log`,
|
||||
pathname: url.pathname,
|
||||
body: undefined,
|
||||
attempt: 1,
|
||||
},
|
||||
() =>
|
||||
this.logger.log(
|
||||
request.namespace,
|
||||
request.pod,
|
||||
request.container,
|
||||
output,
|
||||
{
|
||||
follow: true,
|
||||
tailLines: request.tailLines,
|
||||
sinceSeconds: request.sinceSeconds,
|
||||
sinceTime: request.sinceTime,
|
||||
timestamps: request.timestamps,
|
||||
},
|
||||
),
|
||||
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;
|
||||
},
|
||||
);
|
||||
if (signal.aborted) abort();
|
||||
for await (const chunk of output) yield new Uint8Array(chunk);
|
||||
for await (const chunk of response)
|
||||
yield new Uint8Array(chunk);
|
||||
} catch (error) {
|
||||
if (!signal.aborted) throw logError(error);
|
||||
} finally {
|
||||
signal.removeEventListener("abort", abort);
|
||||
upstream?.abort();
|
||||
output.destroy();
|
||||
if (!signal.aborted)
|
||||
throw error instanceof KubernetesLogError ? error : logError(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user