import { randomUUID } from "node:crypto"; import { execFile } from "node:child_process"; import { isAbsolute, relative, resolve, sep } from "node:path"; import { promisify } from "node:util"; import type { Writable } from "node:stream"; import type { ComposeSpecification, Service } from "../schema/docker.d"; import { BUILD_PROTOCOL_VERSION, type BuildEvent, type BuildRequest, type BuildStatus, type Sha256Digest, } from "../shared/build-protocol"; import { resolveComposeArch } from "./arch"; import { apiRequest, type ApiRequestInit, type ApiRequestOptions } from "./api"; import { DEFAULT_REGISTRY } from "./config"; import { enumerateWorkspace, gitAvailable, serializeWorkspaceManifest, type WorkspaceSnapshot, } from "./workspace"; const execFileAsync = promisify(execFile); const UPLOAD_CHUNK_BYTES = 8 * 1024 * 1024; const DEFAULT_POLL_INTERVAL_MS = 1_000; const BUILD_POLL_REQUEST_TIMEOUT_MS = 300_000; const MAX_BUILD_POLL_ATTEMPTS = 3; const DEFAULT_CONCURRENT_BUILDS = 3; export const MAX_CONCURRENT_REQUESTS = 20; export const MAX_REQUESTS_PER_SECOND = 40; export type SchedulerClock = { now(): number; }; export type SchedulerSleep = (ms: number) => Promise; export class TaskScheduler { private inflight = 0; private maxInflight: number; private maxPerSecond: number; private requestStarts: number[] = []; private waiting: Array<() => void> = []; private rateGate: Promise = Promise.resolve(); private clock: SchedulerClock; private sleep: SchedulerSleep; constructor(options?: { maxInflight?: number; maxPerSecond?: number; clock?: SchedulerClock; sleep?: SchedulerSleep; }) { this.maxInflight = options?.maxInflight ?? MAX_CONCURRENT_REQUESTS; this.maxPerSecond = options?.maxPerSecond ?? MAX_REQUESTS_PER_SECOND; if (!Number.isSafeInteger(this.maxInflight) || this.maxInflight < 1) throw new RangeError("maxInflight must be a positive integer"); if (!Number.isSafeInteger(this.maxPerSecond) || this.maxPerSecond < 1) throw new RangeError("maxPerSecond must be a positive integer"); this.clock = options?.clock ?? { now: () => Date.now() }; this.sleep = options?.sleep ?? ((ms) => Bun.sleep(ms)); } get currentInflight(): number { return this.inflight; } get maxConcurrent(): number { return this.maxInflight; } get currentRequestStarts(): number { this.pruneOldStarts(); return this.requestStarts.length; } private pruneOldStarts(): void { const cutoff = this.clock.now() - 1000; while (this.requestStarts.length > 0 && this.requestStarts[0]! <= cutoff) { this.requestStarts.shift(); } } private async acquire(): Promise { if (this.inflight < this.maxInflight) { this.inflight++; return; } await new Promise((resolve) => { this.waiting.push(resolve); }); } private release(): void { if (this.waiting.length > 0) { const next = this.waiting.shift()!; next(); } else { this.inflight--; } } private waitForRateLimit(): Promise { const reservation = this.rateGate.then(async () => { for (;;) { this.pruneOldStarts(); if (this.requestStarts.length < this.maxPerSecond) { this.requestStarts.push(this.clock.now()); return; } const oldest = this.requestStarts[0]!; await this.sleep(Math.max(1, oldest + 1000 - this.clock.now())); } }); this.rateGate = reservation.catch(() => {}); return reservation; } async run(fn: () => Promise): Promise { await this.acquire(); try { await this.waitForRateLimit(); return await fn(); } finally { this.release(); } } } async function runConcurrent( values: T[], concurrency: number, run: (value: T) => Promise, ): Promise { let index = 0; const worker = async () => { for (;;) { const current = index++; if (current >= values.length) return; await run(values[current]!); } }; await Promise.all( Array.from({ length: Math.min(concurrency, values.length) }, worker), ); } export type ApiRequester = ( path: string, init?: ApiRequestInit, options?: ApiRequestOptions, ) => Promise; export type BuildOptions = { registry?: string; request?: ApiRequester; pollIntervalMs?: number; sleep?: (milliseconds: number) => Promise; snapshot?: WorkspaceSnapshot; workspaceRoot?: string; scheduler?: TaskScheduler; /** Maximum number of service images built at once. */ buildConcurrency?: number; signal?: AbortSignal; }; type BuildPlan = { name: string; image: string; context: string; dockerfile?: string; target?: string; buildArgs: string[]; }; type ProgressReporter = (message: string) => void | Promise; type BuildReporter = { progress?: ProgressReporter; stream?: Writable; /** Per-image output and lifecycle hooks, including for queued images. */ service?: (name: string) => BuildReporter | undefined; settled?: (name: string, error?: unknown) => void; }; export type BuildResult = { built: string[]; changed: string[]; images: Record; }; type SnapshotNegotiation = { workspace: Sha256Digest; missing: Sha256Digest[]; ready: boolean; }; type ImageResult = { image: string; digest: Sha256Digest; reference: string; references?: Record; }; function posixRelative(root: string, path: string): string { return relative(root, path).split(sep).join("/") || "."; } function assertInsideRepo( repoRoot: string, path: string, description: string, service: string, ): string { const value = relative(repoRoot, path); if (value.startsWith(`..${sep}`) || value === ".." || isAbsolute(value)) { throw new Error( `${description} must stay inside the git repo for service ${service}`, ); } return posixRelative(repoRoot, path); } function resolveBuildArgs(service: Service): string[] { if ( !service.build || typeof service.build === "string" || !service.build.args ) return []; if (Array.isArray(service.build.args)) return [...service.build.args]; return Object.entries(service.build.args) .filter(([, value]) => value !== null) .map(([key, value]) => `${key}=${String(value)}`); } function resolveBuildPlan( project: string, name: string, service: Service, cwd: string, repoRoot: string, registry: string, ): BuildPlan | undefined { if (!service.build) return; const build = service.build; const contextInput = typeof build === "string" ? build : (build.context ?? "."); if (contextInput.includes("://")) throw new Error( `Remote build context is not supported for service ${name}`, ); if (typeof build !== "string" && build.dockerfile_inline) throw new Error(`dockerfile_inline is not supported for service ${name}`); const contextPath = resolve(cwd, contextInput); const context = assertInsideRepo( repoRoot, contextPath, "Build context", name, ); const dockerfilePath = typeof build === "string" || !build.dockerfile ? undefined : resolve(contextPath, build.dockerfile); return { name, // The server replaces this requested name with its configured imageName. image: getBuildImageName(project, name, registry), context, dockerfile: dockerfilePath ? assertInsideRepo(repoRoot, dockerfilePath, "Dockerfile", name) : undefined, target: typeof build === "string" ? undefined : build.target, buildArgs: resolveBuildArgs(service), }; } export function parseImageManifestDigest(output: string): string { const manifest = JSON.parse(output) as { digest?: unknown }; if ( typeof manifest.digest !== "string" || !/^sha256:[a-f0-9]{64}$/.test(manifest.digest) ) throw new Error("Registry response did not contain a valid image digest"); return manifest.digest; } export function toPinnedImage(image: string, digest: string): string { if (!/^sha256:[a-f0-9]{64}$/.test(digest)) throw new Error(`Invalid image digest ${digest}`); return `${image}@${digest}`; } export function getBuildImageName( project: string, service: string, registry = DEFAULT_REGISTRY, ): string { return `${registry.replace(/\/+$/, "")}/kuber/${project}-${service}:latest`; } export function imageDigestChanged( before: string | undefined, after: string | undefined, ): boolean { return !before || !after || before !== after; } export async function getRepoRoot(cwd: string): Promise { if (!(await gitAvailable())) return cwd; try { const { stdout } = await execFileAsync("git", [ "-C", cwd, "rev-parse", "--show-toplevel", ]); return stdout.trim(); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return cwd; throw error; } } async function uploadBlob( digest: Sha256Digest, data: Uint8Array, request: ApiRequester, scheduler: TaskScheduler, project?: string, ): Promise { const uploadPath = `/blobs/${encodeURIComponent(digest)}/uploads`; const projectQuery = project ? `?project=${encodeURIComponent(project)}` : ""; const path = `${uploadPath}${projectQuery}`; const progress = await scheduler.run(() => request<{ offset: number; complete: boolean }>(path, { method: "POST", json: { size: data.byteLength }, }), ); let offset = progress.offset; while (!progress.complete && offset < data.byteLength) { const chunk = data.subarray(offset, offset + UPLOAD_CHUNK_BYTES); const uploaded = await scheduler.run(() => request<{ offset: number }>(path, { method: "PATCH", headers: { "content-type": "application/octet-stream", "upload-offset": String(offset), }, body: chunk, }), ); if (uploaded.offset <= offset) throw new Error(`Blob upload for ${digest} made no progress`); offset = uploaded.offset; } if (!progress.complete) { await scheduler.run(() => request(`${uploadPath}/complete${projectQuery}`, { method: "POST", json: {}, }), ); } } export async function uploadWorkspaceSnapshot( snapshot: WorkspaceSnapshot, request: ApiRequester = apiRequest, reporter?: BuildReporter, scheduler?: TaskScheduler, project?: string, ): Promise { const blobs = new Map(snapshot.blobs.map((blob) => [blob.digest, blob.data])); blobs.set(snapshot.digest, serializeWorkspaceManifest(snapshot.manifest)); const requestScheduler = scheduler ?? new TaskScheduler(); for (;;) { const negotiation = await requestScheduler.run(() => request("/snapshots/negotiate", { method: "POST", json: { workspace: snapshot.digest, ...(project && { project }) }, }), ); if (negotiation.ready) return; if (negotiation.missing.length === 0) throw new Error( "Snapshot negotiation is incomplete but reported no missing blobs", ); await runConcurrent( negotiation.missing, requestScheduler.maxConcurrent, async (digest) => { const data = blobs.get(digest); if (!data) throw new Error(`Server requested unknown workspace blob ${digest}`); await reporter?.progress?.(`Uploading ${digest}`); await uploadBlob(digest, data, request, requestScheduler, project); }, ); } } async function reportBuildEvent( event: BuildEvent, reporter?: BuildReporter, reportedStates?: Set, service?: string, ): Promise { if (event.type === "status") { const phase = event.status.phase ?? (event.status.state === "succeeded" || event.status.state === "failed" ? "done" : event.status.state); if (!reportedStates?.has(phase)) { reportedStates?.add(phase); await reporter?.progress?.( `${service ? `${service}: ` : ""}Build ${phase}${event.status.state === "failed" ? `: ${event.status.error ?? "unknown error"}` : ""}`, ); } return; } if (reporter?.stream) reporter.stream.write( service ? `[${service}] ${event.message}` : event.message, ); else await reporter?.progress?.( `${service ? `${service}: ` : ""}${event.message.trimEnd()}`, ); } function isTransientBuildPollError(error: unknown): boolean { if (!(error instanceof Error) || error.name === "AbortError") return false; if (error.name === "TimeoutError" || error instanceof TypeError) return true; const code = "code" in error && typeof error.code === "string" ? error.code : error.cause && typeof error.cause === "object" && "code" in error.cause && typeof error.cause.code === "string" ? error.cause.code : undefined; return ( code === "ECONNABORTED" || code === "ECONNRESET" || code === "ECONNREFUSED" || code === "EAI_AGAIN" || code === "ETIMEDOUT" ); } async function requestBuildPoll( request: ApiRequester, path: string, init: ApiRequestInit | undefined, pollIntervalMs: number, sleep: (milliseconds: number) => Promise, signal?: AbortSignal, ): Promise { for (let attempt = 1; attempt <= MAX_BUILD_POLL_ATTEMPTS; attempt++) { try { signal?.throwIfAborted(); return await request( path, { ...init, signal }, { timeoutMs: BUILD_POLL_REQUEST_TIMEOUT_MS }, ); } catch (error) { if ( !isTransientBuildPollError(error) || attempt === MAX_BUILD_POLL_ATTEMPTS ) { throw error; } signal?.throwIfAborted(); await sleep(pollIntervalMs); } } throw new Error("Build poll retries exhausted"); } async function waitForBuild( id: string, request: ApiRequester, reporters: Array<{ name: string; reporter?: BuildReporter }>, pollIntervalMs: number, sleep: (milliseconds: number) => Promise, initial: BuildStatus, signal?: AbortSignal, prefixService = false, ): Promise { let status = initial; let sequence = 0; let failureLog = ""; // Each physical build belongs only to its destination services. A shared // reporter (the caller's global sink) receives each event just once. const owners = new Map; }>(); for (const { name, reporter } of reporters) if (!owners.has(reporter)) owners.set(reporter, { name, reporter, states: new Set() }); for (;;) { const events = await requestBuildPoll( request, `/builds/${encodeURIComponent(id)}/events?after=${sequence}`, undefined, pollIntervalMs, sleep, signal, ); for (const event of events) { if (event.type === "log" && event.sequence <= sequence) continue; if (event.type === "log") failureLog = (failureLog + event.message).slice(-8_192); for (const { name, reporter, states } of owners.values()) await reportBuildEvent( event, reporter, states, prefixService ? name : undefined, ); if (event.type === "log") sequence = event.sequence; } if (status.state === "succeeded" || status.state === "failed") { const details = failureLog.trim(); if ( status.state === "failed" && details && !status.error?.includes(details) ) return { ...status, error: `${status.error ?? "BuildKit Job failed"}\n${details}`, }; return status; } status = await requestBuildPoll( request, `/builds/${encodeURIComponent(id)}/reconcile`, { method: "POST", json: {}, }, pollIntervalMs, sleep, signal, ); if (status.state !== "succeeded" && status.state !== "failed") { signal?.throwIfAborted(); await sleep(pollIntervalMs); } } } export async function resolveBuildImages( project: string, compose: ComposeSpecification, options: BuildOptions = {}, ): Promise> { const request = options.request ?? apiRequest; const services = Object.entries(compose.services ?? {}) .filter(([, definition]) => definition.build) .map(([service]) => service); const references: string[] = new Array(services.length); let next = 0; await Promise.all( Array.from( { length: Math.min(DEFAULT_CONCURRENT_BUILDS, services.length) }, async () => { for (;;) { const index = next++; if (index >= services.length) return; const service = services[index]!; options.signal?.throwIfAborted(); try { references[index] = ( await request("/images/resolve", { method: "POST", json: { project, service }, signal: options.signal, }) ).reference; } catch (error) { throw new Error( `Cannot resolve a published image for service ${service}. Run kuber up to build it.`, { cause: error }, ); } } }, ), ); return Object.fromEntries( services.map((service, index) => [service, references[index]!]), ); } export async function buildServices( project: string, compose: ComposeSpecification, cwd = process.cwd(), reporter?: BuildReporter, options: BuildOptions = {}, ): Promise { if (!Object.values(compose.services ?? {}).some((service) => service.build)) return { built: [], changed: [], images: {} }; const concurrency = options.buildConcurrency ?? DEFAULT_CONCURRENT_BUILDS; if (!Number.isSafeInteger(concurrency) || concurrency < 1) throw new RangeError("buildConcurrency must be a positive integer"); const request = options.request ?? apiRequest; let repoRoot = options.workspaceRoot; if (!repoRoot && !options.snapshot) { try { repoRoot = await getRepoRoot(cwd); } catch (error) { repoRoot = cwd; } } repoRoot ??= cwd; const snapshot = options.snapshot ?? (await enumerateWorkspace(repoRoot)); const plans = Object.entries(compose.services ?? {}).flatMap( ([name, service]) => { const plan = resolveBuildPlan( project, name, service, cwd, repoRoot, options.registry ?? DEFAULT_REGISTRY, ); return plan ? [plan] : []; }, ); await uploadWorkspaceSnapshot( snapshot, request, reporter, options.scheduler, project, ); const references: string[] = new Array(plans.length); const reporters = new Map( plans.map((plan) => [ plan.name, reporter?.service?.(plan.name) ?? reporter, ]), ); const architecture = resolveComposeArch(compose); const groups = new Map(); plans.forEach((plan, index) => { const key = JSON.stringify({ workspace: snapshot.digest, architecture, context: plan.context, dockerfile: plan.dockerfile, target: plan.target, buildArgs: plan.buildArgs, }); const group = groups.get(key) ?? []; group.push(index); groups.set(key, group); }); const batches = [...groups.values()]; let next = 0; let failure: unknown; let failed = false; const worker = async () => { while (!failed && next < batches.length) { const batch = batches[next++]!; const groupPlans = batch.map((index) => plans[index]!); try { options.signal?.throwIfAborted(); for (const plan of groupPlans) { await reporters .get(plan.name) ?.progress?.( `${reporter?.service ? "" : `${plan.name}: `}Build queued`, ); } const result = await buildPlan(groupPlans); batch.forEach((index) => { const plan = plans[index]!; const reference = index === batch[0] ? result.reference : result.references?.[plan.name]; if (!reference) throw new Error( `Build result missing image for service ${plan.name}`, ); references[index] = reference; }); for (const plan of groupPlans) reporter?.settled?.(plan.name); } catch (error) { if (!failed) failure = error; failed = true; for (const plan of groupPlans) reporter?.settled?.(plan.name, error); } } }; const buildPlan = async (groupPlans: BuildPlan[]): Promise => { const plan = groupPlans[0]!; const id = randomUUID(); const buildRequest: BuildRequest = { version: BUILD_PROTOCOL_VERSION, id, project, service: plan.name, ...(groupPlans.length > 1 && { destinations: groupPlans .slice(1) .map(({ name, image }) => ({ service: name, image })), }), spec: { architecture, image: plan.image, context: plan.context, dockerfile: plan.dockerfile, target: plan.target, buildArgs: plan.buildArgs, workspace: snapshot.digest, }, }; const initial = await request( "/builds", { method: "POST", json: buildRequest, signal: options.signal, }, { timeoutMs: 300_000 }, ); const status = await waitForBuild( id, request, groupPlans.map(({ name }) => ({ name, reporter: reporters.get(name) })), options.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS, options.sleep ?? ((milliseconds) => Bun.sleep(milliseconds)), initial, options.signal, !reporter?.service, ); if (status.state !== "succeeded") throw new Error( `Build failed for service ${plan.name}: ${status.error ?? "unknown error"}`, ); return await request( `/builds/${encodeURIComponent(id)}/result`, { signal: options.signal, }, ); }; await Promise.all( Array.from({ length: Math.min(concurrency, batches.length) }, worker), ); if (failed) throw failure; const images = Object.fromEntries( plans.map((plan, index) => [plan.name, references[index]!]), ); return { built: plans.map((plan) => plan.name), changed: plans.map((plan) => plan.name), images, }; }