Files
kuber/server/kubernetes-exec.ts

276 lines
7.6 KiB
TypeScript

import {
AppsV1Api,
CoreV1Api,
Exec,
KubeConfig,
type V1Deployment,
type V1Status,
} from "@kubernetes/client-node";
import { PassThrough } from "node:stream";
import { logKubernetesRequest } from "../lib/request-log";
import {
type ExecDeployment,
type ExecPod,
type ExecPodContainer,
type KubernetesExecBackend,
type KubernetesExecExit,
type KubernetesExecProcess,
type KubernetesExecRequest,
} from "./exec-service";
function selectorString(labels: Readonly<Record<string, string>>): string {
return Object.entries(labels)
.sort(([left], [right]) => left.localeCompare(right))
.map(([key, value]) => `${key}=${value}`)
.join(",");
}
function matchLabels(deployment: V1Deployment): Record<string, string> {
const matchExpressions = deployment.spec?.selector?.matchExpressions;
if (matchExpressions?.length) return {};
return deployment.spec?.selector?.matchLabels ?? {};
}
async function ownerDeploymentUid(
apps: AppsV1Api,
namespace: string,
ownerReferences:
| Array<{
kind?: string;
name?: string;
uid?: string;
}>
| undefined,
): Promise<string> {
if (!ownerReferences?.length) return "";
for (const ref of ownerReferences) {
if (ref.kind === "ReplicaSet" && ref.name && ref.uid) {
try {
const rs = await apps.readNamespacedReplicaSet({
namespace,
name: ref.name,
});
const rsOwners = rs.metadata?.ownerReferences;
if (rsOwners?.length) {
for (const rsRef of rsOwners) {
if (rsRef.kind === "Deployment" && rsRef.uid) return rsRef.uid;
}
}
} catch {
// ReplicaSet not found or inaccessible; fall through to empty
}
}
}
return "";
}
async function* passthrough(
stream: PassThrough,
signal: AbortSignal,
): AsyncGenerator<Uint8Array> {
try {
for await (const chunk of stream) {
if (signal.aborted) return;
yield chunk instanceof Buffer
? new Uint8Array(chunk.buffer, chunk.byteOffset, chunk.byteLength)
: new Uint8Array(chunk);
}
} finally {
stream.destroy();
}
}
function exitCodeFromStatus(status: V1Status): number | undefined {
const messages: string[] = [];
if (status.message) messages.push(status.message);
for (const cause of status.details?.causes ?? []) {
if (cause.reason) messages.push(cause.reason);
if (cause.message) messages.push(cause.message);
}
for (const message of messages) {
const match = /exit code\s*(\d+)/i.exec(message);
if (match) {
const code = Number(match[1]);
if (Number.isSafeInteger(code) && code >= 0) return code;
}
}
return undefined;
}
export class KubernetesExec implements KubernetesExecBackend {
private readonly apps: AppsV1Api;
private readonly core: CoreV1Api;
private readonly execClient: Exec;
constructor(config: KubeConfig) {
this.apps = config.makeApiClient(AppsV1Api);
this.core = config.makeApiClient(CoreV1Api);
this.execClient = new Exec(config);
}
async getDeployment(
namespace: string,
name: string,
_signal?: AbortSignal,
): Promise<ExecDeployment | undefined> {
try {
const deployment = await this.apps.readNamespacedDeployment({
namespace,
name,
});
const uid = deployment.metadata?.uid;
const labels = deployment.metadata?.labels ?? {};
const selector = matchLabels(deployment);
const containers =
deployment.spec?.template?.spec?.containers
?.map((container) => container.name ?? "")
.filter(Boolean) ?? [];
if (!uid || Object.keys(selector).length === 0) return undefined;
return { name, uid, labels, selector, containers };
} catch {
return undefined;
}
}
async listPods(
namespace: string,
selector: Readonly<Record<string, string>>,
_signal?: AbortSignal,
): Promise<readonly ExecPod[]> {
const result = await this.core.listNamespacedPod({
namespace,
labelSelector: selectorString(selector),
});
const pods: ExecPod[] = [];
for (const pod of result.items) {
const name = pod.metadata?.name;
const uid = pod.metadata?.uid;
if (!name || !uid) continue;
const deploymentUid = await ownerDeploymentUid(
this.apps,
namespace,
pod.metadata?.ownerReferences,
);
const containerStatuses = pod.status?.containerStatuses ?? [];
const containers: ExecPodContainer[] = (pod.spec?.containers ?? []).map(
(container) => {
const status = containerStatuses.find(
(item) => item.name === container.name,
);
return {
name: container.name ?? "",
running: status?.state?.running !== undefined,
ready: status?.ready,
};
},
);
pods.push({
name,
uid,
deploymentUid,
phase: pod.status?.phase,
deletionTimestamp: pod.metadata?.deletionTimestamp
? new Date(pod.metadata.deletionTimestamp).toISOString()
: undefined,
containers,
});
}
return pods;
}
async exec(
request: KubernetesExecRequest,
signal: AbortSignal,
): Promise<KubernetesExecProcess> {
const stdout = new PassThrough();
const stderr = new PassThrough();
const stdin = new PassThrough();
let completed: (exit: KubernetesExecExit) => void = () => {};
let failed: (error: unknown) => void = () => {};
const finished = new Promise<KubernetesExecExit>((resolve, reject) => {
completed = resolve;
failed = reject;
});
let close: () => void = () => {};
const abort = () => {
close();
stdout.destroy();
stderr.destroy();
stdin.destroy();
};
signal.addEventListener("abort", abort, { once: true });
try {
const ws = await logKubernetesRequest(
{
method: "GET",
pathname: `/api/v1/namespaces/${request.namespace}/pods/${request.pod}/exec`,
body: {
container: request.container,
command: request.command,
tty: request.tty,
},
attempt: 1,
},
() =>
this.execClient.exec(
request.namespace,
request.pod,
request.container,
[...request.command],
stdout,
stderr,
stdin,
request.tty,
(status) => {
completed({
exitCode: exitCodeFromStatus(status) ?? 1,
reason: status.reason,
message: status.message,
});
},
),
);
close = () => {
try {
ws.close();
} catch {
// already closed
}
};
if (signal.aborted) abort();
} catch (error) {
signal.removeEventListener("abort", abort);
stdout.destroy();
stderr.destroy();
stdin.destroy();
failed(error);
throw error;
}
return {
stdout: passthrough(stdout, signal),
stderr: passthrough(stderr, signal),
writeStdin(data: Uint8Array) {
if (!stdin.destroyed) stdin.write(data);
},
closeStdin() {
if (!stdin.destroyed) stdin.end();
},
resize(_columns: number, _rows: number) {
// TTY resize is delivered through the exec attach resize stream.
// @kubernetes/client-node does not expose a live resize API during an
// ongoing attach; the initial size is passed at exec start.
},
wait(): Promise<KubernetesExecExit> {
return finished;
},
close() {
abort();
},
};
}
}