Files
kuber/server/build-controller.ts
2026-10-06 15:31:51 +00:00

1247 lines
41 KiB
TypeScript

import { createHash, randomUUID } from "node:crypto";
import { rm } from "node:fs/promises";
import { join } from "node:path";
import {
BUILD_PROTOCOL_VERSION,
assertSha256Digest,
validateBuildpackUri,
type BuildEvent,
type BuildRequest,
type BuildStatus,
type Sha256Digest,
} from "../shared/build-protocol";
import { createBuildJob, type KubernetesJob } from "./build-job";
import {
createBuildpacksJob,
DEFAULT_BUILDPACKS_IMAGE,
validateBuildpacksImage,
} from "./buildpacks-job";
import {
BUILD_RECORD_API_VERSION,
BuildStoreConflictError,
type BuildRecord,
type BuildReconciliationLease,
type BuildStore,
type UploadRecord,
} from "./build-store";
import type { WorkspaceLease, WorkspaceLeaseProvider } from "./operation-store";
import {
materializeWorkspace,
parseWorkspaceManifest,
type MaterializeCas,
type MaterializeTiming,
} from "./materialize";
import { parseImageReference, resolveRegistryDigest } from "./registry";
import { stageBuildpackPackage } from "./buildpack-package";
export const DEFAULT_MAX_BLOB_BYTES = 1024 * 1024 * 1024;
export const DEFAULT_MAX_UPLOAD_CHUNK_BYTES = 8 * 1024 * 1024;
const SNAPSHOT_CHECK_CONCURRENCY = 20;
export const DEFAULT_MAX_LOG_BYTES = 1024 * 1024;
const BUILD_RECONCILE_LEASE_TTL_MS = 30_000;
const BUILD_RECONCILE_HEARTBEAT_MS = BUILD_RECONCILE_LEASE_TTL_MS / 3;
export interface BuildCas extends MaterializeCas {
has(digest: Sha256Digest): Promise<boolean>;
put(data: Uint8Array, expected?: Sha256Digest): Promise<Sha256Digest>;
}
export type JobPhase =
| "queued"
| "creating"
| "starting"
| "running"
| "succeeded"
| "failed";
export interface BuildJobObservation {
phase: JobPhase;
startedAt?: string;
finishedAt?: string;
error?: string;
}
export interface BuildKubernetesOperations {
createJob(job: KubernetesJob): Promise<void>;
getJob(
namespace: string,
name: string,
): Promise<BuildJobObservation | undefined>;
getJobLogs(namespace: string, name: string): Promise<string | Uint8Array>;
deleteJob(namespace: string, name: string): Promise<void>;
}
type ReconcileOptions = { signal?: AbortSignal };
function throwIfAborted(signal?: AbortSignal): void {
if (signal?.aborted)
throw new Error("Build reconciliation scan was cancelled");
}
export interface BuildControllerOptions {
cas: BuildCas;
store: BuildStore;
kubernetes: BuildKubernetesOperations;
namespace: string;
workspaceRoot: string;
workspaceClaimName: string;
cacheImage: string | ((request: BuildRequest) => string);
imageName?: (request: BuildRequest) => string;
pushImage?: (request: BuildRequest) => string;
pushRegistryInsecure?: boolean;
cacheRegistryInsecure?: boolean;
buildkitImage?: string;
buildpacksImage?: string;
stageBuildpack?: typeof stageBuildpackPackage;
serviceAccountName?: string;
registrySecretName?: string;
nodeSelector?: Record<string, string>;
tolerations?: Array<Record<string, unknown>>;
maxBlobBytes?: number;
maxUploadChunkBytes?: number;
maxLogBytes?: number;
now?: () => Date;
materialize?: typeof materializeWorkspace;
resolveDigest?: typeof resolveRegistryDigest;
onReconcileFailure?: (id: string, error: unknown) => void;
reconcileLeases?: WorkspaceLeaseProvider;
reconcileHolder?: string;
reconcileHeartbeatMs?: number;
reconcileSetTimeout?: typeof setTimeout;
reconcileClearTimeout?: typeof clearTimeout;
}
export class BuildValidationError extends Error {
readonly code = "BUILD_INVALID";
}
export class BuildNotFoundError extends Error {
readonly code = "BUILD_NOT_FOUND";
}
export class BuildConflictError extends Error {
readonly code = "BUILD_CONFLICT";
}
export type SnapshotNegotiation = {
workspace: Sha256Digest;
missing: Sha256Digest[];
ready: boolean;
};
export type UploadProgress = {
digest: Sha256Digest;
size: number;
offset: number;
complete: boolean;
};
async function mapConcurrent<T, R>(
values: T[],
run: (value: T) => Promise<R>,
): Promise<R[]> {
const results = new Array<R>(values.length);
let index = 0;
const worker = async () => {
for (;;) {
const current = index++;
if (current >= values.length) return;
results[current] = await run(values[current]!);
}
};
await Promise.all(
Array.from(
{ length: Math.min(SNAPSHOT_CHECK_CONCURRENCY, values.length) },
worker,
),
);
return results;
}
export function buildImageName(
registry: string,
project: string,
service: string,
): string {
const valid = (value: string) =>
value.length <= 63 && /^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?$/.test(value);
if (!valid(project) || !valid(service))
throw new BuildValidationError(
"Build project and service must be valid Kubernetes names",
);
const owner = registry.replace(/\/+$/, "");
if (!owner) throw new BuildValidationError("Build registry is required");
return `${owner}/kuber/${project}-${service}:latest`;
}
function clone<T>(value: T): T {
return structuredClone(value);
}
function jobWorkspaceSubPath(workspaceRoot: string, subPath: string): string {
const segments = workspaceRoot.split("/").filter(Boolean);
const prefix = segments.at(-1);
if (!prefix || prefix === ".") return subPath;
return subPath ? `${prefix}/${subPath}` : prefix;
}
function validateRequest(request: BuildRequest): void {
if (request.version !== BUILD_PROTOCOL_VERSION)
throw new BuildValidationError("Unsupported build protocol version");
if (!request.id || !request.project || !request.service)
throw new BuildValidationError(
"Build ID, project, and service are required",
);
if (Buffer.byteLength(request.id) > 256)
throw new BuildValidationError("Build ID is too long");
if (!request.spec || !["amd64", "arm64"].includes(request.spec.architecture))
throw new BuildValidationError("Invalid build architecture");
if (
request.spec.builder !== undefined &&
request.spec.builder !== "buildkit" &&
request.spec.builder !== "buildpacks"
)
throw new BuildValidationError("Invalid build builder");
if (request.spec.buildpackUri !== undefined) {
if (request.spec.builder !== "buildpacks")
throw new BuildValidationError(
"buildpackUri requires the buildpacks builder",
);
try {
validateBuildpackUri(request.spec.buildpackUri);
} catch (error) {
throw new BuildValidationError((error as Error).message);
}
}
if (request.spec.builder === "buildpacks") {
if (
request.spec.dockerfile !== undefined ||
request.spec.target !== undefined ||
request.spec.buildArgs?.length
)
throw new BuildValidationError(
"Buildpacks does not support dockerfile, target, or buildArgs",
);
}
if (
!Array.isArray(request.spec.buildArgs) ||
!request.spec.buildArgs.every((value) => typeof value === "string")
) {
throw new BuildValidationError("Build arguments must be strings");
}
assertSha256Digest(request.spec.workspace);
if (
[
request.spec.image,
...(Array.isArray(request.destinations)
? request.destinations.map((destination) => destination?.image)
: []),
].some((image) => typeof image === "string" && /[,"\r\n]/.test(image))
)
throw new BuildValidationError(
"Build image cannot contain BuildKit output separators",
);
parseImageReference(request.spec.image);
if (request.destinations !== undefined) {
if (!Array.isArray(request.destinations))
throw new BuildValidationError("Build destinations must be an array");
const services = new Set([request.service]);
const images = new Set([request.spec.image]);
for (const destination of request.destinations) {
if (
!destination ||
typeof destination.service !== "string" ||
!/^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?$/.test(destination.service) ||
destination.service.length > 63 ||
services.has(destination.service) ||
typeof destination.image !== "string" ||
images.has(destination.image)
)
throw new BuildValidationError(
"Build destinations must have unique services and images",
);
parseImageReference(destination.image);
services.add(destination.service);
images.add(destination.image);
}
}
}
function recordStatus(record: BuildRecord): BuildStatus {
const {
version,
id,
state,
phase,
createdAt,
startedAt,
finishedAt,
digest,
error,
} = record.status;
return {
version,
id,
state,
...(phase && { phase }),
createdAt,
...(startedAt && { startedAt }),
...(finishedAt && { finishedAt }),
...(digest && { digest }),
...(error && { error }),
};
}
export class BuildController {
private readonly now: () => Date;
private readonly maxBlobBytes: number;
private readonly maxUploadChunkBytes: number;
private readonly maxLogBytes: number;
private readonly materializer: typeof materializeWorkspace;
private readonly digestResolver: typeof resolveRegistryDigest;
private readonly imageSubmissionLocks = new Map<string, Promise<unknown>>();
private readonly reconcileLocks = new Map<string, Promise<BuildStatus>>();
private readonly reconcileHolder: string;
constructor(private readonly options: BuildControllerOptions) {
this.now = options.now ?? (() => new Date());
this.maxBlobBytes = options.maxBlobBytes ?? DEFAULT_MAX_BLOB_BYTES;
this.maxUploadChunkBytes =
options.maxUploadChunkBytes ?? DEFAULT_MAX_UPLOAD_CHUNK_BYTES;
this.maxLogBytes = options.maxLogBytes ?? DEFAULT_MAX_LOG_BYTES;
this.materializer = options.materialize ?? materializeWorkspace;
this.digestResolver = options.resolveDigest ?? resolveRegistryDigest;
this.reconcileHolder = options.reconcileHolder ?? randomUUID();
}
async negotiateSnapshot(
workspace: Sha256Digest,
): Promise<SnapshotNegotiation> {
assertSha256Digest(workspace);
if (!(await this.options.cas.has(workspace)))
return { workspace, missing: [workspace], ready: false };
const manifest = parseWorkspaceManifest(
await this.options.cas.get(workspace),
);
const digests = [...new Set(manifest.files.map((file) => file.digest))];
const available = await mapConcurrent(digests, (digest) =>
this.options.cas.has(digest),
);
const missing = digests.filter((_digest, index) => !available[index]);
return { workspace, missing, ready: missing.length === 0 };
}
async beginBlobUpload(
digest: Sha256Digest,
size: number,
): Promise<UploadProgress> {
assertSha256Digest(digest);
if (!Number.isSafeInteger(size) || size < 0 || size > this.maxBlobBytes)
throw new BuildValidationError(
`Blob size must be between 0 and ${this.maxBlobBytes}`,
);
if (await this.options.cas.has(digest)) {
const actualSize = (await this.options.cas.get(digest)).byteLength;
if (actualSize !== size)
throw new BuildConflictError(
"Blob size does not match content already in CAS",
);
return { digest, size, offset: size, complete: true };
}
const current = await this.options.store.getUpload(digest);
if (current) {
if (current.spec.size !== size)
throw new BuildConflictError(
"Upload size does not match the existing upload",
);
return { digest, size, offset: current.status.offset, complete: false };
}
const timestamp = this.now().toISOString();
const upload: UploadRecord = {
apiVersion: BUILD_RECORD_API_VERSION,
kind: "BuildUpload",
metadata: {
name: digest.replace(":", "-"),
resourceVersion: "1",
creationTimestamp: timestamp,
},
spec: { digest, size },
status: { offset: 0, data: new Uint8Array() },
};
const stored = await this.options.store.createUpload(upload);
if (stored.spec.size !== size)
throw new BuildConflictError(
"Upload size does not match the existing upload",
);
return { digest, size, offset: stored.status.offset, complete: false };
}
async uploadBlobChunk(
digest: Sha256Digest,
offset: number,
chunk: Uint8Array,
): Promise<UploadProgress> {
assertSha256Digest(digest);
if (
!(chunk instanceof Uint8Array) ||
chunk.byteLength > this.maxUploadChunkBytes
)
throw new BuildValidationError(
`Upload chunks may not exceed ${this.maxUploadChunkBytes} bytes`,
);
const current = await this.options.store.getUpload(digest);
if (!current) {
if (await this.options.cas.has(digest)) {
const size = (await this.options.cas.get(digest)).byteLength;
return { digest, size, offset: size, complete: true };
}
throw new BuildNotFoundError("Blob upload was not initialized");
}
if (offset !== current.status.offset)
throw new BuildConflictError(
`Upload offset mismatch; expected ${current.status.offset}`,
);
const nextOffset = offset + chunk.byteLength;
if (nextOffset > current.spec.size)
throw new BuildValidationError("Upload exceeds the declared blob size");
const data = new Uint8Array(nextOffset);
data.set(current.status.data);
data.set(chunk, offset);
const next = clone(current);
next.metadata.resourceVersion = String(
Number(current.metadata.resourceVersion) + 1,
);
next.status = { offset: nextOffset, data };
await this.options.store.replaceUpload(
next,
current.metadata.resourceVersion,
);
return {
digest,
size: current.spec.size,
offset: nextOffset,
complete: false,
};
}
async completeBlobUpload(digest: Sha256Digest): Promise<UploadProgress> {
assertSha256Digest(digest);
const current = await this.options.store.getUpload(digest);
if (!current) {
if (await this.options.cas.has(digest)) {
const size = (await this.options.cas.get(digest)).byteLength;
return { digest, size, offset: size, complete: true };
}
throw new BuildNotFoundError("Blob upload was not initialized");
}
if (
current.status.offset !== current.spec.size ||
current.status.data.byteLength !== current.spec.size
) {
throw new BuildConflictError(
`Blob upload is incomplete at offset ${current.status.offset}`,
);
}
try {
await this.options.cas.put(current.status.data, digest);
} catch (error) {
throw new BuildValidationError(
error instanceof Error ? error.message : String(error),
);
}
await this.options.store.deleteUpload(digest);
return {
digest,
size: current.spec.size,
offset: current.spec.size,
complete: true,
};
}
async submitBuild(request: BuildRequest): Promise<BuildStatus> {
request = clone(request);
if (
request.destinations !== undefined &&
!Array.isArray(request.destinations)
)
throw new BuildValidationError("Build destinations must be an array");
if (
request.destinations?.some(
(destination) =>
!destination ||
typeof destination.service !== "string" ||
typeof destination.image !== "string",
)
)
throw new BuildValidationError("Invalid build destination");
if (this.options.imageName)
request.spec.image = this.options.imageName(request);
if (this.options.imageName && request.destinations)
request.destinations = request.destinations.map((destination) => ({
service: destination.service,
image: this.options.imageName!({
...request,
service: destination.service,
spec: { ...request.spec, image: destination.image },
destinations: undefined,
}),
}));
validateRequest(request);
if (request.spec.builder === "buildpacks") {
try {
validateBuildpacksImage(
this.options.buildpacksImage ?? DEFAULT_BUILDPACKS_IMAGE,
);
} catch (error) {
throw new BuildValidationError((error as Error).message);
}
}
const imageKey = `${request.project}\0${request.service}\0${request.spec.image}`;
const started = performance.now();
const previous =
this.imageSubmissionLocks.get(imageKey) ?? Promise.resolve();
const current = previous.then(async () => {
const queueWaitMs = performance.now() - started;
const phases: Record<string, number> = {};
const materialize: MaterializeTiming = { casReadMs: 0, fsWriteMs: 0 };
let outcome = "failure";
let jobCreated = false;
const timed = async <T>(
phase: string,
run: () => Promise<T>,
): Promise<T> => {
const start = performance.now();
try {
return await run();
} finally {
phases[phase] = performance.now() - start;
}
};
try {
const status = await this.submitBuildInternal(
request,
timed,
materialize,
() => {
jobCreated = true;
},
);
outcome = "success";
return status;
} finally {
// Fixed-schema numeric telemetry only: never serialize a request or error.
try {
console.info(
JSON.stringify({
event: "build_submission_timing",
buildRequestId:
/^[0-9a-f]{8}-[0-9a-f]{4}-[1-8][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(
request.id,
)
? request.id
: null,
outcome,
jobCreated,
elapsedMs: performance.now() - started,
queueWaitMs,
phases,
materialize: {
fileCount: materialize.fileCount ?? null,
manifestBytes: materialize.manifestBytes ?? null,
fileBytes: materialize.fileBytes ?? null,
casReadMs: materialize.casReadMs,
fsWriteMs: materialize.fsWriteMs,
},
}),
);
} catch {
// Telemetry must not change submission behavior.
}
}
});
const entry = current.catch(() => undefined);
this.imageSubmissionLocks.set(imageKey, entry);
try {
return await current;
} finally {
if (this.imageSubmissionLocks.get(imageKey) === entry)
this.imageSubmissionLocks.delete(imageKey);
}
}
private async submitBuildInternal(
request: BuildRequest,
timed: <T>(phase: string, run: () => Promise<T>) => Promise<T>,
materialize: MaterializeTiming,
markJobCreated: () => void,
): Promise<BuildStatus> {
const imageKey = `${request.project}\0${request.service}\0${request.spec.image}`;
const snapshot = await timed("snapshotMs", () =>
this.negotiateSnapshot(request.spec.workspace),
);
if (!snapshot.ready)
throw new BuildConflictError(
`Workspace snapshot is incomplete: ${snapshot.missing.join(", ")}`,
);
const createdAt = this.now().toISOString();
const hash = createHash("sha256")
.update(request.id)
.digest("hex")
.slice(0, 24);
const jobName = `kuber-build-${hash}`;
const initial: BuildStatus = {
version: BUILD_PROTOCOL_VERSION,
id: request.id,
state: "queued",
phase: "queued",
createdAt,
};
const record: BuildRecord = {
apiVersion: BUILD_RECORD_API_VERSION,
kind: "BuildRecord",
metadata: {
name: request.id,
resourceVersion: "1",
creationTimestamp: createdAt,
labels: { project: request.project, service: request.service },
},
spec: {
request: clone(request),
imageKey,
jobName,
workspaceSubPath: jobName,
},
status: {
...initial,
logBytes: 0,
logOffset: 0,
nextSequence: 1,
events: [{ type: "status", status: initial }],
},
};
let stored: BuildRecord;
try {
const result = await timed("recordCreateMs", () =>
this.options.store.createBuild(record),
);
stored = result.record;
await timed("supersededCleanupMs", async () => {
for (const superseded of result.superseded ?? []) {
// Deletion is best effort and does not wait for the Job or its pods.
await this.options.kubernetes
.deleteJob(this.options.namespace, superseded.spec.jobName)
.catch(() => undefined);
await rm(
join(this.options.workspaceRoot, superseded.spec.workspaceSubPath),
{ recursive: true, force: true },
).catch(() => undefined);
await rm(
join(
this.options.workspaceRoot,
`${superseded.spec.workspaceSubPath}-buildpack`,
),
{ recursive: true, force: true },
).catch(() => undefined);
}
});
if (!result.created) return recordStatus(stored);
const owns = await timed("initialOwnershipMs", async () => {
const current = await this.options.store.getBuild(request.id);
return (
!!current &&
current.status.state === "queued" &&
(!this.options.store.ownsBuild ||
(await this.options.store.ownsBuild(imageKey, request.id)))
);
});
if (!owns)
throw new BuildConflictError(
`Build '${request.id}' was superseded before Job creation`,
);
} catch (error) {
if (error instanceof BuildStoreConflictError)
throw new BuildConflictError(error.message);
throw error;
}
try {
await timed("materializeMs", () =>
this.materializer(
this.options.cas,
request.spec.workspace,
join(this.options.workspaceRoot, stored.spec.workspaceSubPath),
materialize,
),
);
const packageSubPath = `${stored.spec.workspaceSubPath}-buildpack`;
if (request.spec.buildpackUri)
await timed("buildpackStageMs", () =>
(this.options.stageBuildpack ?? stageBuildpackPackage)(
request.spec.buildpackUri!,
request.spec.architecture,
join(this.options.workspaceRoot, packageSubPath),
),
);
if (
this.options.store.ownsBuild &&
!(await timed("finalOwnershipMs", () =>
this.options.store.ownsBuild!(imageKey, request.id),
))
)
throw new BuildConflictError(
`Build '${request.id}' lost image lock before Job creation`,
);
const job = await timed("jobSpecMs", async () => {
let cacheImage =
typeof this.options.cacheImage === "function"
? this.options.cacheImage(request)
: this.options.cacheImage;
if (request.spec.buildpackUri) {
const suffix = createHash("sha256")
.update(
JSON.stringify([
cacheImage,
request.spec.architecture,
request.spec.buildpackUri,
]),
)
.digest("hex")
.slice(0, 24);
const colon = cacheImage.lastIndexOf(":");
const slash = cacheImage.lastIndexOf("/");
if (cacheImage.includes("@"))
throw new BuildValidationError(
"Build cache must be a tag, not a digest",
);
cacheImage = `${colon > slash ? cacheImage.slice(0, colon) : cacheImage}:buildpack-${suffix}`;
}
const pushImage = this.options.pushImage
? this.options.pushImage(request)
: request.spec.image;
const pushImages = request.destinations?.map((destination) =>
this.options.pushImage
? this.options.pushImage({
...request,
service: destination.service,
spec: { ...request.spec, image: destination.image },
destinations: undefined,
})
: destination.image,
);
const jobOptions = {
name: jobName,
namespace: this.options.namespace,
spec: request.spec,
workspaceClaimName: this.options.workspaceClaimName,
workspaceSubPath: jobWorkspaceSubPath(
this.options.workspaceRoot,
stored.spec.workspaceSubPath,
),
cacheImage,
pushImage,
pushImages,
pushRegistryInsecure: this.options.pushRegistryInsecure,
cacheRegistryInsecure:
this.options.cacheRegistryInsecure ??
this.options.pushRegistryInsecure,
buildkitImage: this.options.buildkitImage,
serviceAccountName: this.options.serviceAccountName,
registrySecretName: this.options.registrySecretName,
nodeSelector: this.options.nodeSelector,
tolerations: this.options.tolerations,
};
return request.spec.builder === "buildpacks"
? createBuildpacksJob({
...jobOptions,
buildpacksImage: this.options.buildpacksImage,
...(request.spec.buildpackUri && {
buildpackSubPath: jobWorkspaceSubPath(
this.options.workspaceRoot,
packageSubPath,
),
}),
cacheRegistryInsecure:
this.options.cacheRegistryInsecure ??
this.options.pushRegistryInsecure,
})
: createBuildJob(jobOptions);
});
await timed("jobCreateMs", () => this.options.kubernetes.createJob(job));
markJobCreated();
await timed("recordUpdateMs", () =>
this.updateBuild(stored, (next) => {
next.status.jobCreated = true;
}),
);
return initial;
} catch (error) {
await timed("failureCleanupMs", async () => {
try {
await this.terminateJob(stored.spec.jobName);
} catch {
// Job cleanup is best effort; the record still releases its lock.
}
await rm(
join(this.options.workspaceRoot, stored.spec.workspaceSubPath),
{ recursive: true, force: true },
).catch(() => undefined);
await rm(
join(
this.options.workspaceRoot,
`${stored.spec.workspaceSubPath}-buildpack`,
),
{ recursive: true, force: true },
).catch(() => undefined);
await this.failBuild(
stored.metadata.name,
error instanceof Error ? error.message : String(error),
);
});
throw error;
}
}
async getBuildStatus(id: string): Promise<BuildStatus> {
const record = await this.requireBuild(id);
return recordStatus(record);
}
async getBuildProject(id: string): Promise<string> {
return (await this.requireBuild(id)).spec.request.project;
}
async getBuildEvents(id: string, afterSequence = 0): Promise<BuildEvent[]> {
if (!Number.isSafeInteger(afterSequence) || afterSequence < 0)
throw new BuildValidationError(
"Event sequence must be a non-negative integer",
);
const record = await this.requireBuild(id);
return clone(
record.status.events.filter(
(event) => event.type === "status" || event.sequence > afterSequence,
),
);
}
async reconcileBuild(
id: string,
options: ReconcileOptions = {},
): Promise<BuildStatus> {
const current = this.reconcileLocks.get(id);
if (current) return current;
const reconciliation = this.reconcileBuildWithLease(id, options.signal);
this.reconcileLocks.set(id, reconciliation);
try {
return await reconciliation;
} finally {
if (this.reconcileLocks.get(id) === reconciliation)
this.reconcileLocks.delete(id);
}
}
private async reconcileBuildWithLease(
id: string,
signal?: AbortSignal,
): Promise<BuildStatus> {
throwIfAborted(signal);
let lease: WorkspaceLease | undefined;
let reconcileLease: BuildReconciliationLease | undefined;
let heartbeat: ReturnType<typeof setTimeout> | undefined;
let heartbeatActive = false;
let leaseOwnershipLost = false;
try {
lease = this.options.reconcileLeases
? await this.options.reconcileLeases.acquire(
`build-reconcile:${id}`,
this.reconcileHolder,
BUILD_RECONCILE_LEASE_TTL_MS,
)
: undefined;
throwIfAborted(signal);
if (this.options.reconcileLeases && !lease)
return this.getBuildStatus(id);
reconcileLease = await this.options.store.acquireBuildReconciliationLease(
id,
this.reconcileHolder,
BUILD_RECONCILE_LEASE_TTL_MS,
);
throwIfAborted(signal);
if (!reconcileLease) return this.getBuildStatus(id);
const reconcileLeaseToken = reconcileLease.token;
const requireLeaseOwnership = async () => {
throwIfAborted(signal);
if (leaseOwnershipLost)
throw new BuildConflictError(
"Build reconciliation lease ownership was lost",
);
const ownsLease = lease
? await lease.renew(BUILD_RECONCILE_LEASE_TTL_MS)
: true;
const ownsReconcileLease =
await this.options.store.renewBuildReconciliationLease(
id,
reconcileLeaseToken,
BUILD_RECONCILE_LEASE_TTL_MS,
);
throwIfAborted(signal);
if (!ownsLease || !ownsReconcileLease) {
leaseOwnershipLost = true;
throw new BuildConflictError(
"Build reconciliation lease ownership was lost",
);
}
};
const heartbeatMs = Math.min(
BUILD_RECONCILE_HEARTBEAT_MS,
this.options.reconcileHeartbeatMs &&
Number.isFinite(this.options.reconcileHeartbeatMs) &&
this.options.reconcileHeartbeatMs > 0
? this.options.reconcileHeartbeatMs
: BUILD_RECONCILE_HEARTBEAT_MS,
);
const scheduleHeartbeat = () => {
if (!heartbeatActive || leaseOwnershipLost) return;
heartbeat = (this.options.reconcileSetTimeout ?? setTimeout)(() => {
void requireLeaseOwnership()
.catch(() => {
leaseOwnershipLost = true;
})
.finally(scheduleHeartbeat);
}, heartbeatMs);
};
heartbeatActive = true;
scheduleHeartbeat();
await requireLeaseOwnership();
return await this.reconcileBuildInternal(
id,
requireLeaseOwnership,
reconcileLease,
);
} finally {
heartbeatActive = false;
if (heartbeat !== undefined)
(this.options.reconcileClearTimeout ?? clearTimeout)(heartbeat);
try {
if (reconcileLease)
await this.options.store.releaseBuildReconciliationLease(
id,
reconcileLease.token,
);
} finally {
await lease?.release();
}
}
}
private async reconcileBuildInternal(
id: string,
requireLeaseOwnership: () => Promise<void>,
reconcileLease: BuildReconciliationLease,
): Promise<BuildStatus> {
let record = await this.requireBuild(id);
if (record.status.state === "succeeded" || record.status.state === "failed")
return recordStatus(record);
if (record.status.jobCreated)
record = await this.captureLogs(
record,
requireLeaseOwnership,
reconcileLease,
);
const observation = await this.options.kubernetes.getJob(
this.options.namespace,
record.spec.jobName,
);
if (!observation) return recordStatus(record);
if (
(observation.phase === "queued" ||
observation.phase === "creating" ||
observation.phase === "starting" ||
observation.phase === "running") &&
(record.status.state === "queued" || observation.phase === "running")
) {
await requireLeaseOwnership();
record = await this.setState(
record,
observation.phase === "running" ? "running" : "queued",
observation.phase === "running"
? {
phase: "running",
startedAt: observation.startedAt ?? this.now().toISOString(),
}
: { phase: observation.phase },
reconcileLease,
);
} else if (observation.phase === "failed") {
await requireLeaseOwnership();
record = await this.setState(
record,
"failed",
{
startedAt: record.status.startedAt ?? observation.startedAt,
finishedAt: observation.finishedAt ?? this.now().toISOString(),
error: observation.error ?? "BuildKit Job failed",
phase: "done",
},
reconcileLease,
);
} else if (observation.phase === "succeeded") {
let digest: Sha256Digest;
try {
digest = await this.digestResolver(record.spec.request.spec.image);
assertSha256Digest(digest);
for (const destination of record.spec.request.destinations ?? []) {
const aliasDigest = await this.digestResolver(destination.image);
assertSha256Digest(aliasDigest);
if (aliasDigest !== digest)
throw new Error(
`Pushed image ${destination.image} has a different digest`,
);
}
} catch (error) {
await requireLeaseOwnership();
record = await this.setState(
record,
"failed",
{
finishedAt: this.now().toISOString(),
error: `Unable to resolve pushed image digest: ${error instanceof Error ? error.message : String(error)}`,
phase: "done",
},
reconcileLease,
);
return recordStatus(record);
}
await requireLeaseOwnership();
record = await this.setState(
record,
"succeeded",
{
startedAt: record.status.startedAt ?? observation.startedAt,
finishedAt: observation.finishedAt ?? this.now().toISOString(),
digest,
phase: "done",
},
reconcileLease,
);
}
return recordStatus(record);
}
async reconcilePendingBuilds(options: ReconcileOptions = {}): Promise<void> {
throwIfAborted(options.signal);
const records = await this.options.store.listBuilds();
throwIfAborted(options.signal);
for (const record of records) {
throwIfAborted(options.signal);
if (
!record.status.jobCreated ||
record.status.state === "succeeded" ||
record.status.state === "failed"
)
continue;
try {
await this.reconcileBuild(record.metadata.name, options);
} catch (error) {
if (options.signal?.aborted) throw error;
try {
this.options.onReconcileFailure?.(record.metadata.name, error);
} catch {
// Failure reporting must not interrupt reconciliation of other builds.
}
}
}
throwIfAborted(options.signal);
await this.options.store.pruneBuilds?.(50);
}
async cancelBuild(id: string): Promise<BuildStatus> {
const record = await this.requireBuild(id);
if (record.status.state === "succeeded" || record.status.state === "failed")
return recordStatus(record);
if (record.status.jobCreated) {
try {
await this.terminateJob(record.spec.jobName);
} catch (error) {
await this.failBuild(
record.metadata.name,
`Unable to cancel build safely: ${error instanceof Error ? error.message : String(error)}`,
);
throw error;
}
}
const next = await this.setState(record, "failed", {
finishedAt: this.now().toISOString(),
error: "Build cancelled",
phase: "done",
cancelled: true,
});
return recordStatus(next);
}
async cleanupBuild(id: string): Promise<void> {
const record = await this.requireBuild(id);
if (record.status.state !== "succeeded" && record.status.state !== "failed")
throw new BuildConflictError("An active build cannot be cleaned up");
if (record.status.jobCreated) await this.terminateJob(record.spec.jobName);
await rm(join(this.options.workspaceRoot, record.spec.workspaceSubPath), {
recursive: true,
force: true,
});
await rm(
join(
this.options.workspaceRoot,
`${record.spec.workspaceSubPath}-buildpack`,
),
{ recursive: true, force: true },
);
}
async getBuildResult(id: string): Promise<{
image: string;
digest: Sha256Digest;
reference: string;
references?: Record<string, string>;
}> {
const record = await this.requireBuild(id);
if (record.status.state !== "succeeded" || !record.status.digest)
throw new BuildConflictError("Build has no immutable image result");
const parsed = parseImageReference(record.spec.request.spec.image);
const image = `${parsed.registry}/${parsed.repository}`;
return {
image,
digest: record.status.digest,
reference: `${image}@${record.status.digest}`,
...(record.spec.request.destinations?.length && {
references: Object.fromEntries(
record.spec.request.destinations.map((destination) => {
const alias = parseImageReference(destination.image);
return [
destination.service,
`${alias.registry}/${alias.repository}@${record.status.digest}`,
];
}),
),
}),
};
}
private async requireBuild(id: string): Promise<BuildRecord> {
const record = await this.options.store.getBuild(id);
if (!record) throw new BuildNotFoundError(`Build '${id}' not found`);
return record;
}
private async terminateJob(name: string): Promise<void> {
await this.options.kubernetes.deleteJob(this.options.namespace, name);
}
private async updateBuild(
record: BuildRecord,
change: (next: BuildRecord) => void,
reconcileLease?: BuildReconciliationLease,
): Promise<BuildRecord> {
// A heartbeat may have advanced this record's resourceVersion while an
// unabortable Kubernetes observation was pending.
if (reconcileLease) record = await this.requireBuild(record.metadata.name);
const next = clone(record);
next.metadata.resourceVersion = String(
Number(record.metadata.resourceVersion) + 1,
);
change(next);
if (reconcileLease) {
if (next.status.reconcileLease?.token !== reconcileLease.token)
throw new BuildConflictError(
"Build reconciliation lease ownership was lost",
);
}
await this.options.store.replaceBuild(
next,
record.metadata.resourceVersion,
reconcileLease?.token,
);
return this.requireBuild(record.metadata.name);
}
private async setState(
record: BuildRecord,
state: BuildStatus["state"],
values: Partial<BuildRecord["status"]>,
reconcileLease?: BuildReconciliationLease,
): Promise<BuildRecord> {
if (record.status.state === state && record.status.phase === values.phase)
return record;
return this.updateBuild(
record,
(next) => {
Object.assign(next.status, values, { state });
next.status.events.push({ type: "status", status: recordStatus(next) });
},
reconcileLease,
);
}
private async failBuild(id: string, message: string): Promise<void> {
const current = await this.requireBuild(id);
if (
current.status.state === "succeeded" ||
current.status.state === "failed"
)
return;
await this.setState(current, "failed", {
finishedAt: this.now().toISOString(),
error: message,
phase: "done",
});
}
private async captureLogs(
record: BuildRecord,
requireLeaseOwnership: () => Promise<void>,
reconcileLease: BuildReconciliationLease,
): Promise<BuildRecord> {
let raw: string | Uint8Array;
try {
raw = await this.options.kubernetes.getJobLogs(
this.options.namespace,
record.spec.jobName,
);
} catch {
return record;
}
const bytes = typeof raw === "string" ? Buffer.from(raw) : Buffer.from(raw);
const offset =
bytes.byteLength < record.status.logOffset ? 0 : record.status.logOffset;
if (bytes.byteLength === offset) return record;
await requireLeaseOwnership();
const delta = bytes.subarray(offset);
return this.updateBuild(
record,
(next) => {
next.status.logOffset = bytes.byteLength;
for (const message of delta
.toString("utf8")
.split(/(?<=\n)/)
.filter(Boolean)) {
const event: BuildEvent = {
type: "log",
id: record.metadata.name,
sequence: next.status.nextSequence++,
message,
};
next.status.events.push(event);
next.status.logBytes += Buffer.byteLength(message);
}
while (next.status.logBytes > this.maxLogBytes) {
const index = next.status.events.findIndex(
(event) => event.type === "log",
);
if (index === -1) break;
const [removed] = next.status.events.splice(index, 1);
if (removed?.type === "log")
next.status.logBytes -= Buffer.byteLength(removed.message);
}
},
reconcileLease,
);
}
}