import { afterEach, describe, expect, test } from "bun:test"; import { mkdtemp, rm } 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 { workspaceManifestDigest, type WorkspaceSnapshot, } from "../../lib/workspace"; import type { BuildRequest } from "../../shared/build-protocol"; const directories: string[] = []; 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("negotiates and uploads the manifest through resumable blob routes", async () => { const snapshot = emptySnapshot(); const calls: Array<{ path: string; init?: ApiRequestInit; options?: ApiRequestOptions; }> = []; let negotiations = 0; const request: ApiRequester = async ( 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); 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"); }); 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 ( 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 ( 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 ( 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 ( 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 ( 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(); let negotiations = 0; const request: ApiRequester = async ( 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 ( 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("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 void>(); const waitForStart = (count: number) => new Promise((resolve) => { startWaiters.set(count, resolve); }); const firstTwo = waitForStart(2); const third = waitForStart(3); const fourth = waitForStart(4); const released = new Map void>(); const gates = new Map>(); const ids = new Map(); const output: string[] = []; const stream = new Writable({ write(chunk, _encoding, done) { output.push(String(chunk)); done(); }, }); const request: ApiRequester = async ( 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((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: "." }, ]), ), }, 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("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((resolve) => { bothReady = resolve; }); let fail!: (error: Error) => void; let finish!: () => void; const failing = new Promise((_resolve, reject) => { fail = reject; }); const active = new Promise((resolve) => { finish = resolve; }); const request: ApiRequester = async ( 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: "." }, two: { build: "." }, three: { build: "." }, }, }, 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((resolve) => { ready = resolve; }); const request: ApiRequester = async ( 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((_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: "." }, two: { build: "." }, three: { build: "." }, }, }, 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 ( 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 (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((resolve) => { releaseFirstWave = resolve; }); const run = resolveBuildImages( "shop", { services: Object.fromEntries([ ...services.map((service) => [service, { build: "." }]), ["cache", { image: "redis" }], ]), }, { request: async (_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((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(() => {}); } }); });