diff --git a/src/mgr/kubernetes-runner-job.ts b/src/mgr/kubernetes-runner-job.ts index a100fc4..e6171a4 100644 --- a/src/mgr/kubernetes-runner-job.ts +++ b/src/mgr/kubernetes-runner-job.ts @@ -1,4 +1,3 @@ -import { spawn } from "node:child_process"; import { AgentRunError } from "../common/errors.js"; import { redactJson, redactText } from "../common/redaction.js"; import { isLeaseExpired, isTerminalCommandState, isTerminalRunStatus, summarizeResourceBundleRef, summarizeSessionRef } from "./store.js"; @@ -13,7 +12,8 @@ import { ensureSessionPvc } from "./session-pvc.js"; import { gitTransportSummary } from "../common/git-transport.js"; import { enforceRunnerRetentionBeforeCreate } from "./runner-retention.js"; import type { RunnerRetentionOptions, RunnerRetentionSummary } from "./runner-retention.js"; -import { throwIfAborted, waitForChildProcess } from "./abortable-child.js"; +import { throwIfAborted } from "./abortable-child.js"; +import { createKubernetesClient, type KubernetesClient, type KubernetesResource, type KubernetesTransportResult } from "./kubernetes-transport.js"; const reusableCredentialEnvNames = new Set([ "AUTH_PASSWORD", @@ -248,27 +248,31 @@ async function createKubernetesRunnerJobUnderFence(options: { store: AgentRunSto }; const render = renderRunnerJobManifest({ ...renderOptions, attemptId, ...(runnerId ? { runnerId } : {}), ...(runnerJobId ? { runnerJobId } : {}) }); throwIfAborted(options.signal, "runner dispatch attempt"); - const kubectlCommand = options.defaults.kubectlCommand ?? "kubectl"; + const kubernetes = createKubernetesClient({ + namespace: render.namespace, + ...(options.defaults.kubectlCommand ? { kubectlCommand: options.defaults.kubectlCommand } : {}), + ...(options.signal ? { signal: options.signal } : {}), + }); let transientEnvSecretCreated = false; let transientEnvSecretOwnerAttached = false; let created: JsonRecord | null = null; let adopted = false; try { if (transientEnvSecretName) { - const secretCreate = await kubectlCreateOrAdopt(transientEnvSecretManifest({ namespace: render.namespace, name: transientEnvSecretName, runId: run.id, commandId, attemptId: render.attemptId, runnerId: render.runnerId, jobName: render.jobName, lane: lane ?? "v0.1", items: transientEnv }), kubectlCommand, "runner transient env secret", options.signal, { create: "runner-transient-env-secret-create-failed", adopt: "runner-transient-env-secret-adopt-failed" }); + const secretCreate = await kubernetesCreateOrAdopt(transientEnvSecretManifest({ namespace: render.namespace, name: transientEnvSecretName, runId: run.id, commandId, attemptId: render.attemptId, runnerId: render.runnerId, jobName: render.jobName, lane: lane ?? "v0.1", items: transientEnv }), kubernetes, "runner transient env secret", options.signal, { create: "runner-transient-env-secret-create-failed", adopt: "runner-transient-env-secret-adopt-failed" }); transientEnvSecretCreated = !secretCreate.adopted; } - const jobCreate = await kubectlCreateOrAdopt(render.manifest, kubectlCommand, "runner job", options.signal, { create: "runner-kubernetes-job-create-failed", adopt: "runner-kubernetes-job-adopt-failed" }); + const jobCreate = await kubernetesCreateOrAdopt(render.manifest, kubernetes, "runner job", options.signal, { create: "runner-kubernetes-job-create-failed", adopt: "runner-kubernetes-job-adopt-failed" }); created = jobCreate.object; adopted = jobCreate.adopted; } catch (error) { - if (transientEnvSecretName && transientEnvSecretCreated) await kubectlDeleteSecret(transientEnvSecretName, render.namespace, kubectlCommand, options.signal); + if (transientEnvSecretName && transientEnvSecretCreated) await kubernetesDeleteSecret(transientEnvSecretName, kubernetes); throw error; } - if (!created) throw new AgentRunError("infra-failed", "kubectl did not return created runner job metadata", { httpStatus: 502, details: { reason: "runner-kubernetes-job-create-failed", changeReason: "created-object-missing", valuesPrinted: false } }); + if (!created) throw new AgentRunError("infra-failed", "Kubernetes transport did not return created runner job metadata", { httpStatus: 502, details: { reason: "runner-kubernetes-job-create-failed", changeReason: "created-object-missing", valuesPrinted: false } }); throwIfAborted(options.signal, "runner dispatch attempt"); if (transientEnvSecretName) { - const owner = await kubectlPatchSecretOwnerReference(transientEnvSecretName, render.namespace, { name: render.jobName, uid: objectPath(created, ["metadata", "uid"]) }, kubectlCommand, options.signal); + const owner = await kubernetesPatchSecretOwnerReference(transientEnvSecretName, { name: render.jobName, uid: objectPath(created, ["metadata", "uid"]) }, kubernetes); transientEnvSecretOwnerAttached = owner.ok; if (!owner.ok) render.warnings.push("transientEnv Secret ownerReference patch failed; Kubernetes TTL may not garbage-collect this per-job Secret automatically"); } @@ -518,45 +522,35 @@ function summarizeTransientEnvSecret(name: string, namespace: string, items: Run }; } -async function kubectlCreateOrAdopt(manifest: JsonRecord, kubectlCommand: string, label: string, abortSignal: AbortSignal | undefined, failureReasons: { create: string; adopt: string }): Promise<{ object: JsonRecord; adopted: boolean }> { - throwIfAborted(abortSignal, `kubectl create ${label}`); - const child = spawn(kubectlCommand, ["create", "-f", "-", "-o", "json"], { stdio: ["pipe", "pipe", "pipe"] }); - let stdout = ""; - let stderr = ""; - child.stdout.setEncoding("utf8"); - child.stderr.setEncoding("utf8"); - child.stdout.on("data", (chunk) => { stdout += String(chunk); }); - child.stderr.on("data", (chunk) => { stderr += String(chunk); }); - const completion = waitForChildProcess(child, { ...(abortSignal ? { signal: abortSignal } : {}), label: `kubectl create ${label}` }); - child.stdin.on("error", () => {}); - child.stdin.end(`${JSON.stringify(manifest)}\n`); - const result = await completion; +async function kubernetesCreateOrAdopt(manifest: JsonRecord, kubernetes: KubernetesClient, label: string, abortSignal: AbortSignal | undefined, failureReasons: { create: string; adopt: string }): Promise<{ object: JsonRecord; adopted: boolean }> { + throwIfAborted(abortSignal, `Kubernetes create ${label}`); + const result = await kubernetes.create(manifest, { sensitiveBody: manifest.kind === "Secret" }); if (result.code !== 0) { - if (/alreadyexists|already exists/iu.test(stderr.replace(/\s+/gu, ""))) { + if (result.kubernetesReason === "AlreadyExists" || /alreadyexists|already exists/iu.test(result.stderr.replace(/\s+/gu, ""))) { const metadata = objectRecord(manifest.metadata); const name = optionalString(metadata.name); const namespace = optionalString(metadata.namespace); const kind = optionalString(manifest.kind)?.toLowerCase(); if (!name || !namespace || !kind) throw new AgentRunError("infra-failed", `cannot adopt ${label} without deterministic identity`, { httpStatus: 502, details: { reason: failureReasons.adopt, changeReason: "manifest-identity-missing", valuesPrinted: false } }); - const existing = await kubectlRun(kubectlCommand, ["get", kind, name, "-n", namespace, "-o", "json"], abortSignal); - if (existing.code !== 0) throw new AgentRunError("infra-failed", `kubectl get ${label} failed during adopt`, { httpStatus: 502, details: { reason: failureReasons.adopt, changeReason: "existing-object-observation-failed", stderr: redactText(existing.stderr.slice(-2000)), valuesPrinted: false } }); - const object = parseKubectlObject(existing.stdout, `${label} adopt`, failureReasons.adopt); + const existing = await kubernetes.get(kubernetesResourceForKind(kind), name); + if (existing.code !== 0) throw new AgentRunError("infra-failed", `Kubernetes get ${label} failed during adopt`, { httpStatus: 502, details: { reason: failureReasons.adopt, changeReason: "existing-object-observation-failed", transport: existing.transport, statusCode: existing.statusCode, kubernetesReason: existing.kubernetesReason, stderr: redactText(existing.stderr.slice(-2000)), valuesPrinted: false } }); + const object = parseKubernetesObject(existing.stdout, `${label} adopt`, failureReasons.adopt); assertAdoptedIdentity(manifest, object, label, failureReasons.adopt); return { object, adopted: true }; } - throw new AgentRunError("infra-failed", `kubectl create ${label} failed with code ${result.code}`, { httpStatus: 502, details: redactJson({ reason: failureReasons.create, exitCode: result.code, stderr: redactText(stderr.slice(-4000)), stdout: redactText(stdout.slice(-2000)), signal: result.signal, valuesPrinted: false }) }); + throw new AgentRunError("infra-failed", `Kubernetes create ${label} failed`, { httpStatus: 502, details: redactJson(kubernetesFailureDetails(result, failureReasons.create)) }); } - return { object: parseKubectlObject(stdout, label, failureReasons.create), adopted: false }; + return { object: parseKubernetesObject(result.stdout, label, failureReasons.create), adopted: false }; } -function parseKubectlObject(stdout: string, label: string, reason: string): JsonRecord { +function parseKubernetesObject(stdout: string, label: string, reason: string): JsonRecord { try { const parsed = JSON.parse(stdout) as unknown; if (typeof parsed === "object" && parsed !== null && !Array.isArray(parsed)) return kubernetesObjectSummary(parsed as JsonRecord); } catch (error) { - throw new AgentRunError("infra-failed", `kubectl returned invalid JSON for ${label}: ${error instanceof Error ? error.message : String(error)}`, { httpStatus: 502, details: { reason, changeReason: "invalid-json", stdoutHash: stableHash(stdout), valuesPrinted: false } }); + throw new AgentRunError("infra-failed", `Kubernetes transport returned invalid JSON for ${label}: ${error instanceof Error ? error.message : String(error)}`, { httpStatus: 502, details: { reason, changeReason: "invalid-json", stdoutHash: stableHash(stdout), valuesPrinted: false } }); } - throw new AgentRunError("infra-failed", `kubectl returned non-object JSON for ${label}`, { httpStatus: 502, details: { reason, changeReason: "non-object-json", stdoutHash: stableHash(stdout), valuesPrinted: false } }); + throw new AgentRunError("infra-failed", `Kubernetes transport returned non-object JSON for ${label}`, { httpStatus: 502, details: { reason, changeReason: "non-object-json", stdoutHash: stableHash(stdout), valuesPrinted: false } }); } function kubernetesObjectSummary(object: JsonRecord): JsonRecord { @@ -594,32 +588,39 @@ function objectRecord(value: unknown): JsonRecord { return typeof value === "object" && value !== null && !Array.isArray(value) ? value as JsonRecord : {}; } -async function kubectlDeleteSecret(name: string, namespace: string, kubectlCommand: string, abortSignal?: AbortSignal): Promise { - await kubectlRun(kubectlCommand, ["delete", "secret", name, "-n", namespace, "--ignore-not-found=true"], abortSignal); +async function kubernetesDeleteSecret(name: string, kubernetes: KubernetesClient): Promise { + await kubernetes.delete("secret", name, { ignoreNotFound: true }); } -async function kubectlPatchSecretOwnerReference(name: string, namespace: string, owner: { name: string; uid: string | null }, kubectlCommand: string, abortSignal?: AbortSignal): Promise<{ ok: boolean }> { +async function kubernetesPatchSecretOwnerReference(name: string, owner: { name: string; uid: string | null }, kubernetes: KubernetesClient): Promise<{ ok: boolean }> { if (!owner.uid) return { ok: false }; const patch = { metadata: { ownerReferences: [{ apiVersion: "batch/v1", kind: "Job", name: owner.name, uid: owner.uid, controller: false, blockOwnerDeletion: false }], }, }; - const result = await kubectlRun(kubectlCommand, ["patch", "secret", name, "-n", namespace, "--type", "merge", "-p", JSON.stringify(patch), "-o", "json"], abortSignal); + const result = await kubernetes.patch("secret", name, patch); return { ok: result.code === 0 }; } -async function kubectlRun(kubectlCommand: string, args: string[], abortSignal?: AbortSignal): Promise<{ code: number | null; signal: NodeJS.Signals | null; stdout: string; stderr: string }> { - throwIfAborted(abortSignal, "kubectl"); - const child = spawn(kubectlCommand, args, { stdio: ["ignore", "pipe", "pipe"] }); - let stdout = ""; - let stderr = ""; - child.stdout.setEncoding("utf8"); - child.stderr.setEncoding("utf8"); - child.stdout.on("data", (chunk) => { stdout += String(chunk); }); - child.stderr.on("data", (chunk) => { stderr += String(chunk); }); - const result = await waitForChildProcess(child, { ...(abortSignal ? { signal: abortSignal } : {}), label: "kubectl" }); - return { ...result, stdout, stderr }; +function kubernetesResourceForKind(kind: string): KubernetesResource { + if (kind === "job") return "job"; + if (kind === "secret") return "secret"; + throw new AgentRunError("infra-failed", `cannot adopt unsupported Kubernetes kind ${kind}`, { httpStatus: 502, details: { reason: "runner-kubernetes-adopt-kind-unsupported", kind, valuesPrinted: false } }); +} + +function kubernetesFailureDetails(result: KubernetesTransportResult, reason: string): JsonRecord { + return { + reason, + exitCode: result.code, + signal: result.signal, + transport: result.transport, + statusCode: result.statusCode, + kubernetesReason: result.kubernetesReason, + stderr: redactText(result.stderr.slice(-2000)), + stdoutHash: result.stdout ? stableHash(result.stdout) : null, + valuesPrinted: false, + }; } function stringField(record: JsonRecord, key: string): string { diff --git a/src/mgr/kubernetes-transport.ts b/src/mgr/kubernetes-transport.ts new file mode 100644 index 0000000..7f941ce --- /dev/null +++ b/src/mgr/kubernetes-transport.ts @@ -0,0 +1,444 @@ +import { readFile } from "node:fs/promises"; +import { request as httpsRequest } from "node:https"; +import { spawn } from "node:child_process"; +import type { IncomingHttpHeaders } from "node:http"; +import type { RequestOptions } from "node:https"; +import { AgentRunError } from "../common/errors.js"; +import type { JsonRecord } from "../common/types.js"; +import { redactText } from "../common/redaction.js"; +import { stableHash } from "../common/validation.js"; +import { throwIfAborted, waitForChildProcess } from "./abortable-child.js"; + +const serviceAccountTokenFile = "/var/run/secrets/kubernetes.io/serviceaccount/token"; +const serviceAccountCaFile = "/var/run/secrets/kubernetes.io/serviceaccount/ca.crt"; +const maxCredentialBytes = 1024 * 1024; +const maxResponseBytes = 16 * 1024 * 1024; + +export type KubernetesResource = "job" | "pod" | "event" | "secret" | "pvc"; + +interface KubernetesResourceDescriptor { + cliSingular: string; + cliPlural: string; + apiPrefix: string; + apiPlural: string; +} + +const resourceDescriptors: Record = { + job: { cliSingular: "job", cliPlural: "jobs", apiPrefix: "/apis/batch/v1", apiPlural: "jobs" }, + pod: { cliSingular: "pod", cliPlural: "pods", apiPrefix: "/api/v1", apiPlural: "pods" }, + event: { cliSingular: "event", cliPlural: "events", apiPrefix: "/api/v1", apiPlural: "events" }, + secret: { cliSingular: "secret", cliPlural: "secrets", apiPrefix: "/api/v1", apiPlural: "secrets" }, + pvc: { cliSingular: "pvc", cliPlural: "pvc", apiPrefix: "/api/v1", apiPlural: "persistentvolumeclaims" }, +}; + +export interface KubernetesTransportResult { + code: number | null; + signal: NodeJS.Signals | null; + stdout: string; + stderr: string; + statusCode: number | null; + kubernetesReason: string | null; + transport: "in-cluster-https" | "explicit-kubectl"; + deletedCount?: number; +} + +export interface KubernetesInClusterConfig { + host: string; + port: number; + token: string; + ca: Buffer; +} + +interface KubernetesNativeRequestInput { + config: KubernetesInClusterConfig; + method: "GET" | "POST" | "DELETE" | "PATCH"; + path: string; + body?: string; + contentType?: string; + signal?: AbortSignal; +} + +interface KubernetesNativeResponse { + statusCode: number; + headers: IncomingHttpHeaders; + body: string; +} + +export interface KubernetesTransportTestHooks { + loadConfig?: () => Promise; + request?: (input: KubernetesNativeRequestInput) => Promise; +} + +export interface KubernetesClientOptions { + namespace: string; + kubectlCommand?: string; + signal?: AbortSignal; + testHooks?: KubernetesTransportTestHooks; +} + +export interface KubernetesListOptions { + labelSelector?: string; + fieldSelector?: string; +} + +export interface KubernetesDeleteOptions { + ignoreNotFound?: boolean; + propagationPolicy?: "Background" | "Foreground" | "Orphan"; + preconditions?: { uid: string; resourceVersion: string }; +} + +export interface KubernetesClient { + list(resource: KubernetesResource, options?: KubernetesListOptions): Promise; + get(resource: KubernetesResource, name: string): Promise; + create(manifest: JsonRecord, options?: { sensitiveBody?: boolean }): Promise; + delete(resource: KubernetesResource, name: string, options?: KubernetesDeleteOptions): Promise; + deleteCollection(resource: KubernetesResource, options: KubernetesListOptions & KubernetesDeleteOptions): Promise; + patch(resource: KubernetesResource, name: string, patch: JsonRecord): Promise; +} + +export function createKubernetesClient(options: KubernetesClientOptions): KubernetesClient { + const namespace = requiredIdentity(options.namespace, "namespace"); + return { + list: async (resource, listOptions = {}) => await execute(options, { + operation: `list-${resource}`, + method: "GET", + path: resourcePath(resource, namespace), + query: selectorQuery(listOptions), + kubectl: { + args: ["get", resourceDescriptors[resource].cliPlural, "-n", namespace, ...selectorArgs(listOptions), "-o", "json"], + }, + }), + get: async (resource, name) => { + const resourceName = requiredIdentity(name, "resource name"); + return await execute(options, { + operation: `get-${resource}`, + method: "GET", + path: resourcePath(resource, namespace, resourceName), + kubectl: { args: ["get", resourceDescriptors[resource].cliSingular, resourceName, "-n", namespace, "-o", "json"] }, + }); + }, + create: async (manifest, createOptions = {}) => { + const resource = resourceForManifest(manifest); + const manifestNamespace = manifestNamespaceFor(manifest); + if (manifestNamespace !== namespace) throw transportSchemaError("manifest namespace does not match Kubernetes client namespace", { namespace, manifestNamespace, resource }); + const stdin = `${JSON.stringify(manifest)}\n`; + return await execute(options, { + operation: `create-${resource}`, + method: "POST", + path: resourcePath(resource, namespace), + body: stdin, + contentType: "application/json", + sensitiveBody: createOptions.sensitiveBody === true, + kubectl: { args: ["create", "-f", "-", "-o", "json"], stdin }, + }); + }, + delete: async (resource, name, deleteOptions = {}) => { + const resourceName = requiredIdentity(name, "resource name"); + const deleteBody = kubernetesDeleteBody(deleteOptions); + const useRawDelete = deleteOptions.preconditions !== undefined; + const path = resourcePath(resource, namespace, resourceName); + const result = await execute(options, { + operation: `delete-${resource}`, + method: "DELETE", + path, + body: deleteBody, + contentType: "application/json", + kubectl: useRawDelete + ? { args: ["delete", "--raw", path, "-f", "-"], stdin: deleteBody } + : { args: ["delete", resourceDescriptors[resource].cliSingular, resourceName, "-n", namespace, ...(deleteOptions.ignoreNotFound === true ? ["--ignore-not-found=true"] : [])] }, + }); + return ignoreNotFoundResult(result, deleteOptions.ignoreNotFound === true); + }, + deleteCollection: async (resource, deleteOptions) => { + const query = selectorQuery(deleteOptions); + const kubectlArgs = ["delete", resourceDescriptors[resource].cliSingular, "-n", namespace, ...selectorArgs(deleteOptions), ...(deleteOptions.ignoreNotFound === true ? ["--ignore-not-found=true"] : [])]; + if (options.kubectlCommand) { + const result = await execute(options, { + operation: `delete-collection-${resource}`, + method: "DELETE", + path: resourcePath(resource, namespace), + query, + body: kubernetesDeleteBody(deleteOptions), + contentType: "application/json", + kubectl: { args: kubectlArgs }, + }); + return ignoreNotFoundResult(result, deleteOptions.ignoreNotFound === true); + } + const observed = await execute(options, { + operation: `list-before-delete-collection-${resource}`, + method: "GET", + path: resourcePath(resource, namespace), + query, + kubectl: { args: ["get", resourceDescriptors[resource].cliPlural, "-n", namespace, ...selectorArgs(deleteOptions), "-o", "json"] }, + }); + if (observed.code !== 0) return observed; + const deletedCount = listItemCount(observed.stdout); + const result = await execute(options, { + operation: `delete-collection-${resource}`, + method: "DELETE", + path: resourcePath(resource, namespace), + query, + body: kubernetesDeleteBody(deleteOptions), + contentType: "application/json", + kubectl: { args: kubectlArgs }, + }); + const normalized = ignoreNotFoundResult(result, deleteOptions.ignoreNotFound === true); + return normalized.code === 0 ? { ...normalized, deletedCount } : normalized; + }, + patch: async (resource, name, patch) => { + const resourceName = requiredIdentity(name, "resource name"); + const body = JSON.stringify(patch); + return await execute(options, { + operation: `patch-${resource}`, + method: "PATCH", + path: resourcePath(resource, namespace, resourceName), + body, + contentType: "application/merge-patch+json", + kubectl: { args: ["patch", resourceDescriptors[resource].cliSingular, resourceName, "-n", namespace, "--type", "merge", "-p", body, "-o", "json"] }, + }); + }, + }; +} + +export async function loadInClusterKubernetesConfig(input: { env?: NodeJS.ProcessEnv; tokenFile?: string; caFile?: string } = {}): Promise { + const env = input.env ?? process.env; + const host = env.KUBERNETES_SERVICE_HOST?.trim() ?? ""; + const rawPort = env.KUBERNETES_SERVICE_PORT_HTTPS?.trim() || env.KUBERNETES_SERVICE_PORT?.trim() || ""; + const port = Number(rawPort); + const missing: string[] = []; + if (!host) missing.push("KUBERNETES_SERVICE_HOST"); + if (!rawPort) missing.push("KUBERNETES_SERVICE_PORT_HTTPS|KUBERNETES_SERVICE_PORT"); + if (missing.length > 0 || !Number.isInteger(port) || port < 1 || port > 65_535) { + throw new AgentRunError("infra-failed", "Kubernetes in-cluster service endpoint is unavailable", { + httpStatus: 503, + details: { reason: "kubernetes-in-cluster-endpoint-unavailable", missing, portValid: Number.isInteger(port) && port >= 1 && port <= 65_535, valuesPrinted: false }, + }); + } + const tokenPath = input.tokenFile ?? serviceAccountTokenFile; + const caPath = input.caFile ?? serviceAccountCaFile; + let token: string; + let ca: Buffer; + try { + [token, ca] = await Promise.all([readFile(tokenPath, "utf8"), readFile(caPath)]); + } catch (error) { + throw new AgentRunError("infra-failed", "Kubernetes in-cluster ServiceAccount credentials are unavailable", { + httpStatus: 503, + details: { reason: "kubernetes-service-account-credential-unavailable", errorClass: error instanceof Error ? error.name : "UnknownError", valuesPrinted: false }, + }); + } + token = token.trim(); + if (token.length === 0 || Buffer.byteLength(token, "utf8") > maxCredentialBytes || ca.length === 0 || ca.length > maxCredentialBytes) { + throw new AgentRunError("infra-failed", "Kubernetes in-cluster ServiceAccount credentials are invalid", { + httpStatus: 503, + details: { reason: "kubernetes-service-account-credential-invalid", tokenPresent: token.length > 0, tokenSizeValid: Buffer.byteLength(token, "utf8") <= maxCredentialBytes, caPresent: ca.length > 0, caSizeValid: ca.length <= maxCredentialBytes, valuesPrinted: false }, + }); + } + return { host: stripIpv6Brackets(host), port, token, ca }; +} + +interface ExecuteInput { + operation: string; + method: KubernetesNativeRequestInput["method"]; + path: string; + query?: URLSearchParams; + body?: string; + contentType?: string; + sensitiveBody?: boolean; + kubectl: { args: string[]; stdin?: string }; +} + +async function execute(options: KubernetesClientOptions, input: ExecuteInput): Promise { + throwIfAborted(options.signal, input.operation); + if (options.kubectlCommand) return await runExplicitKubectl(options.kubectlCommand, input.kubectl, options.signal); + const path = input.query && [...input.query.keys()].length > 0 ? `${input.path}?${input.query.toString()}` : input.path; + const loadConfig = options.testHooks?.loadConfig ?? loadInClusterKubernetesConfig; + const request = options.testHooks?.request ?? runNativeHttpsRequest; + let response: KubernetesNativeResponse; + try { + response = await request({ + config: await loadConfig(), + method: input.method, + path, + ...(input.body !== undefined ? { body: input.body } : {}), + ...(input.contentType !== undefined ? { contentType: input.contentType } : {}), + ...(options.signal ? { signal: options.signal } : {}), + }); + } catch (error) { + if (error instanceof AgentRunError) throw error; + const safeMessage = redactText(error instanceof Error ? error.message : String(error)).slice(0, 500); + throw new AgentRunError("infra-failed", "Kubernetes in-cluster HTTPS request failed", { + httpStatus: 503, + details: { reason: "kubernetes-api-request-failed", operation: input.operation, method: input.method, namespace: options.namespace, errorClass: error instanceof Error ? error.name : "UnknownError", messageHash: stableHash(safeMessage), transport: "in-cluster-https", valuesPrinted: false }, + }); + } + const ok = response.statusCode >= 200 && response.statusCode < 300; + const status = kubernetesStatusSummary(response.body, response.statusCode, input.sensitiveBody === true); + return { + code: ok ? 0 : 1, + signal: null, + stdout: ok ? response.body : "", + stderr: ok ? "" : status.text, + statusCode: response.statusCode, + kubernetesReason: status.reason, + transport: "in-cluster-https", + }; +} + +async function runExplicitKubectl(command: string, input: { args: string[]; stdin?: string }, signal?: AbortSignal): Promise { + throwIfAborted(signal, "explicit kubectl test adapter"); + const child = spawn(command, input.args, { stdio: [input.stdin === undefined ? "ignore" : "pipe", "pipe", "pipe"] }); + let stdout = ""; + let stderr = ""; + if (!child.stdout || !child.stderr) throw new AgentRunError("infra-failed", "explicit kubectl test adapter output pipes were not created", { httpStatus: 503 }); + child.stdout.setEncoding("utf8"); + child.stderr.setEncoding("utf8"); + child.stdout.on("data", (chunk) => { stdout += String(chunk); }); + child.stderr.on("data", (chunk) => { stderr += String(chunk); }); + if (input.stdin !== undefined) child.stdin?.end(input.stdin); + const result = await waitForChildProcess(child, { ...(signal ? { signal } : {}), label: "explicit kubectl test adapter" }); + return { ...result, stdout, stderr, statusCode: null, kubernetesReason: null, transport: "explicit-kubectl" }; +} + +async function runNativeHttpsRequest(input: KubernetesNativeRequestInput): Promise { + throwIfAborted(input.signal, "Kubernetes in-cluster HTTPS request"); + const bodyBytes = input.body === undefined ? 0 : Buffer.byteLength(input.body, "utf8"); + const requestOptions: RequestOptions = { + protocol: "https:", + hostname: input.config.host, + port: input.config.port, + method: input.method, + path: input.path, + ca: input.config.ca, + headers: { + Accept: "application/json", + Authorization: `Bearer ${input.config.token}`, + ...(input.body !== undefined ? { "Content-Type": input.contentType ?? "application/json", "Content-Length": String(bodyBytes) } : {}), + }, + }; + return await new Promise((resolve, reject) => { + let settled = false; + const finish = (callback: () => void): void => { + if (settled) return; + settled = true; + input.signal?.removeEventListener("abort", onAbort); + callback(); + }; + const req = httpsRequest(requestOptions, (res) => { + const chunks: Buffer[] = []; + let receivedBytes = 0; + res.on("data", (chunk: Buffer | string) => { + const buffer = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk); + receivedBytes += buffer.length; + if (receivedBytes > maxResponseBytes) { + req.destroy(new Error("Kubernetes API response exceeded the bounded response limit")); + return; + } + chunks.push(buffer); + }); + res.on("end", () => finish(() => resolve({ statusCode: res.statusCode ?? 502, headers: res.headers, body: Buffer.concat(chunks).toString("utf8") }))); + res.on("error", (error) => finish(() => reject(error))); + }); + const onAbort = (): void => { + req.destroy(input.signal?.reason instanceof Error ? input.signal.reason : new Error("Kubernetes request aborted")); + }; + req.on("error", (error) => finish(() => reject(error))); + input.signal?.addEventListener("abort", onAbort, { once: true }); + if (input.body !== undefined) req.write(input.body); + req.end(); + }); +} + +function resourcePath(resource: KubernetesResource, namespace: string, name?: string): string { + const descriptor = resourceDescriptors[resource]; + const base = `${descriptor.apiPrefix}/namespaces/${encodeURIComponent(namespace)}/${descriptor.apiPlural}`; + return name ? `${base}/${encodeURIComponent(name)}` : base; +} + +function selectorQuery(options: KubernetesListOptions): URLSearchParams { + const query = new URLSearchParams(); + if (options.labelSelector) query.set("labelSelector", options.labelSelector); + if (options.fieldSelector) query.set("fieldSelector", options.fieldSelector); + return query; +} + +function selectorArgs(options: KubernetesListOptions): string[] { + return [ + ...(options.labelSelector ? ["-l", options.labelSelector] : []), + ...(options.fieldSelector ? ["--field-selector", options.fieldSelector] : []), + ]; +} + +function resourceForManifest(manifest: JsonRecord): KubernetesResource { + const apiVersion = typeof manifest.apiVersion === "string" ? manifest.apiVersion : ""; + const kind = typeof manifest.kind === "string" ? manifest.kind : ""; + if (apiVersion === "batch/v1" && kind === "Job") return "job"; + if (apiVersion === "v1" && kind === "Secret") return "secret"; + if (apiVersion === "v1" && kind === "PersistentVolumeClaim") return "pvc"; + throw transportSchemaError("manifest kind is not supported by the manager Kubernetes transport", { apiVersion, kind }); +} + +function manifestNamespaceFor(manifest: JsonRecord): string { + const metadata = typeof manifest.metadata === "object" && manifest.metadata !== null && !Array.isArray(manifest.metadata) ? manifest.metadata as JsonRecord : {}; + return requiredIdentity(metadata.namespace, "manifest.metadata.namespace"); +} + +function kubernetesDeleteBody(options: KubernetesDeleteOptions): string { + return JSON.stringify({ + apiVersion: "v1", + kind: "DeleteOptions", + ...(options.propagationPolicy ? { propagationPolicy: options.propagationPolicy } : {}), + ...(options.preconditions ? { preconditions: options.preconditions } : {}), + }); +} + +function listItemCount(stdout: string): number { + try { + const parsed = JSON.parse(stdout) as JsonRecord; + return Array.isArray(parsed.items) ? parsed.items.length : 0; + } catch { + return 0; + } +} + +function ignoreNotFoundResult(result: KubernetesTransportResult, ignoreNotFound: boolean): KubernetesTransportResult { + if (!ignoreNotFound || result.code === 0 || result.kubernetesReason !== "NotFound") return result; + return { ...result, code: 0, stderr: "" }; +} + +function kubernetesStatusSummary(body: string, statusCode: number, sensitiveBody: boolean): { reason: string | null; text: string } { + let reason: string | null = null; + let message: string | null = null; + try { + const parsed = JSON.parse(body) as JsonRecord; + reason = boundedString(parsed.reason, 120); + message = boundedString(parsed.message, 500); + } catch { + // The response body is intentionally not copied into an error envelope. + } + const safeMessage = message ? redactText(message) : null; + const messageHash = safeMessage ? stableHash(safeMessage) : body.length > 0 ? stableHash(body) : null; + const text = [ + `Kubernetes API HTTP ${statusCode}`, + `reason=${reason ?? "Unknown"}`, + ...(messageHash ? [`messageHash=${messageHash}`] : []), + ...(sensitiveBody ? ["requestBody=REDACTED"] : []), + ].join(" "); + return { reason, text }; +} + +function boundedString(value: unknown, maxLength: number): string | null { + return typeof value === "string" && value.length > 0 ? value.slice(0, maxLength) : null; +} + +function requiredIdentity(value: unknown, label: string): string { + if (typeof value !== "string" || value.trim().length === 0) throw transportSchemaError(`${label} is required`, { label }); + return value.trim(); +} + +function transportSchemaError(message: string, details: JsonRecord): AgentRunError { + return new AgentRunError("schema-invalid", message, { httpStatus: 500, details: { reason: "kubernetes-transport-schema-invalid", ...details, valuesPrinted: false } }); +} + +function stripIpv6Brackets(host: string): string { + return host.startsWith("[") && host.endsWith("]") ? host.slice(1, -1) : host; +} diff --git a/src/mgr/runner-retention.ts b/src/mgr/runner-retention.ts index 5b1cdcc..dd990bb 100644 --- a/src/mgr/runner-retention.ts +++ b/src/mgr/runner-retention.ts @@ -1,11 +1,10 @@ -import { spawn } from "node:child_process"; import { AgentRunError } from "../common/errors.js"; import { redactJson, redactText } from "../common/redaction.js"; import type { CommandRecord, JsonRecord, JsonValue, RunRecord } from "../common/types.js"; import { nowIso, stableHash } from "../common/validation.js"; import { isTerminalCommandState, isTerminalRunStatus } from "./store.js"; import type { AgentRunStore } from "./store.js"; -import { throwIfAborted, waitForChildProcess } from "./abortable-child.js"; +import { createKubernetesClient, type KubernetesClient, type KubernetesResource, type KubernetesTransportResult } from "./kubernetes-transport.js"; export interface RunnerRetentionOptions { namespace: string; @@ -347,8 +346,8 @@ export async function reconcileStalePendingRunners(input: { store: AgentRunStore async function runnerSnapshot(store: AgentRunStore, options: RunnerRetentionOptions): Promise { const observed = await Promise.allSettled([ - kubectlGetList(options, "jobs", labelSelector(options.matchLabels)), - kubectlGetList(options, "pods", labelSelector(options.matchLabels)), + kubernetesGetList(options, "jobs", labelSelector(options.matchLabels)), + kubernetesGetList(options, "pods", labelSelector(options.matchLabels)), ]); const failure = observed.find((item): item is PromiseRejectedResult => item.status === "rejected"); if (failure) throw failure.reason; @@ -364,7 +363,7 @@ async function runnerSnapshot(store: AgentRunStore, options: RunnerRetentionOpti let eventObjects: K8sObject[] = []; if (matchingPods.length > 0) { try { - eventObjects = await kubectlGetPodEvents(options); + eventObjects = await kubernetesGetPodEvents(options); } catch (error) { if (options.signal?.aborted) throw error; eventFactsAvailable = false; @@ -622,7 +621,7 @@ async function deleteRunnerResourceCas(options: RunnerRetentionOptions, item: Ru for (const pod of [...item.podObservations].sort((left, right) => left.name.localeCompare(right.name))) { await deleteKubernetesObjectCas(options, item, "pod", pod.name, pod.uid, pod.resourceVersion); } - const currentJob = await kubectlGetObject(options, "job", item.name); + const currentJob = await kubernetesGetObject(options, "job", item.name); const expectedIdentity = { "agentrun.pikastech.local/run-id": item.runId ?? "", "agentrun.pikastech.local/command-id": item.commandId ?? "", @@ -647,13 +646,9 @@ async function deleteKubernetesObjectCas(options: RunnerRetentionOptions, select details: { reason: "runner-retention-k8s-cas-unavailable", namespace: options.namespace, selectedResourceKind: selected.resourceKind, selectedName: selected.name, resourceKind, name, uidPresent: Boolean(uid), resourceVersionPresent: Boolean(resourceVersion), valuesPrinted: false }, }); } - const resourcePath = resourceKind === "job" - ? `/apis/batch/v1/namespaces/${encodeURIComponent(options.namespace)}/jobs/${encodeURIComponent(name)}` - : `/api/v1/namespaces/${encodeURIComponent(options.namespace)}/pods/${encodeURIComponent(name)}`; - const deleteOptions = JSON.stringify({ apiVersion: "v1", kind: "DeleteOptions", propagationPolicy: "Background", preconditions: { uid, resourceVersion } }); - const result = await kubectlRun(options.kubectlCommand ?? "kubectl", ["delete", "--raw", resourcePath, "-f", "-"], options.signal, deleteOptions); + const result = await retentionKubernetes(options).delete(resourceKind, name, { propagationPolicy: "Background", preconditions: { uid, resourceVersion } }); if (result.code === 0) return; - const conflict = /\b(?:Conflict|PreconditionFailed|resourceVersion|UID)\b/iu.test(`${result.stderr}\n${result.stdout}`); + const conflict = result.kubernetesReason === "Conflict" || /\b(?:Conflict|PreconditionFailed|resourceVersion|UID)\b/iu.test(`${result.stderr}\n${result.stdout}`); throw new AgentRunError("infra-failed", `Kubernetes CAS delete ${resourceKind} failed with code ${result.code}`, { httpStatus: conflict ? 409 : 502, details: redactJson({ @@ -668,35 +663,34 @@ async function deleteKubernetesObjectCas(options: RunnerRetentionOptions, select expectedResourceVersion: resourceVersion, exitCode: result.code, signal: result.signal, + transport: result.transport, + statusCode: result.statusCode, + kubernetesReason: result.kubernetesReason, stderr: redactText(result.stderr.slice(-2000)), - stdout: redactText(result.stdout.slice(-1000)), + stdoutHash: result.stdout ? stableHash(result.stdout) : null, valuesPrinted: false, }), }); } -async function kubectlGetList(options: RunnerRetentionOptions, resource: string, selector: string): Promise { - const args = ["get", resource, "-n", options.namespace]; - if (selector) args.push("-l", selector); - args.push("-o", "json"); - const result = await kubectlRun(options.kubectlCommand ?? "kubectl", args, options.signal); - if (result.code !== 0) throw new AgentRunError("infra-failed", `kubectl get ${resource} failed with code ${result.code}`, { httpStatus: 502, details: redactJson({ reason: "runner-kubernetes-observation-failed", namespace: options.namespace, resource, exitCode: result.code, signal: result.signal, stderr: redactText(result.stderr.slice(-2000)), stdout: redactText(result.stdout.slice(-1000)), valuesPrinted: false }) }); - const parsed = parseJsonObject(result.stdout, `kubectl get ${resource}`, "runner-kubernetes-observation-failed"); +async function kubernetesGetList(options: RunnerRetentionOptions, resource: string, selector: string): Promise { + const result = await retentionKubernetes(options).list(retentionResource(resource), selector ? { labelSelector: selector } : {}); + if (result.code !== 0) throw new AgentRunError("infra-failed", `Kubernetes get ${resource} failed`, { httpStatus: 502, details: retentionFailureDetails(result, { reason: "runner-kubernetes-observation-failed", namespace: options.namespace, resource }) }); + const parsed = parseJsonObject(result.stdout, `Kubernetes get ${resource}`, "runner-kubernetes-observation-failed"); const items = parsed.items; return Array.isArray(items) ? items.filter((entry) => typeof entry === "object" && entry !== null && !Array.isArray(entry)).map((entry) => entry as unknown as K8sObject) : []; } -async function kubectlGetObject(options: RunnerRetentionOptions, resource: "job" | "pod", name: string): Promise { - const result = await kubectlRun(options.kubectlCommand ?? "kubectl", ["get", resource, name, "-n", options.namespace, "-o", "json"], options.signal); - if (result.code !== 0) throw new AgentRunError("infra-failed", `kubectl get ${resource}/${name} failed with code ${result.code}`, { httpStatus: 409, details: redactJson({ reason: "runner-retention-k8s-cas-rejected", changeReason: "resource-missing-before-cas", namespace: options.namespace, resourceKind: resource, name, stderr: redactText(result.stderr.slice(-2000)), stdout: redactText(result.stdout.slice(-1000)), valuesPrinted: false }) }); - return parseJsonObject(result.stdout, `kubectl get ${resource}/${name}`, "runner-retention-k8s-cas-rejected") as unknown as K8sObject; +async function kubernetesGetObject(options: RunnerRetentionOptions, resource: "job" | "pod", name: string): Promise { + const result = await retentionKubernetes(options).get(resource, name); + if (result.code !== 0) throw new AgentRunError("infra-failed", `Kubernetes get ${resource}/${name} failed`, { httpStatus: 409, details: retentionFailureDetails(result, { reason: "runner-retention-k8s-cas-rejected", changeReason: "resource-missing-before-cas", namespace: options.namespace, resourceKind: resource, name }) }); + return parseJsonObject(result.stdout, `Kubernetes get ${resource}/${name}`, "runner-retention-k8s-cas-rejected") as unknown as K8sObject; } -async function kubectlGetPodEvents(options: RunnerRetentionOptions): Promise { - const args = ["get", "events", "-n", options.namespace, "--field-selector", "involvedObject.kind=Pod,type=Warning", "-o", "json"]; - const result = await kubectlRun(options.kubectlCommand ?? "kubectl", args, options.signal); - if (result.code !== 0) throw new AgentRunError("infra-failed", `kubectl get events failed with code ${result.code}`, { httpStatus: 502, details: redactJson({ reason: "runner-kubernetes-observation-failed", namespace: options.namespace, resource: "events", exitCode: result.code, signal: result.signal, stderr: redactText(result.stderr.slice(-2000)), stdout: redactText(result.stdout.slice(-1000)), valuesPrinted: false }) }); - const parsed = parseJsonObject(result.stdout, "kubectl get events", "runner-kubernetes-observation-failed"); +async function kubernetesGetPodEvents(options: RunnerRetentionOptions): Promise { + const result = await retentionKubernetes(options).list("event", { fieldSelector: "involvedObject.kind=Pod,type=Warning" }); + if (result.code !== 0) throw new AgentRunError("infra-failed", "Kubernetes get events failed", { httpStatus: 502, details: retentionFailureDetails(result, { reason: "runner-kubernetes-observation-failed", namespace: options.namespace, resource: "events" }) }); + const parsed = parseJsonObject(result.stdout, "Kubernetes get events", "runner-kubernetes-observation-failed"); const items = parsed.items; return Array.isArray(items) ? items.filter((entry) => typeof entry === "object" && entry !== null && !Array.isArray(entry)).map((entry) => entry as unknown as K8sObject) : []; } @@ -705,28 +699,35 @@ async function deleteAssociatedResources(options: RunnerRetentionOptions, jobNam let deleted = 0; const selector = `agentrun.pikastech.local/runner-job=${jobName}`; for (const resource of ["secret", "pvc"] as const) { - const result = await kubectlRun(options.kubectlCommand ?? "kubectl", ["delete", resource, "-n", options.namespace, "-l", selector, "--ignore-not-found=true"], options.signal); - if (result.code !== 0) throw new AgentRunError("infra-failed", `kubectl delete associated ${resource} failed with code ${result.code}`, { httpStatus: 502, details: redactJson({ stderr: redactText(result.stderr.slice(-2000)), stdout: redactText(result.stdout.slice(-1000)), valuesPrinted: false }) }); - deleted += countDeletedLines(result.stdout); + const result = await retentionKubernetes(options).deleteCollection(resource, { labelSelector: selector, ignoreNotFound: true, propagationPolicy: "Background" }); + if (result.code !== 0) throw new AgentRunError("infra-failed", `Kubernetes delete associated ${resource} failed`, { httpStatus: 502, details: retentionFailureDetails(result, { reason: "runner-associated-resource-delete-failed", namespace: options.namespace, resource, jobName }) }); + deleted += result.deletedCount ?? countDeletedLines(result.stdout); } return deleted; } -async function kubectlRun(kubectlCommand: string, args: string[], abortSignal?: AbortSignal, stdin?: string): Promise<{ code: number | null; signal: NodeJS.Signals | null; stdout: string; stderr: string }> { - throwIfAborted(abortSignal, "kubectl"); - const child = spawn(kubectlCommand, args, { stdio: [stdin === undefined ? "ignore" : "pipe", "pipe", "pipe"] }); - let stdout = ""; - let stderr = ""; - const stdoutStream = child.stdout; - const stderrStream = child.stderr; - if (!stdoutStream || !stderrStream) throw new AgentRunError("infra-failed", "kubectl output pipes were not created", { httpStatus: 503 }); - stdoutStream.setEncoding("utf8"); - stderrStream.setEncoding("utf8"); - stdoutStream.on("data", (chunk) => { stdout += String(chunk); }); - stderrStream.on("data", (chunk) => { stderr += String(chunk); }); - if (stdin !== undefined) child.stdin?.end(stdin); - const result = await waitForChildProcess(child, { ...(abortSignal ? { signal: abortSignal } : {}), label: "kubectl" }); - return { ...result, stdout, stderr }; +function retentionKubernetes(options: RunnerRetentionOptions): KubernetesClient { + return createKubernetesClient({ namespace: options.namespace, ...(options.kubectlCommand ? { kubectlCommand: options.kubectlCommand } : {}), ...(options.signal ? { signal: options.signal } : {}) }); +} + +function retentionResource(resource: string): KubernetesResource { + if (resource === "jobs" || resource === "job") return "job"; + if (resource === "pods" || resource === "pod") return "pod"; + throw new AgentRunError("schema-invalid", `runner retention resource ${resource} is unsupported`, { httpStatus: 500, details: { reason: "runner-retention-resource-unsupported", resource, valuesPrinted: false } }); +} + +function retentionFailureDetails(result: KubernetesTransportResult, context: JsonRecord): JsonRecord { + return redactJson({ + ...context, + exitCode: result.code, + signal: result.signal, + transport: result.transport, + statusCode: result.statusCode, + kubernetesReason: result.kubernetesReason, + stderr: redactText(result.stderr.slice(-2000)), + stdoutHash: result.stdout ? stableHash(result.stdout) : null, + valuesPrinted: false, + }); } function compareCandidates(left: RunnerResourceEntry, right: RunnerResourceEntry): number { diff --git a/src/mgr/session-pvc.ts b/src/mgr/session-pvc.ts index b224a2a..41684a4 100644 --- a/src/mgr/session-pvc.ts +++ b/src/mgr/session-pvc.ts @@ -1,10 +1,9 @@ -import { spawn } from "node:child_process"; import { createHash } from "node:crypto"; import type { AgentRunStore } from "./store.js"; import type { JsonRecord } from "../common/types.js"; import { AgentRunError } from "../common/errors.js"; import { redactJson, redactText } from "../common/redaction.js"; -import { throwIfAborted, waitForChildProcess } from "./abortable-child.js"; +import { createKubernetesClient, type KubernetesTransportResult } from "./kubernetes-transport.js"; export interface SessionPvcSpec { pvcName: string; @@ -80,27 +79,36 @@ export function buildSessionPvcSpec(input: { sessionId: string; namespace?: stri }; } -function resolveHandler(options: SessionPvcOptions): KubectlHandler { - if (options.kubectlHandler) return options.kubectlHandler; - const command = options.kubectlCommand ?? process.env.AGENTRUN_KUBECTL ?? "kubectl"; - return ({ args, stdin, signal }) => spawnKubectl(command, args, stdin, signal); +interface SessionPvcKubernetesResult { + stdout: string; + stderr: string; + exitCode: number; + statusCode?: number | null; + kubernetesReason?: string | null; + transport?: string; } -async function spawnKubectl(command: string, args: string[], stdinPayload?: string, signal?: AbortSignal): Promise<{ stdout: string; stderr: string; exitCode: number }> { - throwIfAborted(signal, "kubectl"); - const child = spawn(command, args, { stdio: ["pipe", "pipe", "pipe"] }); - let stdout = ""; - let stderr = ""; - child.stdout.setEncoding("utf8"); - child.stderr.setEncoding("utf8"); - child.stdout.on("data", (chunk) => { stdout += String(chunk); }); - child.stderr.on("data", (chunk) => { stderr += String(chunk); }); - const completion = waitForChildProcess(child, { ...(signal ? { signal } : {}), label: "kubectl" }); - child.stdin.on("error", () => {}); - if (stdinPayload !== undefined) child.stdin.end(`${stdinPayload}\n`); - else child.stdin.end(); - const result = await completion; - return { stdout, stderr, exitCode: result.code ?? 1 }; +interface SessionPvcKubernetes { + create(manifest: JsonRecord): Promise; + get(name: string): Promise; + delete(name: string): Promise; +} + +function resolveKubernetes(options: SessionPvcOptions, namespace: string): SessionPvcKubernetes { + if (options.kubectlHandler) { + const handler = options.kubectlHandler; + return { + create: async (manifest) => await handler({ args: ["create", "-f", "-", "-o", "json"], stdin: JSON.stringify(manifest), ...(options.signal ? { signal: options.signal } : {}) }), + get: async (name) => await handler({ args: ["get", "pvc", name, "-n", namespace, "-o", "json"], ...(options.signal ? { signal: options.signal } : {}) }), + delete: async (name) => await handler({ args: ["delete", "pvc", name, "-n", namespace, "--ignore-not-found"], ...(options.signal ? { signal: options.signal } : {}) }), + }; + } + const client = createKubernetesClient({ namespace, ...(options.kubectlCommand ? { kubectlCommand: options.kubectlCommand } : {}), ...(options.signal ? { signal: options.signal } : {}) }); + return { + create: async (manifest) => pvcTransportResult(await client.create(manifest)), + get: async (name) => pvcTransportResult(await client.get("pvc", name)), + delete: async (name) => pvcTransportResult(await client.delete("pvc", name, { ignoreNotFound: true })), + }; } export async function createSessionPvc(input: { store: AgentRunStore; sessionId: string; namespace?: string; options: SessionPvcOptions }): Promise { @@ -118,10 +126,9 @@ export async function createSessionPvc(input: { store: AgentRunStore; sessionId: }, spec: { accessModes: ["ReadWriteOnce"], storageClassName: spec.storageClassName, resources: { requests: { storage: spec.size } } }, }; - const handler = resolveHandler(input.options); - const result = await handler({ args: ["create", "-f", "-", "-o", "json"], stdin: JSON.stringify(manifest), ...(input.options.signal ? { signal: input.options.signal } : {}) }); + const result = await resolveKubernetes(input.options, spec.namespace).create(manifest); if (result.exitCode !== 0) { - throw new AgentRunError("infra-failed", `kubectl create session PVC failed with code ${result.exitCode}`, { httpStatus: 502, details: redactJson({ stderr: redactText(result.stderr.slice(-4000)) }) }); + throw new AgentRunError("infra-failed", "Kubernetes create session PVC failed", { httpStatus: 502, details: pvcFailureDetails(result, "session-pvc-create-failed") }); } let phase = "Pending"; try { @@ -157,8 +164,7 @@ export async function getSessionPvcSummary(input: { store: AgentRunStore; sessio if (!session) throw new AgentRunError("schema-invalid", `session ${input.sessionId} was not found`, { httpStatus: 404 }); const spec = buildSessionPvcSpec(input); const codexRolloutSubdir = session.codexRolloutSubdir ?? defaultSubdir; - const handler = resolveHandler(input.options); - const result = await handler({ args: ["get", "pvc", spec.pvcName, "-n", spec.namespace, "-o", "json"], ...(input.options.signal ? { signal: input.options.signal } : {}) }); + const result = await resolveKubernetes(input.options, spec.namespace).get(spec.pvcName); if (result.exitCode !== 0) { const notFound = `${result.stdout}\n${result.stderr}`.toLowerCase().includes("notfound") || `${result.stdout}\n${result.stderr}`.toLowerCase().includes("not found"); if (notFound) { @@ -174,7 +180,7 @@ export async function getSessionPvcSummary(input: { store: AgentRunStore; sessio valuesPrinted: false, }; } - throw new AgentRunError("infra-failed", `kubectl get pvc failed with code ${result.exitCode}`, { httpStatus: 502, details: redactJson({ stderr: redactText(result.stderr.slice(-4000)) }) }); + throw new AgentRunError("infra-failed", "Kubernetes get session PVC failed", { httpStatus: 502, details: pvcFailureDetails(result, "session-pvc-observation-failed") }); } let pvc: JsonRecord = {}; try { pvc = JSON.parse(result.stdout) as JsonRecord; } catch { pvc = {}; } @@ -197,8 +203,8 @@ export async function deleteSessionPvc(input: { store: AgentRunStore; sessionId: const session = await input.store.getSession(input.sessionId); if (!session) throw new AgentRunError("schema-invalid", `session ${input.sessionId} was not found`, { httpStatus: 404 }); const spec = buildSessionPvcSpec(input); - const handler = resolveHandler(input.options); - await handler({ args: ["delete", "pvc", spec.pvcName, "-n", spec.namespace, "--ignore-not-found"], ...(input.options.signal ? { signal: input.options.signal } : {}) }); + const result = await resolveKubernetes(input.options, spec.namespace).delete(spec.pvcName); + if (result.exitCode !== 0) throw new AgentRunError("infra-failed", "Kubernetes delete session PVC failed", { httpStatus: 502, details: pvcFailureDetails(result, "session-pvc-delete-failed") }); await input.store.markSessionStorageEvicted({ sessionId: input.sessionId, pvcName: spec.pvcName }); return { pvcName: spec.pvcName, namespace: spec.namespace, storageKind: "evicted" }; } @@ -244,7 +250,6 @@ export async function runSessionStorageGc(input: { store: AgentRunStore; options const now = input.now ?? Date.now(); const max = input.maxSessions ?? 200; const sessions = await input.store.listGcExpiredSessions({ now, limit: max }); - const handler = resolveHandler(input.options); let deleted = 0; let skipped = 0; const deletedPvcNames: string[] = []; @@ -252,7 +257,8 @@ export async function runSessionStorageGc(input: { store: AgentRunStore; options if (session.activeRunId || session.activeCommandId) { skipped++; continue; } if (session.storageKind !== "pvc" || !session.storagePvcName) { skipped++; continue; } try { - await handler({ args: ["delete", "pvc", session.storagePvcName, "-n", session.storageNamespace ?? runtimeNamespace(), "--ignore-not-found"], ...(input.options.signal ? { signal: input.options.signal } : {}) }); + const result = await resolveKubernetes(input.options, session.storageNamespace ?? runtimeNamespace()).delete(session.storagePvcName); + if (result.exitCode !== 0) throw new AgentRunError("infra-failed", "Kubernetes delete expired session PVC failed", { httpStatus: 502, details: pvcFailureDetails(result, "session-pvc-gc-delete-failed") }); await input.store.markSessionStorageEvicted({ sessionId: session.sessionId, pvcName: session.storagePvcName }); deleted++; deletedPvcNames.push(session.storagePvcName); @@ -305,3 +311,26 @@ function parseStorageSize(value: string): number | null { const factor = multipliers[unit] ?? 1; return Math.round(num * factor); } + +function pvcTransportResult(result: KubernetesTransportResult): SessionPvcKubernetesResult { + return { + stdout: result.stdout, + stderr: result.stderr, + exitCode: result.code ?? 1, + statusCode: result.statusCode, + kubernetesReason: result.kubernetesReason, + transport: result.transport, + }; +} + +function pvcFailureDetails(result: SessionPvcKubernetesResult, reason: string): JsonRecord { + return redactJson({ + reason, + exitCode: result.exitCode, + transport: result.transport ?? "explicit-test-handler", + statusCode: result.statusCode ?? null, + kubernetesReason: result.kubernetesReason ?? null, + stderr: redactText(result.stderr.slice(-2000)), + valuesPrinted: false, + }); +} diff --git a/src/selftest/cases/10-mgr-session-pvc.ts b/src/selftest/cases/10-mgr-session-pvc.ts index 35bd55f..7396cf0 100644 --- a/src/selftest/cases/10-mgr-session-pvc.ts +++ b/src/selftest/cases/10-mgr-session-pvc.ts @@ -88,7 +88,22 @@ const selfTest: SelfTestCase = async () => { assert.equal(activeAfter?.storageKind, "pvc"); const restStore = new MemoryAgentRunStore(); - const server = await startManagerServer({ port: 0, host: "127.0.0.1", sourceCommit: "self-test", store: restStore, sessionPvcOptions: { kubectlHandler: fakeKubectl } }); + const server = await startManagerServer({ + port: 0, + host: "127.0.0.1", + sourceCommit: "self-test", + store: restStore, + sessionPvcOptions: { kubectlHandler: fakeKubectl }, + runnerJobDefaults: { + namespace: "agentrun-v02", + managerUrl: "http://agentrun-mgr.agentrun-v02.svc.cluster.local:8080", + runnerApiKeySecretRef: { name: "agentrun-v02-api-key", key: "HWLAB_API_KEY" }, + image: "127.0.0.1:5000/agentrun/agentrun-mgr@sha256:1111111111111111111111111111111111111111111111111111111111111111", + serviceAccountName: "agentrun-v02-runner", + jobNamePrefix: "agentrun-v02-runner", + lane: "v0.2", + }, + }); try { const client = new ManagerClient(server.baseUrl); const create = await client.post("/api/v1/sessions", { sessionId: "sess_rest_create_001", tenantId: "hwlab", projectId: "pikasTech/HWLAB", backendProfile: "codex" }) as { action: string; pvc: { pvcName: string; pvcPhase: string } }; diff --git a/src/selftest/cases/11-kubernetes-transport.ts b/src/selftest/cases/11-kubernetes-transport.ts new file mode 100644 index 0000000..ea3aeb9 --- /dev/null +++ b/src/selftest/cases/11-kubernetes-transport.ts @@ -0,0 +1,144 @@ +import assert from "node:assert/strict"; +import path from "node:path"; +import { writeFile } from "node:fs/promises"; +import { AgentRunError } from "../../common/errors.js"; +import { createKubernetesClient, loadInClusterKubernetesConfig, type KubernetesInClusterConfig } from "../../mgr/kubernetes-transport.js"; +import type { SelfTestCase } from "../harness.js"; + +const selfTest: SelfTestCase = async (context) => { + const tokenFile = path.join(context.tmp, "kubernetes-service-account-token"); + const caFile = path.join(context.tmp, "kubernetes-service-account-ca.crt"); + const token = "service-account-token-must-never-be-rendered"; + const ca = "-----BEGIN CERTIFICATE-----\nselftest-ca-material\n-----END CERTIFICATE-----\n"; + await Promise.all([writeFile(tokenFile, `${token}\n`, "utf8"), writeFile(caFile, ca, "utf8")]); + + const config = await loadInClusterKubernetesConfig({ + env: { KUBERNETES_SERVICE_HOST: "[fd00::10]", KUBERNETES_SERVICE_PORT_HTTPS: "443" }, + tokenFile, + caFile, + }); + assert.equal(config.host, "fd00::10"); + assert.equal(config.port, 443); + assert.equal(config.token, token); + assert.equal(config.ca.toString("utf8"), ca); + + await assert.rejects( + () => loadInClusterKubernetesConfig({ env: {}, tokenFile, caFile }), + (error) => error instanceof AgentRunError + && error.details?.reason === "kubernetes-in-cluster-endpoint-unavailable" + && !JSON.stringify(error).includes(token), + ); + + const requests: Array<{ method: string; path: string; body: string | null; contentType: string | null }> = []; + const client = createKubernetesClient({ + namespace: "agentrun-v02", + testHooks: { + loadConfig: async () => config, + request: async (input) => { + assert.equal(input.config.token, token); + assert.equal(input.config.ca.toString("utf8"), ca); + requests.push({ method: input.method, path: input.path, body: input.body ?? null, contentType: input.contentType ?? null }); + if (input.method === "POST" && input.path.endsWith("/secrets")) { + return { statusCode: 409, headers: {}, body: JSON.stringify({ kind: "Status", reason: "AlreadyExists", message: `Authorization: Bearer ${token}` }) }; + } + if (input.method === "GET" && input.path.includes("persistentvolumeclaims?")) { + return { statusCode: 200, headers: {}, body: JSON.stringify({ items: [{ metadata: { name: "pvc-a" } }, { metadata: { name: "pvc-b" } }] }) }; + } + if (input.method === "GET" && input.path.includes("/jobs?")) { + return { statusCode: 200, headers: {}, body: JSON.stringify({ items: [{ metadata: { name: "runner-a" } }] }) }; + } + if (input.method === "DELETE" && input.path.endsWith("/secrets/missing")) { + return { statusCode: 404, headers: {}, body: JSON.stringify({ kind: "Status", reason: "NotFound", message: "Secret missing was not found" }) }; + } + return { statusCode: input.method === "POST" ? 201 : input.method === "DELETE" ? 200 : 200, headers: {}, body: JSON.stringify({ apiVersion: "v1", kind: "Status", metadata: { name: "bounded" } }) }; + }, + }, + }); + + const listed = await client.list("job", { labelSelector: "app.kubernetes.io/name=agentrun-runner" }); + assert.equal(listed.code, 0); + assert.equal(listed.transport, "in-cluster-https"); + assert.match(requests.at(-1)?.path ?? "", /^\/apis\/batch\/v1\/namespaces\/agentrun-v02\/jobs\?labelSelector=/u); + + const createdJob = await client.create({ apiVersion: "batch/v1", kind: "Job", metadata: { name: "runner-a", namespace: "agentrun-v02" }, spec: {} }); + assert.equal(createdJob.code, 0); + assert.equal(requests.at(-1)?.method, "POST"); + assert.equal(requests.at(-1)?.path, "/apis/batch/v1/namespaces/agentrun-v02/jobs"); + + const createdPvc = await client.create({ apiVersion: "v1", kind: "PersistentVolumeClaim", metadata: { name: "session-a", namespace: "agentrun-v02" }, spec: {} }); + assert.equal(createdPvc.code, 0); + assert.equal(requests.at(-1)?.path, "/api/v1/namespaces/agentrun-v02/persistentvolumeclaims"); + const observedPvc = await client.get("pvc", "session-a"); + assert.equal(observedPvc.code, 0); + assert.equal(requests.at(-1)?.path, "/api/v1/namespaces/agentrun-v02/persistentvolumeclaims/session-a"); + const deletedPvc = await client.delete("pvc", "session-a", { ignoreNotFound: true }); + assert.equal(deletedPvc.code, 0); + assert.equal(requests.at(-1)?.method, "DELETE"); + assert.equal(requests.at(-1)?.path, "/api/v1/namespaces/agentrun-v02/persistentvolumeclaims/session-a"); + + const secret = await client.create({ apiVersion: "v1", kind: "Secret", metadata: { name: "runner-env", namespace: "agentrun-v02" }, stringData: { PRIVATE_VALUE: token } }, { sensitiveBody: true }); + assert.equal(secret.code, 1); + assert.equal(secret.kubernetesReason, "AlreadyExists"); + assert.equal(secret.stderr.includes(token), false); + assert.match(secret.stderr, /requestBody=REDACTED/u); + const ignoredMissing = await client.delete("secret", "missing", { ignoreNotFound: true }); + assert.equal(ignoredMissing.code, 0); + assert.equal(ignoredMissing.stderr, ""); + const rejectedMissing = await client.delete("secret", "missing"); + assert.equal(rejectedMissing.code, 1); + assert.equal(rejectedMissing.kubernetesReason, "NotFound"); + + const casDelete = await client.delete("job", "runner-a", { propagationPolicy: "Background", preconditions: { uid: "uid-a", resourceVersion: "17" } }); + assert.equal(casDelete.code, 0); + assert.equal(requests.at(-1)?.method, "DELETE"); + assert.match(requests.at(-1)?.body ?? "", /"preconditions":\{"uid":"uid-a","resourceVersion":"17"\}/u); + + const patched = await client.patch("secret", "runner-env", { metadata: { ownerReferences: [{ uid: "uid-a" }] } }); + assert.equal(patched.code, 0); + assert.equal(requests.at(-1)?.method, "PATCH"); + assert.equal(requests.at(-1)?.contentType, "application/merge-patch+json"); + + const deleted = await client.deleteCollection("pvc", { labelSelector: "agentrun.pikastech.local/runner-job=runner-a", ignoreNotFound: true, propagationPolicy: "Background" }); + assert.equal(deleted.code, 0); + assert.equal(deleted.deletedCount, 2); + assert.equal(requests.at(-1)?.method, "DELETE"); + assert.match(requests.at(-1)?.path ?? "", /^\/api\/v1\/namespaces\/agentrun-v02\/persistentvolumeclaims\?labelSelector=/u); + + const failingClient = createKubernetesClient({ + namespace: "agentrun-v02", + testHooks: { + loadConfig: async (): Promise => config, + request: async () => { throw new Error(`Authorization: Bearer ${token}`); }, + }, + }); + const requestError = await captureAgentRunError(() => failingClient.get("pod", "runner-a")); + assert.equal(requestError.details?.reason, "kubernetes-api-request-failed"); + assert.equal(JSON.stringify(requestError).includes(token), false); + assert.equal(JSON.stringify(requestError.details).includes(ca), false); + + return { + name: "kubernetes-transport", + tests: [ + "in-cluster-service-account-config", + "native-resource-paths-and-selectors", + "native-job-and-pvc-crud-paths", + "secret-errors-bounded-and-redacted", + "ignore-not-found-delete-semantics", + "cas-delete-and-owner-patch", + "associated-resource-delete-collection", + "transport-failure-redaction", + ], + }; +}; + +async function captureAgentRunError(run: () => Promise): Promise { + try { + await run(); + } catch (error) { + if (error instanceof AgentRunError) return error; + throw error; + } + throw new Error("expected AgentRunError"); +} + +export default selfTest; diff --git a/src/selftest/cases/20-runner-k8s-job.ts b/src/selftest/cases/20-runner-k8s-job.ts index 8a0116d..0134d03 100644 --- a/src/selftest/cases/20-runner-k8s-job.ts +++ b/src/selftest/cases/20-runner-k8s-job.ts @@ -12,7 +12,7 @@ import { assertNoSecretLeak, createRunWithCommand, loadArtificerImageRef, type S const runnerApiKeySecretRef = { name: "agentrun-selftest-api-key", key: "HWLAB_API_KEY" }; const selfTest: SelfTestCase = async (context) => { - const server = await startManagerServer({ port: 0, host: "127.0.0.1", sourceCommit: "self-test", store: new MemoryAgentRunStore() }); + const server = await startManagerServer({ port: 0, host: "127.0.0.1", sourceCommit: "self-test", store: new MemoryAgentRunStore(), runnerJobDefaults: baselineRunnerDefaults() }); try { const client = new ManagerClient(server.baseUrl); const artificerImageRef = await loadArtificerImageRef(context.root); @@ -216,6 +216,9 @@ process.exit(1); runnerApiKeySecretRef, image: "127.0.0.1:5000/agentrun/agentrun-mgr@sha256:1111111111111111111111111111111111111111111111111111111111111111", envIdentity: "selftest-env-identity", + serviceAccountName: "agentrun-v01-runner", + jobNamePrefix: "agentrun-v01-runner", + lane: "v0.1", kubectlCommand: fakeKubectl, unideskSshEndpointEnv: { name: "UNIDESK_MAIN_SERVER_IP", value: "https://unidesk.default.example.test" }, }, @@ -449,7 +452,7 @@ process.exit(1); await new Promise((resolve) => serverWithRetention.server.close(() => resolve())); } const protectedStore = new MemoryAgentRunStore(); - const protectedServer = await startManagerServer({ port: 0, host: "127.0.0.1", sourceCommit: "self-test", store: protectedStore }); + const protectedServer = await startManagerServer({ port: 0, host: "127.0.0.1", sourceCommit: "self-test", store: protectedStore, runnerJobDefaults: baselineRunnerDefaults() }); const protectedKubectl = path.join(context.tmp, "fake-kubectl-retention-protected.js"); const protectedDeleted = path.join(context.tmp, "retention-protected-deleted.json"); await writeFile(protectedDeleted, "[]\n"); @@ -776,3 +779,15 @@ function assertRunnerJobDoesNotMountProfile(manifest: JsonRecord, volumeName: st assert.equal(volumes.some((volume) => volume.name === volumeName), false, `${volumeName} volume must not be mounted for another backendProfile`); assert.equal(mounts.some((mount) => mount.name === volumeName), false, `${volumeName} mount must not exist for another backendProfile`); } + +function baselineRunnerDefaults(): NonNullable[0]>["runnerJobDefaults"]> { + return { + namespace: "agentrun-v02", + managerUrl: "http://agentrun-mgr.agentrun-v02.svc.cluster.local:8080", + runnerApiKeySecretRef, + image: "127.0.0.1:5000/agentrun/agentrun-mgr@sha256:1111111111111111111111111111111111111111111111111111111111111111", + serviceAccountName: "agentrun-v02-runner", + jobNamePrefix: "agentrun-v02-runner", + lane: "v0.2", + }; +}