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 ( 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((resolve) => { release = resolve; }); const started = new Promise((resolve) => { patchStarted = resolve; }); let negotiations = 0; const patches = new Map(); const request: ApiRequester = async ( 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 ( 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 ( 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 (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 ( 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 () => ({ 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(); const submissions: BuildRequest[] = []; const progress: Array<[number, number]> = []; const request: ApiRequester = async ( 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 ( 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 ( 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("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 (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 (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((resolve) => { listening = resolve; }); const calls: string[] = []; const request: ApiRequester = async (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((_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(); const settled: string[] = []; const request: ApiRequester = async ( 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(); const output = new Map(); let started = 0; let release!: () => void; const bothStarted = new Promise((resolve) => { release = resolve; }); const request: ApiRequester = async ( 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 (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 ( 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 (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((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); ready(); return await new Promise((_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 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: { 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(); const ids = new Map(); const settled: Array<[string, unknown]> = []; const statuses = ["creating", "starting", "running", "done"] as const; let polls = 0; const request: ApiRequester = async ( 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((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: { 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((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: { 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 ( 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(() => {}); } }); });