import { BatchV1Api, CoreV1Api, KubernetesObjectApi, type V1DeleteOptions, type V1Job, } from "@kubernetes/client-node"; import { createHash, randomUUID } from "node:crypto"; import { mkdir, open, readFile, rm, writeFile } from "node:fs/promises"; import { join } from "node:path"; import type { Sha256Digest } from "../shared/build-protocol"; import type { BuildJobObservation, BuildKubernetesOperations, } from "./build-controller"; import type { KubernetesJob } from "./build-job"; import { BuildStoreConflictError, buildTimestamp, compareBuildRecords, type BuildRecord, type BuildReconciliationLease, type BuildStore, type CreateBuildResult, type UploadRecord, } from "./build-store"; const TYPE_LABEL = "kuber.astrxl.dev/type"; const IMAGE_LABEL = "kuber.astrxl.dev/image"; type ConfigMap = { apiVersion: "v1"; kind: "ConfigMap"; metadata: { name: string; namespace: string; resourceVersion?: string; creationTimestamp?: string; labels?: Record; }; data?: Record; }; export interface BuildObjectApi { create(value: ConfigMap): Promise; read(value: ConfigMap): Promise; replace(value: ConfigMap): Promise; delete(value: ConfigMap, options?: V1DeleteOptions): Promise; list( apiVersion: string, kind: string, namespace?: string, pretty?: string, exact?: boolean, exportValue?: boolean, fieldSelector?: string, labelSelector?: string, ): Promise<{ items: unknown[] }>; } function hashName(prefix: string, value: string): string { return `${prefix}-${createHash("sha256").update(value).digest("hex").slice(0, 48)}`; } function statusCode(error: unknown): number | undefined { if (!error || typeof error !== "object") return; if ("code" in error && typeof error.code === "number") return error.code; if ("statusCode" in error && typeof error.statusCode === "number") return error.statusCode; } function payload(value: unknown): T | undefined { const object = value as ConfigMap; const raw = object.data?.payload; if (!raw) return; try { const parsed = JSON.parse(raw) as T & { metadata?: { resourceVersion?: string; creationTimestamp?: string }; }; if (parsed.metadata && object.metadata.resourceVersion) parsed.metadata.resourceVersion = object.metadata.resourceVersion; if (parsed.metadata && object.metadata.creationTimestamp) parsed.metadata.creationTimestamp = object.metadata.creationTimestamp; return parsed; } catch { return; } } function map( namespace: string, name: string, type: "build" | "build-upload" | "build-lock", value?: unknown, resourceVersion?: string, labels?: Record, ): ConfigMap { return { apiVersion: "v1", kind: "ConfigMap", metadata: { name, namespace, ...(resourceVersion && { resourceVersion }), labels: { [TYPE_LABEL]: type, ...labels }, }, ...(value !== undefined && { data: { payload: JSON.stringify(value) } }), }; } function terminal(record: BuildRecord): boolean { return ( record.status.state === "succeeded" || record.status.state === "failed" ); } function sameSpec(left: BuildRecord, right: BuildRecord): boolean { return JSON.stringify(left.spec) === JSON.stringify(right.spec); } function newer(left: BuildRecord, right: BuildRecord): boolean { return compareBuildRecords(left, right) > 0; } function supersede(record: BuildRecord, finishedAt: string): BuildRecord { const reason = "Superseded by newer build"; const next = structuredClone(record); next.metadata.resourceVersion = String( Number(record.metadata.resourceVersion) + 1, ); next.status = { ...record.status, state: "failed", finishedAt, error: reason, cancelled: true, events: [ ...record.status.events, { type: "status", status: { version: record.status.version, id: record.status.id, state: "failed", createdAt: record.status.createdAt, ...(record.status.startedAt && { startedAt: record.status.startedAt, }), finishedAt, error: reason, }, }, ], }; return next; } /** Build metadata lives in ConfigMaps; resumable upload bytes live only on the RWX volume. */ export class KubernetesBuildStore implements BuildStore { constructor( private readonly objects: BuildObjectApi, private readonly namespace: string, private readonly uploadRoot: string, ) {} private buildName(id: string): string { return hashName("build", id); } private uploadName(digest: Sha256Digest): string { return hashName("upload", digest); } private uploadPath(digest: Sha256Digest): string { return join(this.uploadRoot, digest.slice("sha256:".length)); } private lockName(imageKey: string): string { return hashName("build-lock", imageKey); } private imageLabel(imageKey: string): string { return createHash("sha256").update(imageKey).digest("hex").slice(0, 63); } private buildLabels(record: BuildRecord): Record { return { ...record.metadata.labels, [IMAGE_LABEL]: this.imageLabel(record.spec.imageKey), }; } private async listImageBuilds(imageKey: string): Promise { const result = await this.objects.list( "v1", "ConfigMap", this.namespace, undefined, undefined, undefined, undefined, `${TYPE_LABEL}=build,${IMAGE_LABEL}=${this.imageLabel(imageKey)}`, ); const records = result.items .map((item) => payload(item)) .filter( (item): item is BuildRecord => item?.kind === "BuildRecord" && item.spec?.imageKey === imageKey, ); // Legacy records lack the image index label. Keep this compatibility scan // bounded to build records, then match all identifying payload fields. const [project, service] = imageKey.split("\0"); if (!project || !service) return records; const legacy = await this.objects.list( "v1", "ConfigMap", this.namespace, undefined, undefined, undefined, undefined, `${TYPE_LABEL}=build`, ); const legacyRecords = legacy.items .map((item) => payload(item)) .filter( (item): item is BuildRecord => item?.kind === "BuildRecord" && item.spec?.request?.project === project && item.spec?.request?.service === service && item.spec?.imageKey === imageKey, ); return [ ...records, ...legacyRecords.filter( (legacyRecord) => !records.some( (record) => record.metadata.name === legacyRecord.metadata.name, ), ), ]; } private async read(value: ConfigMap): Promise { try { return payload(await this.objects.read(value)); } catch (error) { if (statusCode(error) === 404) return; throw error; } } private async readConfigMap( value: ConfigMap, ): Promise { try { return (await this.objects.read(value)) as ConfigMap; } catch (error) { if (statusCode(error) === 404) return; throw error; } } private async delete( value: ConfigMap, expectedResourceVersion?: string, ): Promise { try { await this.objects.delete( value, expectedResourceVersion ? { preconditions: { resourceVersion: expectedResourceVersion } } : undefined, ); } catch (error) { if (statusCode(error) !== 404 && statusCode(error) !== 409) throw error; } } private async supersedeBuild( record: BuildRecord, finishedAt: string, releaseLock = false, ): Promise { const next = supersede(record, finishedAt); await this.replaceBuildRecord( next, record.metadata.resourceVersion, releaseLock, ); return next; } /** * Build records are indexed by image label, rather than scanning all records. * This repairs record-first crashes and makes the record ordering authoritative * even when one of the records never acquired a lock. */ private async reconcileOlderBuilds( candidate: BuildRecord, ): Promise { const records = await this.listImageBuilds(candidate.spec.imageKey); const superseded: BuildRecord[] = []; for (const other of records) { if (other.metadata.name === candidate.metadata.name || terminal(other)) continue; if (newer(other, candidate)) { superseded.push( await this.supersedeBuild(candidate, buildTimestamp(candidate)), ); return superseded; } superseded.push( await this.supersedeBuild(other, buildTimestamp(candidate)), ); } return superseded; } async createBuild(record: BuildRecord): Promise { const existing = await this.getBuild(record.metadata.name); if (existing) { if (!sameSpec(existing, record)) throw new BuildStoreConflictError( "Build ID was already used for a different request", ); if (terminal(existing)) return { record: existing, created: false }; record = existing; } if (!existing) { // Persist before locking so contenders can find and supersede queued records. try { await this.objects.create( map( this.namespace, this.buildName(record.metadata.name), "build", record, undefined, this.buildLabels(record), ), ); } catch (error) { if (statusCode(error) !== 409) throw error; const concurrent = await this.getBuild(record.metadata.name); if (concurrent && sameSpec(concurrent, record)) return this.createBuild(concurrent); throw new BuildStoreConflictError( "Build ID was already used for a different request", ); } const stored = await this.getBuild(record.metadata.name); if (!stored) throw new BuildStoreConflictError("Build record disappeared"); record = stored; } const reconciled = await this.reconcileOlderBuilds(record); if ( reconciled.some((value) => value.metadata.name === record.metadata.name) ) return { record: (await this.getBuild(record.metadata.name))!, created: false, superseded: reconciled, }; const lockName = this.lockName(record.spec.imageKey); for (let attempt = 0; attempt < 8; attempt++) { try { await this.objects.create( map(this.namespace, lockName, "build-lock", { buildId: record.metadata.name, }), ); return { record: (await this.getBuild(record.metadata.name))!, created: true, ...(reconciled.length && { superseded: reconciled }), }; } catch (error) { if (statusCode(error) !== 409) { await this.failCreatedBuild(record, error); throw error; } const lockObject = await this.readConfigMap( map(this.namespace, lockName, "build-lock"), ); if (!lockObject) continue; const lock = payload<{ buildId: string }>(lockObject); if (!lock?.buildId) throw new BuildStoreConflictError("Invalid build lock"); if (lock.buildId === record.metadata.name) return { record: (await this.getBuild(record.metadata.name))!, created: false, ...(reconciled.length && { superseded: reconciled }), }; const active = await this.getBuild(lock.buildId); if (!active) { await this.delete(lockObject, lockObject.metadata.resourceVersion); continue; } if (terminal(active)) { await this.delete(lockObject, lockObject.metadata.resourceVersion); continue; } if (newer(active, record)) { const superseded = await this.supersedeBuild( record, buildTimestamp(record), ); return { record: superseded, created: false }; } try { const nextOld = await this.supersedeBuild( active, buildTimestamp(record), ); await this.objects.replace( map( this.namespace, lockName, "build-lock", { buildId: record.metadata.name }, lockObject.metadata.resourceVersion, ), ); return { record: (await this.getBuild(record.metadata.name))!, created: true, superseded: [...reconciled, nextOld], }; } catch (takeoverError) { if (statusCode(takeoverError) === 409) continue; await this.failCreatedBuild(record, takeoverError); throw takeoverError; } } } await this.failCreatedBuild( record, new BuildStoreConflictError("Build lock acquisition timed out"), ); throw new BuildStoreConflictError("Build lock acquisition timed out"); } private async failCreatedBuild( record: BuildRecord, error: unknown, ): Promise { const current = await this.getBuild(record.metadata.name); if (!current || terminal(current)) return; const next = structuredClone(current); next.metadata.resourceVersion = String( Number(current.metadata.resourceVersion) + 1, ); next.status = { ...current.status, state: "failed", finishedAt: new Date().toISOString(), error: error instanceof Error ? error.message : String(error), events: [ ...current.status.events, { type: "status", status: { version: current.status.version, id: current.status.id, state: "failed", createdAt: current.status.createdAt, finishedAt: next.status.finishedAt, error: next.status.error, }, }, ], }; await this.replaceBuildRecord(next, current.metadata.resourceVersion); } private async replaceBuildRecord( record: BuildRecord, expectedResourceVersion: string, releaseLock = true, ): Promise { const current = await this.getBuild(record.metadata.name); if ( !current || current.metadata.resourceVersion !== expectedResourceVersion || !sameSpec(current, record) ) throw new BuildStoreConflictError( "Build record was concurrently modified", ); try { await this.objects.replace( map( this.namespace, this.buildName(record.metadata.name), "build", record, expectedResourceVersion, this.buildLabels(record), ), ); } catch (error) { if (statusCode(error) === 409) throw new BuildStoreConflictError( "Build record was concurrently modified", ); throw error; } if (terminal(record) && releaseLock) await this.releaseLock(record.spec.imageKey, record.metadata.name); } async getBuild(id: string): Promise { const record = await this.read( map(this.namespace, this.buildName(id), "build"), ); return record?.metadata.name === id ? record : undefined; } async listBuilds(): Promise { const result = await this.objects.list( "v1", "ConfigMap", this.namespace, undefined, undefined, undefined, undefined, `${TYPE_LABEL}=build`, ); return result.items .map((item) => payload(item)) .filter((item): item is BuildRecord => item?.kind === "BuildRecord") .sort(compareBuildRecords); } /** Prune old terminal history using fresh reads and Kubernetes RV preconditions. */ async pruneBuilds(limit = 50): Promise { if (!Number.isSafeInteger(limit) || limit < 0) return; const records = await this.listBuilds(); let remaining = records.length; for (const candidate of records) { if (remaining <= limit) break; if (!terminal(candidate)) continue; const current = await this.readConfigMap( map(this.namespace, this.buildName(candidate.metadata.name), "build"), ); const record = current && payload(current); if ( !current?.metadata.resourceVersion || !record || record.kind !== "BuildRecord" || record.metadata.name !== candidate.metadata.name || current.metadata.resourceVersion !== candidate.metadata.resourceVersion || !terminal(record) ) continue; const expiresAt = record.status.reconcileLease?.expiresAt; if (expiresAt && Date.parse(expiresAt) > Date.now()) continue; // A build ID cannot reacquire an image lock after becoming terminal, so // an absent/mismatched lock cannot newly begin referencing this record. // Treat unreadable or malformed locks as protected rather than guessing. let lock: ConfigMap | undefined; try { lock = await this.readConfigMap( map(this.namespace, this.lockName(record.spec.imageKey), "build-lock"), ); } catch { continue; } if (lock) { const lockRecord = payload<{ buildId?: string }>(lock); if (!lockRecord?.buildId || lockRecord.buildId === record.metadata.name) continue; } await this.delete( map( this.namespace, current.metadata.name, "build", undefined, current.metadata.resourceVersion, ), current.metadata.resourceVersion, ); // Conflicts are swallowed by delete; confirm absence before adjusting the // retained count so concurrent updates never make us over-prune. if (!(await this.readConfigMap(current))) remaining--; } } async replaceBuild( record: BuildRecord, expectedResourceVersion: string, reconcileLeaseToken?: string, ): Promise { const current = await this.getBuild(record.metadata.name); if ( !current || current.metadata.resourceVersion !== expectedResourceVersion ) throw new BuildStoreConflictError( "Build record was concurrently modified", ); if (!sameSpec(current, record)) throw new BuildStoreConflictError("Build specification is immutable"); if ( reconcileLeaseToken && current.status.reconcileLease?.token !== reconcileLeaseToken ) { throw new BuildStoreConflictError( "Build reconciliation lease ownership was lost", ); } if ( current.status.state === "succeeded" && (record.status.state !== "succeeded" || current.status.digest !== record.status.digest) ) throw new BuildStoreConflictError( "A successful image digest is immutable", ); try { await this.objects.replace( map( this.namespace, this.buildName(record.metadata.name), "build", record, expectedResourceVersion, this.buildLabels(record), ), ); } catch (error) { if (statusCode(error) === 409) throw new BuildStoreConflictError( "Build record was concurrently modified", ); throw error; } if (terminal(record)) { await this.releaseLock(record.spec.imageKey, record.metadata.name); } } async acquireBuildReconciliationLease( id: string, holder: string, ttlMs: number, ): Promise { if (!holder || !Number.isFinite(ttlMs) || ttlMs <= 0) throw new BuildStoreConflictError("Invalid build reconciliation lease"); for (let attempt = 0; attempt < 8; attempt++) { const current = await this.getBuild(id); if (!current) return; if ( current.status.reconcileLease && Date.parse(current.status.reconcileLease.expiresAt) > Date.now() ) { return; } const lease: BuildReconciliationLease = { holder, token: randomUUID(), expiresAt: new Date(Date.now() + ttlMs).toISOString(), }; const next = structuredClone(current); next.metadata.resourceVersion = String( Number(current.metadata.resourceVersion) + 1, ); next.status.reconcileLease = lease; try { // The old fence token and resourceVersion identify one persisted record. // This CAS installs the successor fence before it can reconcile. await this.replaceBuild(next, current.metadata.resourceVersion); return lease; } catch (error) { if (!(error instanceof BuildStoreConflictError) || attempt === 7) throw error; } } throw new BuildStoreConflictError( "Build reconciliation lease acquisition timed out", ); } async renewBuildReconciliationLease( id: string, token: string, ttlMs: number, ): Promise { if (!token || !Number.isFinite(ttlMs) || ttlMs <= 0) return false; for (let attempt = 0; attempt < 8; attempt++) { const current = await this.getBuild(id); if ( !current || current.status.reconcileLease?.token !== token || Date.parse(current.status.reconcileLease.expiresAt) <= Date.now() ) { return false; } const next = structuredClone(current); next.metadata.resourceVersion = String( Number(current.metadata.resourceVersion) + 1, ); next.status.reconcileLease!.expiresAt = new Date( Date.now() + ttlMs, ).toISOString(); try { await this.replaceBuild(next, current.metadata.resourceVersion, token); return true; } catch (error) { if (!(error instanceof BuildStoreConflictError)) throw error; } } return false; } async releaseBuildReconciliationLease( id: string, token: string, ): Promise { const current = await this.getBuild(id); if (current?.status.reconcileLease?.token !== token) return; const next = structuredClone(current); next.metadata.resourceVersion = String( Number(current.metadata.resourceVersion) + 1, ); delete next.status.reconcileLease; try { await this.replaceBuild(next, current.metadata.resourceVersion, token); } catch (error) { if (!(error instanceof BuildStoreConflictError)) throw error; } } async ownsBuild(imageKey: string, buildId: string): Promise { const lock = await this.read<{ buildId: string }>( map(this.namespace, this.lockName(imageKey), "build-lock"), ); return lock?.buildId === buildId; } private async releaseLock(imageKey: string, buildId: string): Promise { const name = this.lockName(imageKey); const object = await this.readConfigMap( map(this.namespace, name, "build-lock"), ); if (!object || payload<{ buildId: string }>(object)?.buildId !== buildId) return; if (!object.metadata.resourceVersion) return; await this.delete( map( this.namespace, name, "build-lock", undefined, object.metadata.resourceVersion, ), object.metadata.resourceVersion, ); } async getUpload(digest: Sha256Digest): Promise { const record = await this.read( map(this.namespace, this.uploadName(digest), "build-upload"), ); if (!record || record.spec.digest !== digest) return; try { record.status.data = new Uint8Array( await readFile(this.uploadPath(digest)), ); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; record.status.data = new Uint8Array(); } if (record.status.data.byteLength !== record.status.offset) throw new Error(`Upload file for ${digest} does not match its record`); return record; } async createUpload(record: UploadRecord): Promise { const existing = await this.getUpload(record.spec.digest); if (existing) return existing; await mkdir(this.uploadRoot, { recursive: true, mode: 0o700 }); const path = this.uploadPath(record.spec.digest); const handle = await open(path, "wx", 0o600).catch((error) => { if ((error as NodeJS.ErrnoException).code === "EEXIST") return; throw error; }); await handle?.close(); const metadata = structuredClone(record); metadata.status.data = new Uint8Array(); try { const created = (await this.objects.create( map( this.namespace, this.uploadName(record.spec.digest), "build-upload", metadata, ), )) as ConfigMap; const result = payload(created) ?? metadata; result.status.data = new Uint8Array(); return result; } catch (error) { if (statusCode(error) === 409) return (await this.getUpload(record.spec.digest))!; await rm(path, { force: true }); throw error; } } async replaceUpload( record: UploadRecord, expectedResourceVersion: string, ): Promise { const metadata = structuredClone(record); metadata.status.data = new Uint8Array(); try { await this.objects.replace( map( this.namespace, this.uploadName(record.spec.digest), "build-upload", metadata, expectedResourceVersion, ), ); await writeFile(this.uploadPath(record.spec.digest), record.status.data, { mode: 0o600, }); } catch (error) { if (statusCode(error) === 409) throw new BuildStoreConflictError( "Upload record was concurrently modified", ); throw error; } } async deleteUpload(digest: Sha256Digest): Promise { await this.delete( map(this.namespace, this.uploadName(digest), "build-upload"), ); await rm(this.uploadPath(digest), { force: true }); } } export class KubernetesBuildOperations implements BuildKubernetesOperations { constructor( private readonly batch: BatchV1Api, private readonly core: CoreV1Api, ) {} async createJob(job: KubernetesJob): Promise { await this.batch.createNamespacedJob({ namespace: job.metadata.namespace, body: job as unknown as V1Job, fieldManager: "kuber-server", fieldValidation: "Strict", }); } async getJob( namespace: string, name: string, ): Promise { let job: V1Job; try { job = await this.batch.readNamespacedJob({ namespace, name }); } catch (error) { if (statusCode(error) === 404) return; throw error; } const failed = job.status?.conditions?.find( (condition) => condition.type === "Failed" && condition.status === "True", ); const complete = job.status?.conditions?.find( (condition) => condition.type === "Complete" && condition.status === "True", ); let phase: BuildJobObservation["phase"] = "queued"; let podError: string | undefined; let containerStartedAt: string | undefined; if (!failed && !complete) { const pods = await this.core.listNamespacedPod({ namespace, labelSelector: `job-name=${name}`, }); const pod = pods.items .sort( (a, b) => (a.metadata?.creationTimestamp?.getTime() ?? 0) - (b.metadata?.creationTimestamp?.getTime() ?? 0), ) .at(-1); const container = pod?.status?.containerStatuses?.find( (entry) => entry.name === "buildkit", ); phase = pod?.status?.phase === "Running" && container?.state?.running ? "running" : !pod || !pod.spec?.nodeName ? "creating" : "starting"; podError = container?.state?.terminated?.message ?? container?.state?.waiting?.message; containerStartedAt = container?.state?.running?.startedAt?.toISOString(); } else phase = failed ? "failed" : "succeeded"; return { phase, startedAt: phase === "running" ? containerStartedAt : undefined, finishedAt: job.status?.completionTime?.toISOString(), ...((failed?.message || podError) && { error: failed?.message ?? podError, }), }; } async getJobLogs(namespace: string, name: string): Promise { const pods = await this.core.listNamespacedPod({ namespace, labelSelector: `job-name=${name}`, }); const pod = pods.items .sort( (a, b) => (a.metadata?.creationTimestamp?.getTime() ?? 0) - (b.metadata?.creationTimestamp?.getTime() ?? 0), ) .at(-1); if (!pod?.metadata?.name) return ""; return this.core.readNamespacedPodLog({ namespace, name: pod.metadata.name, container: "buildkit", }); } async deleteJob(namespace: string, name: string): Promise { try { await this.batch.deleteNamespacedJob({ namespace, name, propagationPolicy: "Background", }); } catch (error) { if (statusCode(error) !== 404) throw error; } } } export function buildObjectApi(objects: KubernetesObjectApi): BuildObjectApi { return objects as unknown as BuildObjectApi; }