Files
kuber/server/exec-service.ts
2026-09-03 11:28:30 +07:00

730 lines
21 KiB
TypeScript

export const EXEC_MANAGED_BY_LABEL = "app.kubernetes.io/managed-by";
export const EXEC_MANAGED_BY_VALUE = "kuber";
export const EXEC_WORKSPACE_UID_LABEL = "kuber.dev/workspace-uid";
export const DEFAULT_MAX_COMMAND_ARGUMENTS = 64;
export const DEFAULT_MAX_COMMAND_BYTES = 32 * 1024;
export const DEFAULT_MAX_ARGUMENT_BYTES = 8 * 1024;
export const DEFAULT_OUTPUT_CAP_BYTES = 1024 * 1024;
export const DEFAULT_INTERACTIVE_QUEUE_CAPACITY = 128;
export const DEFAULT_MAX_STDIN_FRAME_BYTES = 64 * 1024;
const DNS_LABEL = /^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?$/;
const DNS_SUBDOMAIN =
/^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?(?:\.[a-z0-9](?:[-a-z0-9]*[a-z0-9])?)*$/;
const textEncoder = new TextEncoder();
export type ExecWorkspace = {
project: string;
uid: string;
};
export type ExecDeployment = {
name: string;
uid: string;
labels: Readonly<Record<string, string>>;
selector: Readonly<Record<string, string>>;
containers: readonly string[];
};
export type ExecPodContainer = {
name: string;
running: boolean;
ready?: boolean;
};
export type ExecPod = {
name: string;
uid: string;
/** The controlling Deployment UID, resolved through the ReplicaSet owner. */
deploymentUid: string;
phase?: string;
deletionTimestamp?: string;
containers: readonly ExecPodContainer[];
};
export type ResolvedExecTarget = {
namespace: string;
deployment: string;
deploymentUid: string;
pod: string;
podUid: string;
container: string;
};
export type KubernetesExecRequest = ResolvedExecTarget & {
command: readonly string[];
tty: boolean;
};
export type ExecChunk = string | Uint8Array;
export type KubernetesExecExit = {
exitCode: number;
reason?: string;
message?: string;
};
export interface KubernetesExecProcess {
stdout: AsyncIterable<ExecChunk>;
stderr: AsyncIterable<ExecChunk>;
writeStdin(data: Uint8Array): void | Promise<void>;
closeStdin(): void | Promise<void>;
resize(columns: number, rows: number): void | Promise<void>;
wait(): Promise<KubernetesExecExit>;
close(): void | Promise<void>;
}
/** All cluster-specific behavior, including @kubernetes/client-node, lives here. */
export interface KubernetesExecBackend {
getDeployment(
namespace: string,
name: string,
signal?: AbortSignal,
): Promise<ExecDeployment | undefined>;
listPods(
namespace: string,
selector: Readonly<Record<string, string>>,
signal?: AbortSignal,
): Promise<readonly ExecPod[]>;
exec(
request: KubernetesExecRequest,
signal: AbortSignal,
): Promise<KubernetesExecProcess>;
}
export type ExecErrorCode =
| "EXEC_INVALID"
| "EXEC_TARGET_NOT_FOUND"
| "EXEC_TARGET_FORBIDDEN"
| "EXEC_TARGET_NOT_READY"
| "EXEC_OUTPUT_LIMIT"
| "EXEC_ABORTED"
| "EXEC_FAILED";
export class ExecServiceError extends Error {
constructor(
readonly code: ExecErrorCode,
message: string,
) {
super(message);
this.name = "ExecServiceError";
}
}
export class ExecOutputLimitError extends ExecServiceError {
constructor(
readonly stream: "stdout" | "stderr",
readonly limitBytes: number,
) {
super(
"EXEC_OUTPUT_LIMIT",
`${stream} exceeded the ${limitBytes} byte output limit`,
);
this.name = "ExecOutputLimitError";
}
}
export type ExecInput = {
workspace: ExecWorkspace;
deployment: string;
container?: string;
command: readonly string[];
signal?: AbortSignal;
};
export type NonTtyExecInput = ExecInput & {
stdoutCapBytes?: number;
stderrCapBytes?: number;
};
export type NonTtyExecResult = ResolvedExecTarget &
KubernetesExecExit & {
stdout: string;
stderr: string;
};
export type ExecClientFrame =
| { type: "stdin"; data: string | Uint8Array; eof?: boolean }
| { type: "resize"; columns: number; rows: number }
| { type: "close" };
export type ExecServerFrame =
| { type: "stdout"; data: Uint8Array }
| { type: "stderr"; data: Uint8Array }
| ({ type: "exit" } & KubernetesExecExit)
| { type: "error"; code: ExecErrorCode; message: string };
export interface InteractiveExecSession extends AsyncIterable<ExecServerFrame> {
readonly target: ResolvedExecTarget;
send(frame: ExecClientFrame): Promise<void>;
close(): Promise<void>;
}
export type InteractiveExecInput = ExecInput & {
signal: AbortSignal;
tty?: boolean;
queueCapacity?: number;
};
export type ExecServiceOptions = {
maxCommandArguments?: number;
maxCommandBytes?: number;
maxArgumentBytes?: number;
maxOutputCapBytes?: number;
defaultOutputCapBytes?: number;
interactiveQueueCapacity?: number;
maxStdinFrameBytes?: number;
};
type NormalizedLimits = Required<ExecServiceOptions>;
function positiveInteger(value: number, name: string): number {
if (!Number.isSafeInteger(value) || value <= 0) {
throw new Error(`${name} must be a positive integer`);
}
return value;
}
function limitsFrom(options: ExecServiceOptions): NormalizedLimits {
const limits = {
maxCommandArguments: positiveInteger(
options.maxCommandArguments ?? DEFAULT_MAX_COMMAND_ARGUMENTS,
"maxCommandArguments",
),
maxCommandBytes: positiveInteger(
options.maxCommandBytes ?? DEFAULT_MAX_COMMAND_BYTES,
"maxCommandBytes",
),
maxArgumentBytes: positiveInteger(
options.maxArgumentBytes ?? DEFAULT_MAX_ARGUMENT_BYTES,
"maxArgumentBytes",
),
maxOutputCapBytes: positiveInteger(
options.maxOutputCapBytes ?? DEFAULT_OUTPUT_CAP_BYTES,
"maxOutputCapBytes",
),
defaultOutputCapBytes: positiveInteger(
options.defaultOutputCapBytes ?? DEFAULT_OUTPUT_CAP_BYTES,
"defaultOutputCapBytes",
),
interactiveQueueCapacity: positiveInteger(
options.interactiveQueueCapacity ?? DEFAULT_INTERACTIVE_QUEUE_CAPACITY,
"interactiveQueueCapacity",
),
maxStdinFrameBytes: positiveInteger(
options.maxStdinFrameBytes ?? DEFAULT_MAX_STDIN_FRAME_BYTES,
"maxStdinFrameBytes",
),
};
if (limits.defaultOutputCapBytes > limits.maxOutputCapBytes) {
throw new Error("defaultOutputCapBytes cannot exceed maxOutputCapBytes");
}
return limits;
}
function invalid(message: string): never {
throw new ExecServiceError("EXEC_INVALID", message);
}
function validateName(
value: string,
kind: string,
maxLength = 63,
pattern = DNS_LABEL,
): void {
if (!value || value.length > maxLength || !pattern.test(value)) {
invalid(
`${kind} must be a valid lowercase DNS name of at most ${maxLength} characters`,
);
}
}
function validateInput(input: ExecInput, limits: NormalizedLimits): void {
validateName(input.workspace.project, "Workspace project");
if (!input.workspace.uid?.trim() || input.workspace.uid.length > 256) {
invalid("Workspace UID must be between 1 and 256 characters");
}
validateName(input.deployment, "Deployment name", 253, DNS_SUBDOMAIN);
if (input.container !== undefined)
validateName(input.container, "Container name");
if (!Array.isArray(input.command) || input.command.length === 0) {
invalid("Command is required");
}
if (input.command.length > limits.maxCommandArguments) {
invalid(
`Command cannot contain more than ${limits.maxCommandArguments} arguments`,
);
}
let commandBytes = 0;
for (const argument of input.command) {
if (typeof argument !== "string" || argument.includes("\0")) {
invalid("Command arguments must be strings without null bytes");
}
const bytes = textEncoder.encode(argument).byteLength;
if (bytes > limits.maxArgumentBytes) {
invalid(`Command argument exceeds ${limits.maxArgumentBytes} bytes`);
}
commandBytes += bytes;
}
if (commandBytes > limits.maxCommandBytes) {
invalid(`Command exceeds ${limits.maxCommandBytes} bytes`);
}
}
function outputCap(
value: number | undefined,
stream: "stdout" | "stderr",
limits: NormalizedLimits,
): number {
const cap = value ?? limits.defaultOutputCapBytes;
if (
!Number.isSafeInteger(cap) ||
cap <= 0 ||
cap > limits.maxOutputCapBytes
) {
invalid(
`${stream}CapBytes must be between 1 and ${limits.maxOutputCapBytes}`,
);
}
return cap;
}
function aborted(): ExecServiceError {
return new ExecServiceError("EXEC_ABORTED", "Exec request was aborted");
}
function asBytes(chunk: ExecChunk): Uint8Array {
return typeof chunk === "string" ? textEncoder.encode(chunk) : chunk;
}
async function collect(
chunks: AsyncIterable<ExecChunk>,
stream: "stdout" | "stderr",
cap: number,
): Promise<string> {
const values: Uint8Array[] = [];
let size = 0;
for await (const chunk of chunks) {
const bytes = asBytes(chunk);
size += bytes.byteLength;
if (size > cap) throw new ExecOutputLimitError(stream, cap);
values.push(bytes);
}
const output = new Uint8Array(size);
let offset = 0;
for (const value of values) {
output.set(value, offset);
offset += value.byteLength;
}
return new TextDecoder().decode(output);
}
type QueueReader<T> = (value: IteratorResult<T>) => void;
class BoundedQueue<T> implements AsyncIterable<T> {
private readonly values: T[] = [];
private readonly readers: QueueReader<T>[] = [];
private readonly writers: Array<() => void> = [];
private closed = false;
constructor(private readonly capacity: number) {}
async push(value: T): Promise<boolean> {
while (!this.closed) {
const reader = this.readers.shift();
if (reader) {
reader({ value, done: false });
return true;
}
if (this.values.length < this.capacity) {
this.values.push(value);
return true;
}
await new Promise<void>((resolve) => this.writers.push(resolve));
}
return false;
}
close(discard = false): void {
if (this.closed) return;
this.closed = true;
if (discard) this.values.length = 0;
if (this.values.length === 0) {
for (const reader of this.readers.splice(0))
reader({ value: undefined, done: true });
}
for (const writer of this.writers.splice(0)) writer();
}
[Symbol.asyncIterator](): AsyncIterator<T> {
return {
next: () => {
const value = this.values.shift();
if (value !== undefined) {
this.writers.shift()?.();
return Promise.resolve({ value, done: false });
}
if (this.closed)
return Promise.resolve({ value: undefined, done: true });
return new Promise<IteratorResult<T>>((resolve) =>
this.readers.push(resolve),
);
},
};
}
}
function safeError(error: unknown): { code: ExecErrorCode; message: string } {
if (error instanceof ExecServiceError) {
return { code: error.code, message: error.message };
}
return {
code: "EXEC_FAILED",
message: error instanceof Error ? error.message : "Kubernetes exec failed",
};
}
async function closeQuietly(process: KubernetesExecProcess): Promise<void> {
try {
await process.close();
} catch {
// Closing an already-ended Kubernetes transport is harmless.
}
}
class DuplexSession implements InteractiveExecSession {
private readonly queue: BoundedQueue<ExecServerFrame>;
private readonly done: Promise<void>;
private closedByClient = false;
private readonly onExternalAbort: () => void;
constructor(
readonly target: ResolvedExecTarget,
private readonly process: KubernetesExecProcess,
private readonly externalSignal: AbortSignal,
private readonly controller: AbortController,
private readonly maxStdinFrameBytes: number,
queueCapacity: number,
) {
this.queue = new BoundedQueue(queueCapacity);
this.onExternalAbort = () => this.stop(true);
if (externalSignal.aborted) this.onExternalAbort();
else
externalSignal.addEventListener("abort", this.onExternalAbort, {
once: true,
});
this.done = this.run();
}
[Symbol.asyncIterator](): AsyncIterator<ExecServerFrame> {
return this.queue[Symbol.asyncIterator]();
}
async send(frame: ExecClientFrame): Promise<void> {
if (this.controller.signal.aborted) throw aborted();
switch (frame.type) {
case "stdin": {
if (
typeof frame.data !== "string" &&
!(frame.data instanceof Uint8Array)
) {
invalid("stdin data must be a string or Uint8Array");
}
const data = asBytes(frame.data);
if (data.byteLength > this.maxStdinFrameBytes) {
invalid(`stdin frame exceeds ${this.maxStdinFrameBytes} bytes`);
}
if (data.byteLength) await this.process.writeStdin(data);
if (frame.eof) await this.process.closeStdin();
return;
}
case "resize":
if (
!Number.isSafeInteger(frame.columns) ||
!Number.isSafeInteger(frame.rows) ||
frame.columns < 1 ||
frame.rows < 1 ||
frame.columns > 65_535 ||
frame.rows > 65_535
) {
invalid("Terminal dimensions must be integers between 1 and 65535");
}
await this.process.resize(frame.columns, frame.rows);
return;
case "close":
await this.close();
return;
default:
invalid("Unknown interactive exec frame");
}
}
async close(): Promise<void> {
this.closedByClient = true;
this.stop(true);
await this.done;
}
private stop(discard: boolean): void {
this.controller.abort();
this.queue.close(discard);
void closeQuietly(this.process);
}
private async pump(
stream: "stdout" | "stderr",
chunks: AsyncIterable<ExecChunk>,
): Promise<void> {
for await (const chunk of chunks) {
if (this.controller.signal.aborted) return;
await this.queue.push({ type: stream, data: asBytes(chunk) });
}
}
private async run(): Promise<void> {
try {
const [, , status] = await Promise.all([
this.pump("stdout", this.process.stdout),
this.pump("stderr", this.process.stderr),
this.process.wait(),
]);
if (!this.controller.signal.aborted) {
await this.queue.push({ type: "exit", ...status });
}
} catch (error) {
if (!this.externalSignal.aborted && !this.closedByClient) {
await this.queue.push({ type: "error", ...safeError(error) });
}
} finally {
this.externalSignal.removeEventListener("abort", this.onExternalAbort);
this.controller.abort();
await closeQuietly(this.process);
this.queue.close();
}
}
}
export class ExecService {
private readonly limits: NormalizedLimits;
constructor(
private readonly backend: KubernetesExecBackend,
options: ExecServiceOptions = {},
) {
this.limits = limitsFrom(options);
}
async resolveTarget(input: ExecInput): Promise<ResolvedExecTarget> {
validateInput(input, this.limits);
if (input.signal?.aborted) throw aborted();
const namespace = input.workspace.project;
const deployment = await this.backend.getDeployment(
namespace,
input.deployment,
input.signal,
);
if (!deployment) {
throw new ExecServiceError(
"EXEC_TARGET_NOT_FOUND",
`Deployment ${input.deployment} was not found`,
);
}
if (deployment.name !== input.deployment) {
throw new ExecServiceError(
"EXEC_TARGET_NOT_FOUND",
`Deployment ${input.deployment} was not found`,
);
}
if (
deployment.labels[EXEC_MANAGED_BY_LABEL] !== EXEC_MANAGED_BY_VALUE ||
deployment.labels[EXEC_WORKSPACE_UID_LABEL] !== input.workspace.uid
) {
throw new ExecServiceError(
"EXEC_TARGET_FORBIDDEN",
`Deployment ${input.deployment} is not owned by this workspace`,
);
}
if (!deployment.uid || Object.keys(deployment.selector).length === 0) {
throw new ExecServiceError(
"EXEC_TARGET_NOT_READY",
`Deployment ${input.deployment} has no usable pod selector`,
);
}
const requestedContainer = input.container ?? deployment.containers[0];
if (!requestedContainer) {
throw new ExecServiceError(
"EXEC_TARGET_NOT_READY",
`Deployment ${input.deployment} has no container`,
);
}
const pods = await this.backend.listPods(
namespace,
deployment.selector,
input.signal,
);
const running = pods
.filter(
(pod) =>
pod.phase === "Running" &&
!pod.deletionTimestamp &&
pod.deploymentUid === deployment.uid &&
Boolean(pod.name && pod.uid) &&
pod.containers.some(
(container) =>
container.name === requestedContainer && container.running,
),
)
.sort((left, right) => {
const leftReady = left.containers.find(
(container) => container.name === requestedContainer,
)?.ready
? 1
: 0;
const rightReady = right.containers.find(
(container) => container.name === requestedContainer,
)?.ready
? 1
: 0;
return rightReady - leftReady || left.name.localeCompare(right.name);
});
const pod = running[0];
if (!pod) {
const hasOwnedRunningPod = pods.some(
(candidate) =>
candidate.phase === "Running" &&
!candidate.deletionTimestamp &&
candidate.deploymentUid === deployment.uid,
);
throw new ExecServiceError(
input.container && hasOwnedRunningPod
? "EXEC_TARGET_NOT_FOUND"
: "EXEC_TARGET_NOT_READY",
input.container && hasOwnedRunningPod
? `Container ${requestedContainer} was not found or running in a deployment pod`
: `Deployment ${input.deployment} has no running owned pod/container`,
);
}
return {
namespace,
deployment: deployment.name,
deploymentUid: deployment.uid,
pod: pod.name,
podUid: pod.uid,
container: requestedContainer,
};
}
async execute(input: NonTtyExecInput): Promise<NonTtyExecResult> {
const stdoutCap = outputCap(input.stdoutCapBytes, "stdout", this.limits);
const stderrCap = outputCap(input.stderrCapBytes, "stderr", this.limits);
const target = await this.resolveTarget(input);
const controller = new AbortController();
let process: KubernetesExecProcess | undefined;
const onAbort = () => {
controller.abort();
if (process) void closeQuietly(process);
};
if (input.signal?.aborted) onAbort();
else input.signal?.addEventListener("abort", onAbort, { once: true });
try {
if (controller.signal.aborted) throw aborted();
process = await this.backend.exec(
{ ...target, command: [...input.command], tty: false },
controller.signal,
);
const stopOnFailure = <T>(promise: Promise<T>): Promise<T> =>
promise.catch((error) => {
controller.abort();
if (process) void closeQuietly(process);
throw error;
});
const stdout = stopOnFailure(
collect(process.stdout, "stdout", stdoutCap),
);
const stderr = stopOnFailure(
collect(process.stderr, "stderr", stderrCap),
);
const results = await Promise.allSettled([
stdout,
stderr,
process.wait(),
]);
const failure = results.find(
(result): result is PromiseRejectedResult =>
result.status === "rejected",
);
if (failure) {
controller.abort();
await closeQuietly(process);
if (input.signal?.aborted) throw aborted();
if (failure.reason instanceof ExecServiceError) throw failure.reason;
throw new ExecServiceError(
"EXEC_FAILED",
safeError(failure.reason).message,
);
}
const [stdoutResult, stderrResult, statusResult] = results as [
PromiseFulfilledResult<string>,
PromiseFulfilledResult<string>,
PromiseFulfilledResult<KubernetesExecExit>,
];
if (input.signal?.aborted) throw aborted();
return {
...target,
...statusResult.value,
stdout: stdoutResult.value,
stderr: stderrResult.value,
};
} finally {
input.signal?.removeEventListener("abort", onAbort);
controller.abort();
if (process) await closeQuietly(process);
}
}
async openInteractive(
input: InteractiveExecInput,
): Promise<InteractiveExecSession> {
const queueCapacity = positiveInteger(
input.queueCapacity ?? this.limits.interactiveQueueCapacity,
"queueCapacity",
);
validateInput(input, this.limits);
if (input.signal.aborted) throw aborted();
const target = await this.resolveTarget(input);
if (input.signal.aborted) throw aborted();
const controller = new AbortController();
const onAbort = () => controller.abort();
input.signal.addEventListener("abort", onAbort, { once: true });
try {
const process = await this.backend.exec(
{ ...target, command: [...input.command], tty: input.tty ?? true },
controller.signal,
);
input.signal.removeEventListener("abort", onAbort);
return new DuplexSession(
target,
process,
input.signal,
controller,
this.limits.maxStdinFrameBytes,
queueCapacity,
);
} catch (error) {
input.signal.removeEventListener("abort", onAbort);
if (input.signal.aborted) throw aborted();
throw error;
}
}
}
export function createExecService(
backend: KubernetesExecBackend,
options?: ExecServiceOptions,
): ExecService {
return new ExecService(backend, options);
}