Files
kuber/tests/lib/build-api.test.ts
2026-10-06 15:31:51 +00:00

2092 lines
66 KiB
TypeScript

import { afterEach, describe, expect, test } from "bun:test";
import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { Writable } from "node:stream";
import type { ApiRequestInit, ApiRequestOptions } from "../../lib/api";
import {
buildServices,
resolveBuildImages,
uploadWorkspaceSnapshot,
type ApiRequester,
TaskScheduler,
MAX_CONCURRENT_REQUESTS,
MAX_REQUESTS_PER_SECOND,
} from "../../lib/build";
import {
enumerateWorkspace,
serializeWorkspaceManifest,
workspaceManifestDigest,
type WorkspaceSnapshot,
} from "../../lib/workspace";
import type { BuildRequest } from "../../shared/build-protocol";
import { ApiStreamUnsupportedError } from "../../lib/api";
const directories: string[] = [];
test("auto URI is submitted and grouped by exact package identity while bare auto stays distinct", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-auto-uri-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
await writeFile(join(root, "package.json"), "{}\n");
const uri =
"https://github.com/example/bun/releases/download/v1/buildpack.cnb";
const submitted: BuildRequest[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path.includes("/events")) return [] as T;
if (path === "/builds") {
const build = init!.json as BuildRequest;
submitted.push(build);
return { state: "succeeded", id: build.id } as T;
}
if (path.endsWith("/result")) {
const build = submitted.find((build) => path.includes(build.id))!;
return {
reference: `image:${build.service}`,
references: Object.fromEntries(
(build.destinations ?? []).map(({ service }) => [
service,
`image:${service}`,
]),
),
} as T;
}
throw new Error(path);
};
const result = await buildServices(
"shop",
{
services: {
one: { build: `auto:${uri}` },
same: { build: `auto:${uri}` },
other: { build: `auto:${uri.replace("/v1/", "/v2/")}` },
bare: { build: "auto" },
},
},
root,
undefined,
{ request, sleep: async () => {}, workspaceRoot: root },
);
expect(result.built).toHaveLength(4);
expect(submitted).toHaveLength(3);
expect(
submitted.every(
({ spec }) =>
spec.builder === "buildpacks" &&
spec.context === "." &&
spec.dockerfile === undefined,
),
).toBe(true);
expect(
submitted.find(({ service }) => service === "one")?.destinations?.[0]
?.service,
).toBe("same");
expect(
submitted.find(({ service }) => service === "one")?.spec.buildpackUri,
).toBe(uri);
expect(
submitted.find(({ service }) => service === "bare")?.spec.buildpackUri,
).toBeUndefined();
await expect(
buildServices(
"shop",
{ services: { bad: { build: "auto:http://localhost/package.cnb" } } },
root,
undefined,
{ request, workspaceRoot: root },
),
).rejects.toThrow(/URI/);
});
function emptySnapshot(): WorkspaceSnapshot {
const manifest = { version: 1 as const, files: [] };
return { manifest, digest: workspaceManifestDigest(manifest), blobs: [] };
}
function createSnapshotWithBlobs(
count: number,
blobSize: number,
): WorkspaceSnapshot {
const blobs = Array.from({ length: count }, (_, i) => {
const data = new Uint8Array(blobSize);
data.fill(i + 1);
const digest = `sha256:${"0".repeat(62)}${String(i + 1).padStart(2, "0")}`;
return {
digest: digest as import("../../shared/build-protocol").Sha256Digest,
data,
};
});
const manifest = {
version: 1 as const,
files: blobs
.map((blob, i) => ({
path: `file${i}.txt`,
type: "file" as const,
digest: blob.digest,
size: blob.data.byteLength,
mode: 0o644 as const,
}))
.sort((left, right) =>
Buffer.from(left.path).compare(Buffer.from(right.path)),
),
};
return {
manifest,
digest: workspaceManifestDigest(manifest),
blobs,
};
}
afterEach(async () => {
await Promise.all(
directories
.splice(0)
.map((path) => rm(path, { recursive: true, force: true })),
);
});
describe("authenticated build API pipeline", () => {
test("reports aggregate unique missing bytes per chunk through renegotiation", async () => {
const snapshot = createSnapshotWithBlobs(2, 1024);
const [first, second] = snapshot.blobs;
const total =
first!.data.byteLength +
second!.data.byteLength +
serializeWorkspaceManifest(snapshot.manifest).byteLength;
const reports: Array<[number, number]> = [];
let release!: () => void;
let patchStarted!: () => void;
const held = new Promise<void>((resolve) => {
release = resolve;
});
const started = new Promise<void>((resolve) => {
patchStarted = resolve;
});
let negotiations = 0;
const patches = new Map<string, number>();
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
negotiations++;
return (
negotiations === 1
? {
ready: false,
missing: [
first!.digest,
first!.digest,
second!.digest,
snapshot.digest,
],
}
: negotiations === 2
? { ready: false, missing: [first!.digest, snapshot.digest] }
: { ready: true, missing: [] }
) as T;
}
const digest = decodeURIComponent(path.split("/")[2]!);
if (init?.method === "PATCH") {
if (digest === first!.digest) {
patchStarted();
await held;
}
patches.set(digest, (patches.get(digest) ?? 0) + 1);
return { offset: (init.body as Uint8Array).byteLength } as T;
}
if (init?.method === "POST" && !path.includes("/complete"))
return {
offset: patches.has(digest)
? digest === snapshot.digest
? total - 2048
: 1024
: 0,
complete: patches.has(digest),
} as T;
return { complete: true } as T;
};
const upload = uploadWorkspaceSnapshot(snapshot, request, {
upload: (bytes, size) => reports.push([bytes, size]),
});
await started;
expect(reports.at(-1)?.[1]).toBe(total);
expect(reports.at(-1)?.[0]).toBeLessThan(total);
release();
await upload;
expect(patches.get(first!.digest)).toBe(1);
expect(patches.get(second!.digest)).toBe(1);
expect(patches.get(snapshot.digest)).toBe(1);
expect(reports.at(-1)).toEqual([total, total]);
expect(
reports.every(([bytes, size]) => bytes <= size && size === total),
).toBe(true);
});
test("settles the upload with its original error before starting a build", async () => {
const snapshot = emptySnapshot();
const failure = new Error("upload interrupted");
let builds = 0;
const settled: unknown[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate")
return { ready: false, missing: [snapshot.digest] } as T;
if (path === "/builds") {
builds++;
throw new Error("unexpected build");
}
if (init?.method === "PATCH") throw failure;
return { offset: 0, complete: false } as T;
};
await expect(
buildServices(
"shop",
{ services: { web: { build: "." } } },
process.cwd(),
{
uploadSettled: (error) => {
settled.push(error);
},
},
{ request, snapshot },
),
).rejects.toBe(failure);
expect(settled).toEqual([failure]);
expect(builds).toBe(0);
});
test("keeps builder groups separate even when snapshot digests coincide", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-same-digest-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
const snapshot = emptySnapshot();
const submissions: BuildRequest[] = [];
const uploads: string[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
uploads.push((init!.json as { workspace: string }).workspace);
return { ready: true } as T;
}
if (path === "/builds") {
const build = init!.json as BuildRequest;
submissions.push(build);
return { state: "queued", id: build.id } as T;
}
if (path.includes("/events")) return [] as T;
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.endsWith("/result")) return { reference: "image:built" } as T;
throw new Error(path);
};
await buildServices(
"shop",
{
services: {
docker: { build: "." },
auto: { build: "auto" },
},
},
root,
undefined,
{ request, snapshot, sleep: async () => {}, pollIntervalMs: 0 },
);
expect(uploads).toEqual([snapshot.digest, snapshot.digest]);
expect(submissions).toHaveLength(2);
expect(submissions.map(({ spec }) => spec.workspace)).toEqual([
snapshot.digest,
snapshot.digest,
]);
expect(submissions[0]?.spec).not.toHaveProperty("builder");
expect(submissions[1]?.spec).toHaveProperty("builder", "buildpacks");
});
test("auto fails before uploading when no Git repository is available", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-auto-gitless-"));
directories.push(root);
const calls: string[] = [];
const request: ApiRequester = async <T>(path: string) => {
calls.push(path);
return { ready: true } as T;
};
await expect(
buildServices(
"shop",
{ services: { web: { build: "auto" } } },
root,
undefined,
{ request, snapshot: emptySnapshot() },
),
).rejects.toThrow(/requires Git and a Git repository/);
expect(calls).toEqual([]);
});
test("auto uses a contained selected subproject context and Git-filtered snapshot", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-auto-subproject-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
await writeFile(join(root, ".gitignore"), "*.secret\n");
await mkdir(join(root, "apps", "web"), { recursive: true });
await writeFile(join(root, "apps", "web", "package.json"), "{}\n");
await writeFile(join(root, "apps", "web", "app.secret"), "private");
const submitted: BuildRequest[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
const build = init!.json as BuildRequest;
submitted.push(build);
return { state: "queued", id: build.id } as T;
}
if (path.includes("/events")) return [] as T;
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.endsWith("/result")) return { reference: "image:built" } as T;
throw new Error(path);
};
await buildServices(
"shop",
{
services: {
selected: { build: "auto", "x-kuber-build-context": "apps/web" },
root: { build: "auto" },
},
},
root,
undefined,
{ request, sleep: async () => {}, pollIntervalMs: 0 },
);
expect(submitted.map(({ spec }) => spec.context)).toEqual([
"apps/web",
".",
]);
expect(
submitted.every(
({ spec }) =>
(spec as BuildRequest["spec"] & { builder?: string }).builder ===
"buildpacks",
),
).toBe(true);
const autoSnapshot = await enumerateWorkspace(root, "auto");
expect(autoSnapshot.manifest.files.map(({ path }) => path)).toContain(
"apps/web/package.json",
);
expect(autoSnapshot.manifest.files.map(({ path }) => path)).not.toContain(
"apps/web/app.secret",
);
});
test("auto rejects traversal and absolute selected contexts", async () => {
const root = await mkdtemp(
join(tmpdir(), "kuber-auto-context-validation-"),
);
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
const request: ApiRequester = async <T>() => ({ ready: true }) as T;
for (const context of [
"../outside",
"apps/../../outside",
"/tmp/outside",
"C:/outside",
]) {
await expect(
buildServices(
"shop",
{
services: {
web: { build: "auto", "x-kuber-build-context": context },
},
},
root,
undefined,
{ request, snapshot: emptySnapshot() },
),
).rejects.toThrow(/x-kuber-build-context/);
}
});
test("uploads distinct auto and Dockerfile snapshots and groups only matching builders", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-mixed-build-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
await writeFile(join(root, ".gitignore"), "tracked.secret\n.env*\n");
await writeFile(join(root, "tracked.secret"), "private");
expect(
await Bun.spawn(["git", "-C", root, "add", "-f", "tracked.secret"])
.exited,
).toBe(0);
await writeFile(join(root, "Dockerfile"), "FROM scratch\n");
await writeFile(join(root, ".dockerignore"), "app.txt\n");
await writeFile(join(root, "app.txt"), "app");
const regular = await enumerateWorkspace(root);
const auto = await enumerateWorkspace(root, "auto");
const negotiations: string[] = [];
const uploadedManifests = new Map<string, string>();
const submissions: BuildRequest[] = [];
const progress: Array<[number, number]> = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
const workspace = (init!.json as { workspace: string }).workspace;
negotiations.push(workspace);
return {
workspace,
missing:
negotiations.filter((digest) => digest === workspace).length === 1
? [workspace]
: [],
ready:
negotiations.filter((digest) => digest === workspace).length !== 1,
} as T;
}
if (path.includes("/uploads")) {
if (init?.method === "PATCH") {
const digest = decodeURIComponent(path.split("/")[2]!);
uploadedManifests.set(
digest,
Buffer.from(init.body as Uint8Array).toString(),
);
return { offset: (init.body as Uint8Array).byteLength } as T;
}
return { offset: 0, complete: false } as T;
}
if (path === "/builds") {
const build = init?.json as BuildRequest;
submissions.push(build);
return { state: "queued", id: build.id } as T;
}
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.includes("/events")) return [] as T;
if (path.endsWith("/result")) {
const build = submissions.find(({ id }) => path.includes(id))!;
return {
reference: `image:${build.service}`,
references: Object.fromEntries(
(build.destinations ?? []).map(({ service }) => [
service,
`image:${service}`,
]),
),
} as T;
}
throw new Error(path);
};
const result = await buildServices(
"shop",
{
services: {
docker: { build: "." },
auto: { build: "auto" },
second: { build: "auto" },
},
},
root,
{
upload: (uploaded, total) => {
progress.push([uploaded, total]);
},
},
{ request, sleep: async () => {}, pollIntervalMs: 0 },
);
expect(negotiations).toEqual([
regular.digest,
auto.digest,
regular.digest,
auto.digest,
]);
const totalBytes =
serializeWorkspaceManifest(regular.manifest).byteLength +
serializeWorkspaceManifest(auto.manifest).byteLength;
expect(progress[0]).toEqual([0, totalBytes]);
expect(progress.at(-1)).toEqual([totalBytes, totalBytes]);
expect(
progress.every(
([uploaded, total]) => uploaded <= total && total === totalBytes,
),
).toBe(true);
expect([...uploadedManifests.keys()]).toEqual([
regular.digest,
auto.digest,
]);
expect(
JSON.parse(uploadedManifests.get(regular.digest)!).files.map(
(file: { path: string }) => file.path,
),
).toContain("tracked.secret");
expect(
JSON.parse(uploadedManifests.get(auto.digest)!).files.map(
(file: { path: string }) => file.path,
),
).not.toContain("tracked.secret");
expect(submissions).toHaveLength(2);
expect(submissions[0]?.spec).toMatchObject({
workspace: regular.digest,
context: ".",
});
expect(submissions[0]?.spec).not.toHaveProperty("builder");
expect(submissions[1]?.spec).toMatchObject({
workspace: auto.digest,
context: ".",
builder: "buildpacks",
});
expect(submissions[1]?.spec).not.toHaveProperty("dockerfile");
expect(submissions[1]?.destinations).toEqual([
{ service: "second", image: expect.any(String) },
]);
expect(result.images).toEqual({
docker: "image:docker",
auto: "image:auto",
second: "image:second",
});
});
test("negotiates and uploads the manifest through resumable blob routes", async () => {
const snapshot = emptySnapshot();
const controller = new AbortController();
const calls: Array<{
path: string;
init?: ApiRequestInit;
options?: ApiRequestOptions;
}> = [];
let negotiations = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
options?: ApiRequestOptions,
) => {
calls.push({ path, init, options });
if (path === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: [snapshot.digest],
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
if (init?.method === "POST" && !path.endsWith("/complete"))
return { offset: 0, complete: false } as T;
if (init?.method === "PATCH")
return { offset: (init.body as Uint8Array).byteLength } as T;
return { complete: true } as T;
};
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
undefined,
undefined,
controller.signal,
);
expect(
calls.map(({ path, init }) => [init?.method ?? "GET", path]),
).toEqual([
["POST", "/snapshots/negotiate"],
["POST", `/blobs/${encodeURIComponent(snapshot.digest)}/uploads`],
["PATCH", `/blobs/${encodeURIComponent(snapshot.digest)}/uploads`],
[
"POST",
`/blobs/${encodeURIComponent(snapshot.digest)}/uploads/complete`,
],
["POST", "/snapshots/negotiate"],
]);
expect(new Headers(calls[2]!.init?.headers).get("upload-offset")).toBe("0");
expect(calls.filter(({ path }) => path === "/snapshots/negotiate")).toEqual(
[
expect.objectContaining({
init: expect.objectContaining({ signal: controller.signal }),
options: { timeoutMs: 300_000 },
}),
expect.objectContaining({
init: expect.objectContaining({ signal: controller.signal }),
options: { timeoutMs: 300_000 },
}),
],
);
expect(calls.every(({ init }) => init?.signal === controller.signal)).toBe(
true,
);
});
test("uses project-scoped URLs for every resumable blob upload endpoint", async () => {
const snapshot = emptySnapshot();
const project = "shop & staging";
const uploadPath = `/blobs/${encodeURIComponent(snapshot.digest)}/uploads`;
const calls: Array<{ method: string; path: string }> = [];
let negotiations = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
const url = new URL(path, "https://kuber.test");
const method = init?.method ?? "GET";
calls.push({ method, path });
if (url.pathname === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: [snapshot.digest],
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
if (url.pathname === uploadPath && method === "POST")
return { offset: 0, complete: false } as T;
if (url.pathname === uploadPath && method === "PATCH")
return { offset: (init!.body as Uint8Array).byteLength } as T;
if (url.pathname === `${uploadPath}/complete` && method === "POST")
return { complete: true } as T;
throw new Error(`Unexpected request ${method} ${path}`);
};
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
undefined,
project,
);
const projectQuery = `?project=${encodeURIComponent(project)}`;
expect(calls).toEqual([
{ method: "POST", path: "/snapshots/negotiate" },
{ method: "POST", path: `${uploadPath}${projectQuery}` },
{ method: "PATCH", path: `${uploadPath}${projectQuery}` },
{ method: "POST", path: `${uploadPath}/complete${projectQuery}` },
{ method: "POST", path: "/snapshots/negotiate" },
]);
});
describe("TaskScheduler", () => {
test("uses the production request limits", () => {
expect(MAX_CONCURRENT_REQUESTS).toBe(20);
expect(MAX_REQUESTS_PER_SECOND).toBe(40);
});
test("limits concurrent in-flight operations", async () => {
let maxInflight = 0;
let currentInflight = 0;
const scheduler = new TaskScheduler({
maxInflight: 5,
maxPerSecond: 100,
});
const tasks = Array.from({ length: 10 }, () =>
scheduler.run(async () => {
currentInflight++;
maxInflight = Math.max(maxInflight, currentInflight);
await Bun.sleep(10);
currentInflight--;
}),
);
await Promise.all(tasks);
expect(maxInflight).toBeLessThanOrEqual(5);
});
test("rate limits request starts to maxPerSecond", async () => {
const requestStarts: number[] = [];
let currentTime = 0;
const clock = { now: () => currentTime };
const sleep = async (ms: number) => {
currentTime += ms;
};
const scheduler = new TaskScheduler({
maxInflight: 100,
maxPerSecond: 4,
clock,
sleep,
});
const tasks = Array.from({ length: 8 }, () =>
scheduler.run(async () => {
requestStarts.push(currentTime);
}),
);
await Promise.all(tasks);
for (const start of requestStarts) {
expect(
requestStarts.filter(
(candidate) => candidate >= start && candidate < start + 1000,
).length,
).toBeLessThanOrEqual(4);
}
});
});
describe("concurrent blob uploads", () => {
test("uploads multiple blobs concurrently", async () => {
const snapshot = createSnapshotWithBlobs(5, 1024);
let negotiations = 0;
let inflight = 0;
let maxInflight = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: snapshot.blobs.map((blob) => blob.digest),
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
inflight += 1;
maxInflight = Math.max(maxInflight, inflight);
await Bun.sleep(5);
inflight -= 1;
if (init?.method === "POST" && !path.endsWith("/complete"))
return { offset: 0, complete: false } as T;
if (init?.method === "PATCH") {
return { offset: (init.body as Uint8Array).byteLength } as T;
}
return { complete: true } as T;
};
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
new TaskScheduler({ maxInflight: 5, maxPerSecond: 100 }),
);
expect(maxInflight).toBeGreaterThan(1);
expect(negotiations).toBe(2);
});
test("completes active blobs before initializing the full backlog", async () => {
const snapshot = createSnapshotWithBlobs(12, 1024);
const methods: string[] = [];
let negotiations = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: snapshot.blobs.map((blob) => blob.digest),
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
methods.push(
path.endsWith("/complete") ? "complete" : (init?.method ?? "GET"),
);
if (init?.method === "POST" && !path.endsWith("/complete"))
return { offset: 0, complete: false } as T;
if (init?.method === "PATCH")
return { offset: (init.body as Uint8Array).byteLength } as T;
return { complete: true } as T;
};
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
new TaskScheduler({ maxInflight: 4, maxPerSecond: 100 }),
);
expect(methods.indexOf("complete")).toBeLessThan(
methods.lastIndexOf("POST"),
);
});
test("inflight never exceeds configured maximum", async () => {
const snapshot = createSnapshotWithBlobs(30, 1024);
let currentInflight = 0;
let maxObservedInflight = 0;
const maxInflight = 10;
let negotiations = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: snapshot.blobs.map((blob) => blob.digest),
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
currentInflight += 1;
maxObservedInflight = Math.max(maxObservedInflight, currentInflight);
await Bun.sleep(2);
currentInflight -= 1;
if (init?.method === "POST" && !path.endsWith("/complete"))
return { offset: 0, complete: false } as T;
if (init?.method === "PATCH") {
return { offset: (init.body as Uint8Array).byteLength } as T;
}
return { complete: true } as T;
};
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
new TaskScheduler({ maxInflight, maxPerSecond: 1000 }),
);
expect(maxObservedInflight).toBeLessThanOrEqual(maxInflight);
});
test("request starts are rate-limited to maxPerSecond", async () => {
const snapshot = createSnapshotWithBlobs(8, 1024);
const requestStarts: Array<{ path: string; time: number }> = [];
let currentTime = 0;
const clock = { now: () => currentTime };
const sleep = async (ms: number) => {
currentTime += ms;
};
let negotiations = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
requestStarts.push({ path, time: currentTime });
if (path === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: snapshot.blobs.map((blob) => blob.digest),
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
if (init?.method === "POST" && !path.endsWith("/complete")) {
return { offset: 0, complete: false } as T;
}
if (init?.method === "PATCH") {
return { offset: (init.body as Uint8Array).byteLength } as T;
}
if (path.endsWith("/complete")) {
return { complete: true } as T;
}
return { complete: true } as T;
};
const maxPerSecond = 4;
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
new TaskScheduler({
maxInflight: 100,
maxPerSecond,
clock,
sleep,
}),
);
for (const { time } of requestStarts) {
expect(
requestStarts.filter(
({ time: candidate }) =>
candidate >= time && candidate < time + 1000,
).length,
).toBeLessThanOrEqual(maxPerSecond);
}
});
test("chunks of one blob remain ordered during concurrent uploads", async () => {
const blobSize = 25 * 1024 * 1024;
const snapshot = createSnapshotWithBlobs(3, blobSize);
const chunkOrders = new Map<string, number[]>();
let negotiations = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") {
negotiations += 1;
return (
negotiations === 1
? {
workspace: snapshot.digest,
missing: snapshot.blobs.map((blob) => blob.digest),
ready: false,
}
: { workspace: snapshot.digest, missing: [], ready: true }
) as T;
}
if (init?.method === "POST" && !path.endsWith("/complete")) {
return { offset: 0, complete: false } as T;
}
if (init?.method === "PATCH") {
const digest = decodeURIComponent(path.split("/")[2]!);
const offset = Number(
new Headers(init.headers).get("upload-offset") ?? "0",
);
const chunkIndex = Math.floor(offset / (8 * 1024 * 1024));
if (!chunkOrders.has(digest)) chunkOrders.set(digest, []);
chunkOrders.get(digest)!.push(chunkIndex);
return { offset: offset + (init.body as Uint8Array).byteLength } as T;
}
return { complete: true } as T;
};
await uploadWorkspaceSnapshot(
snapshot,
request,
undefined,
new TaskScheduler({ maxInflight: 3, maxPerSecond: 100 }),
);
for (const [, chunks] of chunkOrders) {
for (let i = 1; i < chunks.length; i++) {
expect(chunks[i]!).toBeGreaterThan(chunks[i - 1]!);
}
}
});
});
test("submits, reconciles, reports logs, and returns the server image reference", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-build-api-"));
directories.push(root);
const git = Bun.spawn(["git", "init", "-q", root]);
expect(await git.exited).toBe(0);
const snapshot = emptySnapshot();
const calls: Array<{
path: string;
init?: ApiRequestInit;
options?: ApiRequestOptions;
}> = [];
let submitted: BuildRequest | undefined;
const output: string[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
options?: ApiRequestOptions,
) => {
calls.push({ path, init, options });
if (path === "/snapshots/negotiate")
return { workspace: snapshot.digest, missing: [], ready: true } as T;
if (path === "/builds") {
submitted = init?.json as BuildRequest;
return {
version: 1,
id: submitted.id,
state: "queued",
createdAt: "2026-01-01T00:00:00Z",
} as T;
}
if (path.includes("/events"))
return [
{
type: "log",
id: submitted!.id,
sequence: 1,
message: "build log\n",
},
] as T;
if (path.endsWith("/reconcile"))
return {
version: 1,
id: submitted!.id,
state: "succeeded",
createdAt: "2026-01-01T00:00:00Z",
digest: `sha256:${"a".repeat(64)}`,
} as T;
if (path.endsWith("/result"))
return {
image: "registry.server/kuber/shop-web",
digest: `sha256:${"a".repeat(64)}`,
reference: `registry.server/kuber/shop-web@sha256:${"a".repeat(64)}`,
} as T;
throw new Error(`Unexpected request ${path}`);
};
const result = await buildServices(
"shop",
{
services: { web: { build: { context: ".", args: { MODE: "prod" } } } },
},
root,
{
progress: (message) => {
output.push(message);
},
},
{ request, snapshot, sleep: async () => {}, pollIntervalMs: 0 },
);
expect(submitted?.spec).toMatchObject({
architecture: "arm64",
context: ".",
buildArgs: ["MODE=prod"],
workspace: snapshot.digest,
});
expect(result).toEqual({
built: ["web"],
changed: ["web"],
images: {
web: `registry.server/kuber/shop-web@sha256:${"a".repeat(64)}`,
},
});
expect(output).toContain("web: build log");
expect(calls.find(({ path }) => path === "/builds")?.options).toEqual({
timeoutMs: 300_000,
});
expect(
calls
.filter(
({ path }) => path.includes("/events") || path.endsWith("/reconcile"),
)
.map(({ options }) => options),
).toEqual([
{ timeoutMs: 300_000 },
{ timeoutMs: 300_000 },
{ timeoutMs: 300_000 },
]);
expect(calls.some(({ path }) => path.endsWith("/result"))).toBe(true);
});
test("streams one physical build, resumes cursor, and never polls JSON while SSE is active", async () => {
const snapshot = emptySnapshot();
const calls: string[] = [];
const cursors: number[] = [];
const output: string[] = [];
const request: ApiRequester = async <T>(path: string) => {
calls.push(path);
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") return { state: "queued" } as T;
if (path.endsWith("/result"))
return {
reference: "image:web",
references: { worker: "image:worker" },
} as T;
throw new Error(`Unexpected JSON request ${path}`);
};
const result = await buildServices(
"shop",
{
services: { web: { build: "." }, worker: { build: "." } },
},
process.cwd(),
{
stream: new Writable({
write(chunk, _, done) {
output.push(String(chunk));
done();
},
}),
},
{
request,
snapshot,
pollIntervalMs: 0,
sleep: async () => {},
streamEvents: async function* (_path, after) {
cursors.push(after);
if (after === 0) {
yield {
id: 1,
event: {
type: "log",
sequence: 1,
id: "build",
message: "once\n",
},
};
yield {
event: {
type: "status",
status: {
version: 1,
id: "build",
state: "running",
createdAt: "now",
},
},
};
} else {
yield {
id: 1,
event: {
type: "log",
sequence: 1,
id: "build",
message: "once\n",
},
};
yield {
event: {
type: "status",
status: {
version: 1,
id: "build",
state: "succeeded",
createdAt: "now",
},
},
};
}
},
},
);
expect(result.images).toEqual({ web: "image:web", worker: "image:worker" });
expect(cursors).toEqual([0, 1]);
expect(output).toEqual(["[web] once\n"]);
expect(
calls.filter(
(path) => path.includes("/events") || path.includes("/reconcile"),
),
).toEqual([]);
});
test("falls back to JSON events/reconcile only after SSE is unsupported", async () => {
const calls: string[] = [];
const request: ApiRequester = async <T>(path: string) => {
calls.push(path);
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") return { state: "queued" } as T;
if (path.includes("/events"))
return [{ type: "log", sequence: 1, message: "fallback\n" }] as T;
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.endsWith("/result")) return { reference: "image:web" } as T;
throw new Error(path);
};
expect(
(
await buildServices(
"shop",
{ services: { web: { build: "." } } },
process.cwd(),
undefined,
{
request,
snapshot: emptySnapshot(),
streamEvents: async function* () {
yield* [];
throw new ApiStreamUnsupportedError();
},
},
)
).images,
).toEqual({ web: "image:web" });
expect(
calls.filter(
(path) => path.includes("/events") || path.includes("/reconcile"),
),
).toEqual([
expect.stringContaining("/events?after=0"),
expect.stringContaining("/reconcile"),
expect.stringContaining("/events?after=1"),
]);
});
test("aborts an active stream and does not issue JSON polling requests", async () => {
const controller = new AbortController();
const reason = new DOMException("Stopped", "AbortError");
let listening!: () => void;
const ready = new Promise<void>((resolve) => {
listening = resolve;
});
const calls: string[] = [];
const request: ApiRequester = async <T>(path: string) => {
calls.push(path);
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") return { state: "queued" } as T;
throw new Error(path);
};
const result = buildServices(
"shop",
{ services: { web: { build: "." } } },
process.cwd(),
undefined,
{
request,
snapshot: emptySnapshot(),
signal: controller.signal,
streamEvents: async function* (_path, _after, signal) {
yield* [];
listening();
await new Promise<void>((_resolve, reject) =>
signal?.addEventListener("abort", () => reject(signal.reason), {
once: true,
}),
);
},
},
);
await ready;
controller.abort(reason);
await expect(result).rejects.toBe(reason);
expect(calls).toEqual(["/snapshots/negotiate", "/builds"]);
});
test("builds equivalent services once and attributes events and image results to each", async () => {
const snapshot = emptySnapshot();
const submitted: BuildRequest[] = [];
const messages = new Map<string, string[]>();
const settled: string[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
submitted.push(init!.json as BuildRequest);
return { state: "queued" } as T;
}
if (path.includes("/events"))
return [
{ type: "status", status: { state: "running", phase: "running" } },
{ type: "log", sequence: 1, message: "shared log\n" },
] as T;
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.endsWith("/result"))
return {
reference: "registry/kuber/shop-web@sha256:abc",
references: { worker: "registry/kuber/shop-worker@sha256:abc" },
} as T;
throw new Error(path);
};
const result = await buildServices(
"shop",
{
services: { web: { build: "." }, worker: { build: { context: "." } } },
},
process.cwd(),
{
service: (name) => {
const output: string[] = [];
messages.set(name, output);
return {
progress: (message) => {
output.push(message);
},
};
},
settled: (name) => {
settled.push(name);
},
},
{ request, snapshot, sleep: async () => {}, pollIntervalMs: 0 },
);
expect(submitted).toHaveLength(1);
expect(submitted[0]?.destinations).toEqual([
{
service: "worker",
image: "registry.neko-piranha.ts.net/kuber/shop-worker:latest",
},
]);
expect(result).toEqual({
built: ["web", "worker"],
changed: ["web", "worker"],
images: {
web: "registry/kuber/shop-web@sha256:abc",
worker: "registry/kuber/shop-worker@sha256:abc",
},
});
expect(settled).toEqual(["web", "worker"]);
for (const name of ["web", "worker"])
expect(messages.get(name)?.join(" ")).toContain("shared log");
});
test("routes overlapping build logs only to their destination children", async () => {
const snapshot = emptySnapshot();
const submissions = new Map<string, BuildRequest>();
const output = new Map<string, string[]>();
let started = 0;
let release!: () => void;
const bothStarted = new Promise<void>((resolve) => {
release = resolve;
});
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
const build = init!.json as BuildRequest;
submissions.set(build.id, build);
if (++started === 2) release();
return { state: "queued" } as T;
}
const id = path.split("/")[2]!;
const build = submissions.get(id)!;
if (path.includes("/events"))
return path.includes("after=0")
? ([
{
type: "status",
status: { state: "running", phase: "running" },
},
{ type: "log", sequence: 1, message: `${build.service} log\n` },
] as T)
: ([] as T);
if (path.endsWith("/reconcile")) {
await bothStarted;
return { state: "succeeded" } as T;
}
if (path.endsWith("/result"))
return {
reference: `image:${build.service}`,
references: { worker: "image:worker" },
} as T;
throw new Error(path);
};
const result = await buildServices(
"shop",
{
services: {
web: { build: "." },
worker: { build: "." },
admin: { build: { context: ".", args: { ROLE: "admin" } } },
},
},
process.cwd(),
{
service: (name) => {
const lines: string[] = [];
output.set(name, lines);
return {
progress: (message) => {
lines.push(message);
},
stream: new Writable({
write(chunk, _encoding, done) {
lines.push(String(chunk));
done();
},
}),
};
},
},
{
request,
snapshot,
buildConcurrency: 2,
pollIntervalMs: 0,
sleep: async () => {},
},
);
expect(submissions.size).toBe(2);
expect(
[...submissions.values()]
.find(({ service }) => service === "web")
?.destinations?.map(({ service }) => service),
).toEqual(["worker"]);
expect(result.images).toEqual({
web: "image:web",
worker: "image:worker",
admin: "image:admin",
});
for (const name of ["web", "worker"]) {
expect(
output.get(name)?.filter((line) => line === "web log\n"),
).toHaveLength(1);
expect(output.get(name)?.join("")).not.toContain("admin log");
expect(output.get(name)).toContain("Build running");
}
expect(
output.get("admin")?.filter((line) => line === "admin log\n"),
).toHaveLength(1);
expect(output.get("admin")?.join("")).not.toContain("web log");
});
test("reports a grouped build once to a shared reporter", async () => {
const output: string[] = [];
const request: ApiRequester = async <T>(path: string) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") return { state: "queued" } as T;
if (path.includes("/events"))
return path.includes("after=0")
? ([
{ type: "log", sequence: 1, message: "one physical build\n" },
] as T)
: ([] as T);
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.endsWith("/result"))
return {
reference: "image:web",
references: { worker: "image:worker" },
} as T;
throw new Error(path);
};
await buildServices(
"shop",
{ services: { web: { build: "." }, worker: { build: "." } } },
process.cwd(),
{
stream: new Writable({
write(chunk, _encoding, done) {
output.push(String(chunk));
done();
},
}),
},
{
request,
snapshot: emptySnapshot(),
sleep: async () => {},
pollIntervalMs: 0,
},
);
expect(output).toEqual(["[web] one physical build\n"]);
});
test("does not merge builds differing in context, Dockerfile, target, or build arguments", async () => {
const snapshot = emptySnapshot();
const submitted: BuildRequest[] = [];
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
submitted.push(init!.json as BuildRequest);
return { state: "queued" } as T;
}
if (path.includes("/events")) return [] as T;
if (path.endsWith("/reconcile")) return { state: "succeeded" } as T;
if (path.endsWith("/result"))
return { reference: "image@sha256:abc" } as T;
throw new Error(path);
};
await buildServices(
"shop",
{
services: {
base: { build: "." },
context: { build: { context: "nested" } },
dockerfile: { build: { context: ".", dockerfile: "Otherfile" } },
target: { build: { context: ".", target: "test" } },
args: { build: { context: ".", args: { MODE: "test" } } },
},
},
process.cwd(),
undefined,
{ request, snapshot },
);
expect(submitted).toHaveLength(5);
expect(submitted.every((build) => !build.destinations)).toBe(true);
});
test("marks every service failed when a shared BuildKit job fails", async () => {
const snapshot = emptySnapshot();
const settled: Array<[string, unknown]> = [];
let submissions = 0;
const request: ApiRequester = async <T>(path: string) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
submissions++;
return { state: "queued" } as T;
}
if (path.includes("/events"))
return [{ type: "log", sequence: 1, message: "build failed\n" }] as T;
if (path.endsWith("/reconcile"))
return { state: "failed", error: "export failed" } as T;
throw new Error(path);
};
await expect(
buildServices(
"shop",
{
services: {
web: { build: "." },
worker: { build: "." },
},
},
process.cwd(),
{
settled: (name, error) => {
settled.push([name, error]);
},
},
{ request, snapshot, sleep: async () => {}, pollIntervalMs: 0 },
),
).rejects.toThrow("export failed\nbuild failed");
expect(submissions).toBe(1);
expect(settled.map(([name]) => name)).toEqual(["web", "worker"]);
expect(settled.every(([, error]) => error instanceof Error)).toBe(true);
});
test("cancels a shared build once and does not start queued groups", async () => {
const snapshot = emptySnapshot();
const controller = new AbortController();
const reason = new DOMException("Cancelled", "AbortError");
const started: BuildRequest[] = [];
const settled: string[] = [];
let ready!: () => void;
const submitted = new Promise<void>((resolve) => {
ready = resolve;
});
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
started.push(init!.json as BuildRequest);
ready();
return await new Promise<T>((_resolve, reject) => {
init?.signal?.addEventListener(
"abort",
() => reject(init.signal!.reason),
{ once: true },
);
});
}
throw new Error(path);
};
const result = buildServices(
"shop",
{
services: {
web: { build: "." },
worker: { build: "." },
different: { build: { context: ".", args: { MODE: "different" } } },
},
},
process.cwd(),
{
settled: (name) => {
settled.push(name);
},
},
{ request, snapshot, buildConcurrency: 1, signal: controller.signal },
);
await submitted;
controller.abort(reason);
await expect(result).rejects.toBe(reason);
expect(started).toHaveLength(1);
expect(started[0]!.destinations?.map(({ service }) => service)).toEqual([
"worker",
]);
expect(settled).toEqual(["web", "worker"]);
});
test("bounds overlapping builds, attributes logs, and returns images in compose order", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-build-api-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
const snapshot = emptySnapshot();
const started: string[] = [];
const startWaiters = new Map<number, () => void>();
const waitForStart = (count: number) =>
new Promise<void>((resolve) => {
startWaiters.set(count, resolve);
});
const firstTwo = waitForStart(2);
const third = waitForStart(3);
const fourth = waitForStart(4);
const released = new Map<string, () => void>();
const gates = new Map<string, Promise<void>>();
const ids = new Map<string, string>();
const output: string[] = [];
const stream = new Writable({
write(chunk, _encoding, done) {
output.push(String(chunk));
done();
},
});
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate")
return { workspace: snapshot.digest, missing: [], ready: true } as T;
if (path === "/builds") {
const build = init?.json as BuildRequest;
started.push(build.service);
startWaiters.get(started.length)?.();
ids.set(build.id, build.service);
gates.set(
build.service,
new Promise<void>((resolve) => released.set(build.service, resolve)),
);
return { id: build.id, state: "queued" } as T;
}
const id = path.split("/")[2]!;
const service = ids.get(id)!;
if (path.includes("/events"))
return [{ type: "log", sequence: 1, message: `log ${service}\n` }] as T;
if (path.endsWith("/reconcile")) {
await gates.get(service);
return { id, state: "succeeded" } as T;
}
if (path.endsWith("/result"))
return { reference: `image:${service}` } as T;
throw new Error(`Unexpected request ${path}`);
};
const run = buildServices(
"shop",
{
services: Object.fromEntries(
["one", "two", "three", "four"].map((name) => [
name,
{ build: { context: ".", args: { SERVICE: name } } },
]),
),
},
root,
{ stream },
{ request, snapshot, buildConcurrency: 2, sleep: async () => {} },
);
try {
await firstTwo;
expect(started).toEqual(["one", "two"]);
released.get("two")!();
await third;
expect(started).toEqual(["one", "two", "three"]);
released.get("three")!();
await fourth;
released.get("four")!();
released.get("one")!();
expect(await run).toEqual({
built: ["one", "two", "three", "four"],
changed: ["one", "two", "three", "four"],
images: {
one: "image:one",
two: "image:two",
three: "image:three",
four: "image:four",
},
});
expect(new Set(ids.keys()).size).toBe(4);
expect(output).toEqual(
expect.arrayContaining([
"[one] log one\n",
"[two] log two\n",
"[three] log three\n",
"[four] log four\n",
]),
);
} finally {
for (const release of released.values()) release();
await run.catch(() => {});
}
});
test("keeps per-image logs and lifecycle updates separate, including a failed build", async () => {
const snapshot = emptySnapshot();
const outputs = new Map<string, string[]>();
const ids = new Map<string, string>();
const settled: Array<[string, unknown]> = [];
const statuses = ["creating", "starting", "running", "done"] as const;
let polls = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
const build = init!.json as BuildRequest;
ids.set(build.id, build.service);
return { state: "queued" } as T;
}
const id = path.split("/")[2]!;
const service = ids.get(id)!;
if (path.includes("/events")) {
const phase = statuses[Math.min(polls, 3)]!;
return [
{
type: "status",
status: {
state:
phase === "done"
? service === "server"
? "failed"
: "succeeded"
: phase === "running"
? "running"
: "queued",
phase,
error:
service === "server" && phase === "done"
? "stack trace"
: undefined,
},
},
{ type: "log", sequence: polls + 1, message: `${service} log\n` },
] as T;
}
if (path.endsWith("/reconcile")) {
const phase = statuses[Math.min(polls++, 3)]!;
return {
state:
phase === "done"
? service === "server"
? "failed"
: "succeeded"
: phase === "running"
? "running"
: "queued",
phase,
error: "stack trace",
} as T;
}
if (path.endsWith("/result"))
return { reference: `image:${service}` } as T;
throw new Error(path);
};
await expect(
buildServices(
"shop",
{
services: {
client: { build: "." },
server: { build: { context: ".", args: { SERVICE: "server" } } },
},
},
process.cwd(),
{
service: (name) => {
const output: string[] = [];
outputs.set(name, output);
return {
progress: (message) => {
output.push(message);
},
stream: new Writable({
write(chunk, _encoding, done) {
output.push(String(chunk));
done();
},
}),
};
},
settled: (name, error) => {
settled.push([name, error]);
},
},
{
request,
snapshot,
sleep: async () => {},
pollIntervalMs: 0,
buildConcurrency: 1,
},
),
).rejects.toThrow(
"Build failed for service server: stack trace\nserver log",
);
expect(outputs.get("client")).toEqual(
expect.arrayContaining([
"Build queued",
"Build creating",
"Build starting",
"Build running",
"Build done",
"client log\n",
]),
);
expect(outputs.get("client")?.join("")).not.toContain("server log");
expect(outputs.get("server")?.join("")).toContain("server log");
expect(settled.map(([name]) => name)).toEqual(["client", "server"]);
expect(settled[1]?.[1]).toBeInstanceOf(Error);
expect(ids.size).toBe(2); // Different build arguments require separate jobs.
});
test("drains active builds and does not start queued services after failure", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-build-api-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
const snapshot = emptySnapshot();
const started: string[] = [];
let bothReady!: () => void;
const bothStarted = new Promise<void>((resolve) => {
bothReady = resolve;
});
let fail!: (error: Error) => void;
let finish!: () => void;
const failing = new Promise<never>((_resolve, reject) => {
fail = reject;
});
const active = new Promise<void>((resolve) => {
finish = resolve;
});
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
const service = (init!.json as BuildRequest).service;
started.push(service);
if (started.length === 2) bothReady();
if (service === "one") return await failing;
await active;
return { state: "succeeded" } as T;
}
if (path.includes("/events")) return [] as T;
if (path.endsWith("/result")) return { reference: "image:two" } as T;
throw new Error(`Unexpected request ${path}`);
};
const run = buildServices(
"shop",
{
services: {
one: { build: { context: ".", args: { SERVICE: "one" } } },
two: { build: { context: ".", args: { SERVICE: "two" } } },
three: { build: { context: ".", args: { SERVICE: "three" } } },
},
},
root,
undefined,
{ request, snapshot, buildConcurrency: 2 },
);
const outcome = run.then(
() => undefined,
(error: unknown) => error,
);
await bothStarted;
fail(new Error("first build failed"));
await Promise.resolve();
expect(started).toEqual(["one", "two"]);
finish();
expect(await outcome).toMatchObject({ message: "first build failed" });
expect(started).toEqual(["one", "two"]);
});
test("passes cancellation to active build requests and skips queued builds", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-build-api-"));
directories.push(root);
expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0);
const snapshot = emptySnapshot();
const controller = new AbortController();
const cancelled = new DOMException("Cancelled", "AbortError");
const started: string[] = [];
let ready!: () => void;
const bothStarted = new Promise<void>((resolve) => {
ready = resolve;
});
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
) => {
if (path === "/snapshots/negotiate") return { ready: true } as T;
if (path === "/builds") {
started.push((init!.json as BuildRequest).service);
expect(init?.signal).toBe(controller.signal);
if (started.length === 2) ready();
return await new Promise<T>((_resolve, reject) => {
init?.signal?.addEventListener(
"abort",
() => reject(init.signal!.reason),
{ once: true },
);
});
}
throw new Error(`Unexpected request ${path}`);
};
const outcome = buildServices(
"shop",
{
services: {
one: { build: { context: ".", args: { SERVICE: "one" } } },
two: { build: { context: ".", args: { SERVICE: "two" } } },
three: { build: { context: ".", args: { SERVICE: "three" } } },
},
},
root,
undefined,
{ request, snapshot, buildConcurrency: 2, signal: controller.signal },
).then(
() => undefined,
(error: unknown) => error,
);
await bothStarted;
controller.abort(cancelled);
expect(await outcome).toBe(cancelled);
expect(started).toEqual(["one", "two"]);
});
test("retries a timed out poll request with the build request timeout", async () => {
const root = await mkdtemp(join(tmpdir(), "kuber-build-api-"));
directories.push(root);
const git = Bun.spawn(["git", "init", "-q", root]);
expect(await git.exited).toBe(0);
const snapshot = emptySnapshot();
const pollCalls: Array<{ path: string; options?: ApiRequestOptions }> = [];
let buildId = "";
let eventRequests = 0;
const request: ApiRequester = async <T>(
path: string,
init?: ApiRequestInit,
options?: ApiRequestOptions,
) => {
if (path === "/snapshots/negotiate")
return { workspace: snapshot.digest, missing: [], ready: true } as T;
if (path === "/builds") {
const buildRequest = init?.json as BuildRequest | undefined;
if (!buildRequest) throw new Error("Expected build request");
buildId = buildRequest.id;
return {
version: 1,
id: buildId,
state: "queued",
createdAt: "2026-01-01T00:00:00Z",
} as T;
}
if (path.includes("/events")) {
pollCalls.push({ path, options });
eventRequests += 1;
if (eventRequests === 1)
throw new DOMException("request timed out", "TimeoutError");
return [] as T;
}
if (path.endsWith("/reconcile")) {
pollCalls.push({ path, options });
return {
version: 1,
id: buildId,
state: "succeeded",
createdAt: "2026-01-01T00:00:00Z",
digest: `sha256:${"a".repeat(64)}`,
} as T;
}
if (path.endsWith("/result"))
return {
image: "registry.server/kuber/shop-web",
digest: `sha256:${"a".repeat(64)}`,
reference: `registry.server/kuber/shop-web@sha256:${"a".repeat(64)}`,
} as T;
throw new Error(`Unexpected request ${path}`);
};
await buildServices(
"shop",
{ services: { web: { build: "." } } },
root,
undefined,
{ request, snapshot, sleep: async () => {}, pollIntervalMs: 0 },
);
expect(eventRequests).toBe(3);
expect(pollCalls).toEqual([
expect.objectContaining({ options: { timeoutMs: 300_000 } }),
expect.objectContaining({ options: { timeoutMs: 300_000 } }),
expect.objectContaining({ options: { timeoutMs: 300_000 } }),
expect.objectContaining({ options: { timeoutMs: 300_000 } }),
]);
});
test("no-build image resolution uses only the resolve route", async () => {
const calls: string[] = [];
const images = await resolveBuildImages(
"shop",
{ services: { web: { build: "." }, cache: { image: "redis" } } },
{
request: async <T>(path: string, init?: ApiRequestInit) => {
calls.push(`${init?.method}:${path}:${JSON.stringify(init?.json)}`);
return { reference: "registry/web@sha256:immutable" } as T;
},
},
);
expect(images).toEqual({ web: "registry/web@sha256:immutable" });
expect(calls).toEqual([
'POST:/images/resolve:{"project":"shop","service":"web"}',
]);
});
test("resolves build images concurrently with bounded concurrency and stable order", async () => {
const services = ["one", "two", "three", "four", "five"];
const started: string[] = [];
const pending: Array<() => void> = [];
let active = 0;
let maximumActive = 0;
let releaseFirstWave!: () => void;
const firstWave = new Promise<void>((resolve) => {
releaseFirstWave = resolve;
});
const run = resolveBuildImages(
"shop",
{
services: Object.fromEntries([
...services.map((service) => [service, { build: "." }]),
["cache", { image: "redis" }],
]),
},
{
request: async <T>(_path: string, init?: ApiRequestInit) => {
const service = ((init?.json ?? {}) as { service: string }).service;
started.push(service);
active++;
maximumActive = Math.max(maximumActive, active);
if (started.length === 3) releaseFirstWave();
await new Promise<void>((resolve) => pending.push(resolve));
active--;
return { reference: `image:${service}` } as T;
},
},
);
try {
await firstWave;
expect(started).toEqual(services.slice(0, 3));
expect(maximumActive).toBe(3);
for (let i = 0; i < services.length; i++) {
const release = pending.shift();
if (!release) throw new Error("Expected a pending image request");
release();
await Promise.resolve();
await Promise.resolve();
}
const images = await run;
expect(Object.keys(images)).toEqual(services);
expect(images).toEqual(
Object.fromEntries(
services.map((service) => [service, `image:${service}`]),
),
);
expect(started).toEqual(services);
} finally {
for (const release of pending) release();
await run.catch(() => {});
}
});
});