import { AppsV1Api, CoreV1Api, Exec, KubeConfig, type V1Deployment, type V1Status, } from "@kubernetes/client-node"; import { PassThrough } from "node:stream"; import { type ExecDeployment, type ExecPod, type ExecPodContainer, type KubernetesExecBackend, type KubernetesExecExit, type KubernetesExecProcess, type KubernetesExecRequest, } from "./exec-service"; function selectorString( labels: Readonly>, ): string { return Object.entries(labels) .sort(([left], [right]) => left.localeCompare(right)) .map(([key, value]) => `${key}=${value}`) .join(","); } function matchLabels( deployment: V1Deployment, ): Record { 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 { 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 { 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 { 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>, _signal?: AbortSignal, ): Promise { 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 { 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((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 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 { return finished; }, close() { abort(); }, }; } }