Files
kuber/lib/build.ts
T
2026-09-27 10:40:50 +00:00

700 lines
19 KiB
TypeScript

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<void>;
export class TaskScheduler {
private inflight = 0;
private maxInflight: number;
private maxPerSecond: number;
private requestStarts: number[] = [];
private waiting: Array<() => void> = [];
private rateGate: Promise<void> = 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<void> {
if (this.inflight < this.maxInflight) {
this.inflight++;
return;
}
await new Promise<void>((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<void> {
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<T>(fn: () => Promise<T>): Promise<T> {
await this.acquire();
try {
await this.waitForRateLimit();
return await fn();
} finally {
this.release();
}
}
}
async function runConcurrent<T>(
values: T[],
concurrency: number,
run: (value: T) => Promise<void>,
): Promise<void> {
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 = <T>(
path: string,
init?: ApiRequestInit,
options?: ApiRequestOptions,
) => Promise<T>;
export type BuildOptions = {
registry?: string;
request?: ApiRequester;
pollIntervalMs?: number;
sleep?: (milliseconds: number) => Promise<void>;
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<void>;
type BuildReporter = {
progress?: ProgressReporter;
stream?: Writable;
};
export type BuildResult = {
built: string[];
changed: string[];
images: Record<string, string>;
};
type SnapshotNegotiation = {
workspace: Sha256Digest;
missing: Sha256Digest[];
ready: boolean;
};
type ImageResult = {
image: string;
digest: Sha256Digest;
reference: string;
};
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<string> {
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<void> {
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<void> {
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<SnapshotNegotiation>("/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<BuildStatus["state"]>,
service?: string,
): Promise<number> {
if (event.type === "status") {
if (!reportedStates?.has(event.status.state)) {
reportedStates?.add(event.status.state);
await reporter?.progress?.(
`${service ? `${service}: ` : ""}Build ${event.status.state}`,
);
}
return 0;
}
if (reporter?.stream)
reporter.stream.write(
service ? `[${service}] ${event.message}` : event.message,
);
else
await reporter?.progress?.(
`${service ? `${service}: ` : ""}${event.message.trimEnd()}`,
);
return event.sequence;
}
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<T>(
request: ApiRequester,
path: string,
init: ApiRequestInit | undefined,
pollIntervalMs: number,
sleep: (milliseconds: number) => Promise<void>,
signal?: AbortSignal,
): Promise<T> {
for (let attempt = 1; attempt <= MAX_BUILD_POLL_ATTEMPTS; attempt++) {
try {
signal?.throwIfAborted();
return await request<T>(
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,
reporter: BuildReporter | undefined,
pollIntervalMs: number,
sleep: (milliseconds: number) => Promise<void>,
initial: BuildStatus,
signal?: AbortSignal,
service?: string,
): Promise<BuildStatus> {
let status = initial;
let sequence = 0;
const reportedStates = new Set<BuildStatus["state"]>();
for (;;) {
const events = await requestBuildPoll<BuildEvent[]>(
request,
`/builds/${encodeURIComponent(id)}/events?after=${sequence}`,
undefined,
pollIntervalMs,
sleep,
signal,
);
for (const event of events)
sequence = Math.max(
sequence,
await reportBuildEvent(event, reporter, reportedStates, service),
);
if (status.state === "succeeded" || status.state === "failed")
return status;
status = await requestBuildPoll<BuildStatus>(
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<Record<string, string>> {
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<ImageResult>(
"/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<BuildResult> {
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);
let next = 0;
let failure: unknown;
let failed = false;
const worker = async () => {
while (!failed && next < plans.length) {
const index = next++;
const plan = plans[index]!;
try {
options.signal?.throwIfAborted();
await reporter?.progress?.(`Building ${plan.name}`);
references[index] = await buildPlan(plan);
} catch (error) {
if (!failed) failure = error;
failed = true;
}
}
};
const buildPlan = async (plan: BuildPlan): Promise<string> => {
const id = randomUUID();
const buildRequest: BuildRequest = {
version: BUILD_PROTOCOL_VERSION,
id,
project,
service: plan.name,
spec: {
architecture: resolveComposeArch(compose),
image: plan.image,
context: plan.context,
dockerfile: plan.dockerfile,
target: plan.target,
buildArgs: plan.buildArgs,
workspace: snapshot.digest,
},
};
const initial = await request<BuildStatus>(
"/builds",
{
method: "POST",
json: buildRequest,
signal: options.signal,
},
{ timeoutMs: 300_000 },
);
const status = await waitForBuild(
id,
request,
reporter,
options.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS,
options.sleep ?? ((milliseconds) => Bun.sleep(milliseconds)),
initial,
options.signal,
plan.name,
);
if (status.state !== "succeeded")
throw new Error(
`Build failed for service ${plan.name}: ${status.error ?? "unknown error"}`,
);
return (
await request<ImageResult>(`/builds/${encodeURIComponent(id)}/result`, {
signal: options.signal,
})
).reference;
};
await Promise.all(
Array.from({ length: Math.min(concurrency, plans.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,
};
}