import { afterEach, describe, expect, spyOn, test } from "bun:test"; import { createHash } from "node:crypto"; import { lstat, mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { BuildConflictError, BuildController, BuildValidationError, type BuildJobObservation, type BuildKubernetesOperations, } from "../../server/build-controller"; import { createApp } from "../../server/app"; import { hashToken, MemoryAuthStore } from "../../server/auth"; import { MemoryBuildStore, type BuildRecord, type BuildStore, } from "../../server/build-store"; import { FilesystemCas } from "../../server/cas"; import type { KubernetesJob } from "../../server/build-job"; import { BUILD_PROTOCOL_VERSION, type BuildRequest, type Sha256Digest, type WorkspaceManifest, } from "../../shared/build-protocol"; const roots: string[] = []; afterEach(async () => { await Promise.all( roots.splice(0).map((root) => rm(root, { recursive: true, force: true })), ); }); function digest(value: Uint8Array | string): Sha256Digest { return `sha256:${createHash("sha256").update(value).digest("hex")}`; } class FakeKubernetes implements BuildKubernetesOperations { jobs: KubernetesJob[] = []; deleted: string[] = []; observation: BuildJobObservation | undefined = { phase: "queued" }; logs = ""; createError?: Error; preserveAfterDelete = false; async createJob(job: KubernetesJob) { if (this.createError) throw this.createError; this.jobs.push(job); } async getJob() { return this.observation; } async getJobLogs() { return this.logs; } async deleteJob(_namespace: string, name: string) { this.deleted.push(name); if (!this.preserveAfterDelete) this.observation = undefined; } } class SupersededBeforeJobStore extends MemoryBuildStore { private superseded = false; override async ownsBuild( imageKey: string, buildId: string, ): Promise { if (!this.superseded) { this.superseded = true; const current = await this.getBuild(buildId); if (current) { const newer = structuredClone(current); newer.metadata.name = "newer-build"; newer.metadata.creationTimestamp = "2026-09-02T00:00:01.000Z"; newer.spec.request.id = "newer-build"; await super.createBuild(newer); } } return super.ownsBuild(imageKey, buildId); } } class TakeoverAtFencedWriteStore extends MemoryBuildStore { onTakeover?: () => void; private tookOver = false; override async replaceBuild( record: BuildRecord, expectedResourceVersion: string, reconcileLeaseToken?: string, ): Promise { if (reconcileLeaseToken && !this.tookOver) { this.tookOver = true; const current = (await this.getBuild(record.metadata.name))!; current.status.reconcileLease!.expiresAt = new Date(0).toISOString(); await super.replaceBuild(current, current.metadata.resourceVersion); this.onTakeover?.(); } return super.replaceBuild( record, expectedResourceVersion, reconcileLeaseToken, ); } } async function fixture( maxLogBytes = 1024, store: BuildStore = new MemoryBuildStore(), ) { const root = await mkdtemp(join(tmpdir(), "kuber-controller-")); roots.push(root); const cas = new FilesystemCas(join(root, "cas")); const kubernetes = new FakeKubernetes(); const source = Buffer.from("FROM scratch\n"); const sourceDigest = await cas.put(source); const manifest: WorkspaceManifest = { version: BUILD_PROTOCOL_VERSION, files: [ { path: "Dockerfile", type: "file", digest: sourceDigest, size: source.byteLength, mode: 0o644, }, ], }; const workspace = await cas.put(Buffer.from(JSON.stringify(manifest))); let now = 0; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", maxLogBytes, now: () => new Date(Date.UTC(2026, 8, 2, 0, 0, now++)), resolveDigest: async () => `sha256:${"f".repeat(64)}`, }); const request: BuildRequest = { version: BUILD_PROTOCOL_VERSION, id: "request-one", project: "demo", service: "web", spec: { architecture: "amd64", image: "registry.test/demo/web:latest", context: ".", buildArgs: [], workspace, }, }; return { root, cas, store, kubernetes, controller, request, workspace, sourceDigest, }; } async function internalFixture( maxLogBytes = 1024, resolveDigest?: (image: string) => Promise, ) { const root = await mkdtemp(join(tmpdir(), "kuber-controller-internal-")); roots.push(root); const cas = new FilesystemCas(join(root, "cas")); const store = new MemoryBuildStore(); const kubernetes = new FakeKubernetes(); const source = Buffer.from("FROM scratch\n"); const sourceDigest = await cas.put(source); const manifest: WorkspaceManifest = { version: BUILD_PROTOCOL_VERSION, files: [ { path: "Dockerfile", type: "file", digest: sourceDigest, size: source.byteLength, mode: 0o644, }, ], }; const workspace = await cas.put(Buffer.from(JSON.stringify(manifest))); let now = 0; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: (request) => `cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/cache-${request.project}-${request.service}`, imageName: (request) => `registry.neko-piranha.ts.net/kuber/${request.project}-${request.service}:latest`, pushImage: (request) => `cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/${request.project}-${request.service}:latest`, pushRegistryInsecure: true, maxLogBytes, now: () => new Date(Date.UTC(2026, 8, 2, 0, 0, now++)), resolveDigest: resolveDigest ?? (async () => `sha256:${"f".repeat(64)}`), }); const request: BuildRequest = { version: BUILD_PROTOCOL_VERSION, id: "request-internal", project: "demo", service: "web", spec: { architecture: "amd64", image: "registry.neko-piranha.ts.net/kuber/demo-web:latest", context: ".", buildArgs: [], workspace, }, }; return { root, cas, store, kubernetes, controller, request, workspace, sourceDigest, }; } describe("build controller", () => { test("publishes multiple service destinations in one Job and resolves every canonical reference", async () => { const resolved: string[] = []; const { controller, request, kubernetes } = await internalFixture( 1024, async (image) => { resolved.push(image); return `sha256:${"f".repeat(64)}`; }, ); request.destinations = [ { service: "worker", image: "external.example/other:latest" }, ]; await controller.submitBuild(request); expect(kubernetes.jobs).toHaveLength(1); expect( (kubernetes.jobs[0]!.spec as any).template.spec.containers[0].args, ).toContain( '--output=type=image,"name=cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/demo-web:latest,cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/demo-worker:latest",push=true,registry.insecure=true', ); kubernetes.observation = { phase: "succeeded" }; expect(await controller.reconcileBuild(request.id)).toMatchObject({ state: "succeeded", }); expect(resolved).toEqual([ "registry.neko-piranha.ts.net/kuber/demo-web:latest", "registry.neko-piranha.ts.net/kuber/demo-worker:latest", ]); expect(await controller.getBuildResult(request.id)).toMatchObject({ reference: `registry.neko-piranha.ts.net/kuber/demo-web@sha256:${"f".repeat(64)}`, references: { worker: `registry.neko-piranha.ts.net/kuber/demo-worker@sha256:${"f".repeat(64)}`, }, }); }); test("fails the shared job if an alias resolves to a different digest", async () => { const { controller, request, kubernetes } = await internalFixture( 1024, async (image) => `sha256:${(image.includes("worker") ? "e" : "f").repeat(64)}`, ); request.destinations = [ { service: "worker", image: "registry.test/worker:latest" }, ]; await controller.submitBuild(request); kubernetes.observation = { phase: "succeeded" }; expect(await controller.reconcileBuild(request.id)).toMatchObject({ state: "failed", error: expect.stringContaining("different digest"), }); await expect(controller.getBuildResult(request.id)).rejects.toBeInstanceOf( BuildConflictError, ); }); test("rejects duplicate alias destinations and unsafe exporter names", async () => { const { controller, request, kubernetes } = await fixture(); request.destinations = [ { service: "web", image: "registry.test/demo/other:latest" }, ]; await expect(controller.submitBuild(request)).rejects.toBeInstanceOf( BuildValidationError, ); request.destinations = [ { service: "worker", image: "registry.test/demo/web:latest" }, ]; await expect(controller.submitBuild(request)).rejects.toBeInstanceOf( BuildValidationError, ); request.destinations = [ { service: "worker", image: "registry.test/demo/worker:latest,push=false", }, ]; await expect(controller.submitBuild(request)).rejects.toBeInstanceOf( BuildValidationError, ); expect(kubernetes.jobs).toHaveLength(0); }); test("negotiates snapshots and resumes verified blob uploads", async () => { const { controller, cas } = await fixture(); const content = Buffer.from("resumable"); const expected = digest(content); expect( await controller.beginBlobUpload(expected, content.byteLength), ).toMatchObject({ offset: 0, complete: false }); expect( await controller.uploadBlobChunk(expected, 0, content.subarray(0, 3)), ).toMatchObject({ offset: 3 }); expect( await controller.beginBlobUpload(expected, content.byteLength), ).toMatchObject({ offset: 3 }); await expect( controller.uploadBlobChunk(expected, 0, content), ).rejects.toBeInstanceOf(BuildConflictError); await controller.uploadBlobChunk(expected, 3, content.subarray(3)); expect(await controller.completeBlobUpload(expected)).toMatchObject({ complete: true, offset: content.byteLength, }); expect( await controller.uploadBlobChunk(expected, 0, Buffer.from("retry")), ).toMatchObject({ complete: true, size: content.byteLength, offset: content.byteLength, }); expect(Buffer.from(await cas.get(expected))).toEqual(content); const bad = `sha256:${"0".repeat(64)}` as Sha256Digest; await controller.beginBlobUpload(bad, 1); await controller.uploadBlobChunk(bad, 0, Buffer.from("x")); await expect(controller.completeBlobUpload(bad)).rejects.toBeInstanceOf( BuildValidationError, ); }); test("reports the manifest and source blobs missing during negotiation", async () => { const { controller, cas } = await fixture(); const absent = `sha256:${"1".repeat(64)}` as Sha256Digest; expect(await controller.negotiateSnapshot(absent)).toEqual({ workspace: absent, missing: [absent], ready: false, }); const manifest = Buffer.from( JSON.stringify({ version: 1, files: [ { path: "x", type: "file", digest: absent, size: 1, mode: 420 }, ], }), ); const workspace = await cas.put(manifest); expect(await controller.negotiateSnapshot(workspace)).toEqual({ workspace, missing: [absent], ready: false, }); }); test("submits once, supersedes competing image builds, captures bounded logs, and resolves immutable results", async () => { const { controller, kubernetes, request } = await fixture(8); expect(await controller.submitBuild(request)).toMatchObject({ state: "queued", }); expect( await controller.submitBuild(structuredClone(request)), ).toMatchObject({ state: "queued" }); expect(kubernetes.jobs).toHaveLength(1); const jobSpec = kubernetes.jobs[0]!.spec as any; expect(jobSpec.template.spec.tolerations).toEqual([ { key: "arch", operator: "Equal", value: "amd64", effect: "NoExecute", }, ]); expect(jobSpec.template.spec.containers[0].volumeMounts).toContainEqual( expect.objectContaining({ name: "workspace", mountPath: "/workspace", subPath: `workspaces/${kubernetes.jobs[0]!.metadata.name}`, }), ); const replacement = { ...request, id: "request-two" }; expect(await controller.submitBuild(replacement)).toMatchObject({ state: "queued", }); expect(kubernetes.jobs).toHaveLength(2); await expect( controller.submitBuild({ ...request, project: "changed" }), ).rejects.toBeInstanceOf(BuildConflictError); kubernetes.logs = "old-line\nnew-line\n"; kubernetes.observation = { phase: "running", startedAt: "2026-09-02T00:00:03.000Z", }; expect(await controller.reconcileBuild(replacement.id)).toMatchObject({ state: "running", }); const logs = (await controller.getBuildEvents(replacement.id)).filter( (event) => event.type === "log", ); expect( logs.reduce( (bytes, event) => bytes + Buffer.byteLength(event.message), 0, ), ).toBeLessThanOrEqual(8); kubernetes.logs += "done\n"; kubernetes.observation = { phase: "succeeded", finishedAt: "2026-09-02T00:00:04.000Z", }; const status = await controller.reconcileBuild(replacement.id); expect(status).toMatchObject({ state: "succeeded", digest: `sha256:${"f".repeat(64)}`, }); expect(await controller.getBuildResult(replacement.id)).toEqual({ image: "registry.test/demo/web", digest: `sha256:${"f".repeat(64)}`, reference: `registry.test/demo/web@sha256:${"f".repeat(64)}`, }); expect(await controller.reconcileBuild(replacement.id)).toEqual(status); }); test("persists creating and starting phases without claiming the build has started", async () => { const { controller, kubernetes, request } = await fixture(); expect(await controller.submitBuild(request)).toMatchObject({ state: "queued", phase: "queued", }); for (const phase of ["creating", "starting"] as const) { kubernetes.observation = { phase }; expect(await controller.reconcileBuild(request.id)).toMatchObject({ state: "queued", phase, }); } kubernetes.observation = { phase: "running" }; expect(await controller.reconcileBuild(request.id)).toMatchObject({ state: "running", phase: "running", }); kubernetes.observation = { phase: "failed", error: "buildkit crashed\nstack line", }; expect(await controller.reconcileBuild(request.id)).toMatchObject({ state: "failed", phase: "done", error: "buildkit crashed\nstack line", }); expect( (await controller.getBuildEvents(request.id)) .filter((event) => event.type === "status") .map((event) => event.type === "status" && event.status.phase), ).toEqual(["queued", "creating", "starting", "running", "done"]); }); test("concurrent controllers converge on one same-ID record and Job", async () => { const first = await fixture(); const second = new BuildController({ cas: first.cas, store: first.store, kubernetes: first.kubernetes, namespace: "builds", workspaceRoot: join(first.root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }); await Promise.all([ first.controller.submitBuild(first.request), second.submitBuild(structuredClone(first.request)), ]); expect(first.kubernetes.jobs).toHaveLength(1); expect(await first.store.getBuild(first.request.id)).toMatchObject({ spec: { request: { id: first.request.id } }, }); }); test("starts a replacement without waiting for superseded Job deletion", async () => { const { controller, kubernetes, request, root, store } = await fixture(); await controller.submitBuild(request); const oldJobName = kubernetes.jobs[0]!.metadata.name; kubernetes.preserveAfterDelete = true; await expect( controller.submitBuild({ ...request, id: "replacement" }), ).resolves.toMatchObject({ state: "queued", }); expect(kubernetes.deleted).toEqual([oldJobName]); expect(kubernetes.jobs).toHaveLength(2); await expect( lstat(join(root, "workspaces", oldJobName)), ).rejects.toMatchObject({ code: "ENOENT" }); expect((await store.getBuild(request.id))?.status).toMatchObject({ state: "failed", error: "Superseded by newer build", cancelled: true, }); }); test("deletes an ambiguous superseded Job even when status was not persisted", async () => { const { controller, kubernetes, request, store } = await fixture(); await controller.submitBuild(request); const oldJobName = kubernetes.jobs[0]!.metadata.name; const old = (await store.getBuild(request.id))!; old.status.jobCreated = false; await store.replaceBuild(old, old.metadata.resourceVersion); await expect( controller.submitBuild({ ...request, id: "replacement" }), ).resolves.toMatchObject({ state: "queued" }); expect(kubernetes.deleted).toEqual([oldJobName]); }); test("does not create a Job when the candidate is superseded before Job creation", async () => { const store = new SupersededBeforeJobStore(); const { controller, kubernetes, request } = await fixture(1024, store); await expect(controller.submitBuild(request)).rejects.toBeInstanceOf( BuildConflictError, ); expect(kubernetes.jobs).toHaveLength(0); }); test("maps supersession before Job creation to HTTP 409", async () => { const auth = new MemoryAuthStore(); await auth.putUser({ username: "operator", passwordHash: "hash", roles: ["operator"], }); await auth.putSession({ tokenHash: hashToken("token"), username: "operator", authVersion: 1, expiresAt: "2030-01-01T00:00:00.000Z", }); const store = new SupersededBeforeJobStore(); const { controller, kubernetes, request: buildRequest, } = await fixture(1024, store); const app = createApp({ store: auth, builds: controller }); const response = await app( new Request("https://kuber.test/api/v2/builds", { method: "POST", headers: { authorization: "Bearer token", "content-type": "application/json", }, body: JSON.stringify(buildRequest), }), ); expect(response.status).toBe(409); expect(await response.json()).toMatchObject({ code: "BUILD_CONFLICT" }); expect(kubernetes.jobs).toHaveLength(0); }); test("releases the image lock after an ambiguous Job creation failure", async () => { const { controller, kubernetes, request, store } = await fixture(); kubernetes.createError = new Error("create response lost"); kubernetes.preserveAfterDelete = true; await expect(controller.submitBuild(request)).rejects.toThrow( "create response lost", ); expect((await store.getBuild(request.id))?.status).toMatchObject({ state: "failed", }); kubernetes.createError = undefined; await expect( controller.submitBuild({ ...request, id: "replacement" }), ).resolves.toMatchObject({ state: "queued" }); }); test("reconciles only active records with created Jobs and isolates failures", async () => { const { controller, request, store } = await fixture(); await controller.submitBuild(request); const candidate = (await store.getBuild(request.id))!; const terminal = structuredClone(candidate); terminal.metadata.name = "terminal"; terminal.spec.request.id = "terminal"; terminal.spec.imageKey = "terminal"; terminal.status.state = "succeeded"; await store.createBuild(terminal); const noJob = structuredClone(candidate); noJob.metadata.name = "no-job"; noJob.spec.request.id = "no-job"; noJob.spec.imageKey = "no-job"; noJob.status.jobCreated = false; await store.createBuild(noJob); const secondCandidate = structuredClone(candidate); secondCandidate.metadata.name = "a-second-candidate"; secondCandidate.spec.request.id = "a-second-candidate"; secondCandidate.spec.imageKey = "a-second-candidate"; await store.createBuild(secondCandidate); const reconciled: string[] = []; spyOn(controller, "reconcileBuild").mockImplementation(async (id) => { reconciled.push(id); if (id === "a-second-candidate") throw new Error("temporary failure"); return controller.getBuildStatus(id); }); await controller.reconcilePendingBuilds(); expect(reconciled).toEqual(["a-second-candidate", request.id]); }); test("shares one reconciliation between an authenticated request and background scan", async () => { const { controller, kubernetes, request } = await fixture(); await controller.submitBuild(request); kubernetes.logs = "build output\n"; let startLogRead!: () => void; let releaseLogRead!: () => void; const logReadStarted = new Promise((resolve) => { startLogRead = resolve; }); const logReadReleased = new Promise((resolve) => { releaseLogRead = resolve; }); const getJobLogs = spyOn(kubernetes, "getJobLogs").mockImplementation( async () => { startLogRead(); await logReadReleased; return kubernetes.logs; }, ); const auth = new MemoryAuthStore(); await auth.putUser({ username: "operator", passwordHash: "hash", roles: ["operator"], }); await auth.putSession({ tokenHash: hashToken("token"), username: "operator", authVersion: 1, expiresAt: "2030-01-01T00:00:00.000Z", }); const app = createApp({ store: auth, builds: controller }); const requestReconcile = app( new Request(`https://kuber.test/api/v2/builds/${request.id}/reconcile`, { method: "POST", headers: { authorization: "Bearer token" }, }), ); await logReadStarted; const backgroundReconcile = controller.reconcilePendingBuilds(); releaseLogRead(); expect((await requestReconcile).status).toBe(200); await backgroundReconcile; expect(getJobLogs).toHaveBeenCalledTimes(1); expect( (await controller.getBuildEvents(request.id)).filter( (event) => event.type === "log", ), ).toHaveLength(1); }); test("does not overlap a hung aborted scan and resumes after it settles", async () => { const { cas, kubernetes, request, root, store } = await fixture(); const submittingController = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }); await submittingController.submitBuild(request); let acquired!: () => void; const leaseAcquired = new Promise((resolve) => { acquired = resolve; }); let firstAcquire = true; let releaseFirstAcquire!: () => void; const firstAcquireReleased = new Promise((resolve) => { releaseFirstAcquire = resolve; }); const lease = { workspaceId: "", holder: "", expiresAt: "", renew: async () => true, release: async () => {}, }; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileLeases: { acquire: async () => { if (!firstAcquire) return lease; firstAcquire = false; acquired(); await firstAcquireReleased; }, }, }); const aborted = new AbortController(); const first = controller.reconcilePendingBuilds({ signal: aborted.signal }); await leaseAcquired; aborted.abort(); const getJob = spyOn(kubernetes, "getJob"); const second = controller.reconcilePendingBuilds(); await Promise.resolve(); expect(getJob).not.toHaveBeenCalled(); releaseFirstAcquire(); await expect(first).rejects.toThrow( "Build reconciliation scan was cancelled", ); await second; await controller.reconcilePendingBuilds(); expect(getJob).toHaveBeenCalledTimes(1); }); test("skips reconciliation while another replica holds the shared lease", async () => { const { cas, kubernetes, request, root, store } = await fixture(); const submittingController = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }); await submittingController.submitBuild(request); kubernetes.logs = "replica must not read this\n"; const getJobLogs = spyOn(kubernetes, "getJobLogs"); const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileLeases: { acquire: async () => undefined }, }); expect(await controller.reconcileBuild(request.id)).toMatchObject({ state: "queued", }); expect(getJobLogs).not.toHaveBeenCalled(); }); test("heartbeats a hung observation so its per-build lease cannot be taken over", async () => { const { cas, kubernetes, request, root, store } = await fixture(); await new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }).submitBuild(request); let releaseObservation!: () => void; let observationStarted!: () => void; const observationPending = new Promise((resolve) => { releaseObservation = resolve; }); const observationStartedPromise = new Promise((resolve) => { observationStarted = resolve; }); kubernetes.getJob = async () => { observationStarted(); await observationPending; return { phase: "queued" }; }; const callbacks: Array<() => void> = []; const now = spyOn(Date, "now").mockReturnValue(0); const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileLeases: { acquire: async () => ({ workspaceId: `build-reconcile:${request.id}`, holder: "replica-one", expiresAt: "", renew: async () => true, release: async () => {}, }), }, reconcileSetTimeout: ((callback: () => void) => { callbacks.push(callback); return 0 as unknown as ReturnType; }) as typeof setTimeout, reconcileClearTimeout: (() => {}) as typeof clearTimeout, }); const reconciliation = controller.reconcileBuild(request.id); await observationStartedPromise; now.mockReturnValue(10_000); callbacks.shift()!(); await Promise.resolve(); await Promise.resolve(); now.mockReturnValue(35_000); expect( await store.acquireBuildReconciliationLease( request.id, "replica-two", 30_000, ), ).toBeUndefined(); expect( (await store.getBuild(request.id))?.status.reconcileLease?.expiresAt, ).toBe(new Date(40_000).toISOString()); releaseObservation(); await reconciliation; now.mockRestore(); }); test("heartbeat loss fences persistence after a hung log read", async () => { const { cas, kubernetes, request, root, store } = await fixture(); await new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }).submitBuild(request); let releaseLogs!: () => void; let logsStarted!: () => void; const logsPending = new Promise((resolve) => { releaseLogs = resolve; }); const logsStartedPromise = new Promise((resolve) => { logsStarted = resolve; }); kubernetes.getJobLogs = async () => { logsStarted(); await logsPending; return "stale output\n"; }; let renewals = 0; const callbacks: Array<() => void> = []; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileLeases: { acquire: async () => ({ workspaceId: `build-reconcile:${request.id}`, holder: "replica-one", expiresAt: "", renew: async () => ++renewals === 1, release: async () => {}, }), }, reconcileSetTimeout: ((callback: () => void) => { callbacks.push(callback); return 0 as unknown as ReturnType; }) as typeof setTimeout, reconcileClearTimeout: (() => {}) as typeof clearTimeout, }); const reconciliation = controller.reconcileBuild(request.id); await logsStartedPromise; callbacks.shift()!(); await Promise.resolve(); await Promise.resolve(); releaseLogs(); await expect(reconciliation).rejects.toThrow( "Build reconciliation lease ownership was lost", ); expect((await store.getBuild(request.id))?.status).toMatchObject({ logOffset: 0, logBytes: 0, events: [{ type: "status" }], }); }); test("does not persist logs after reconciliation lease ownership is lost", async () => { const { cas, kubernetes, request, root, store } = await fixture(); const submittingController = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }); await submittingController.submitBuild(request); kubernetes.logs = "unpersisted output\n"; let renewals = 0; let released = false; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileLeases: { acquire: async () => ({ workspaceId: `build-reconcile:${request.id}`, holder: "replica-one", expiresAt: "2030-01-01T00:00:00.000Z", renew: async () => ++renewals === 1, release: async () => { released = true; }, }), }, }); await expect(controller.reconcileBuild(request.id)).rejects.toThrow( "Build reconciliation lease ownership was lost", ); expect(released).toBe(true); expect((await store.getBuild(request.id))?.status).toMatchObject({ logOffset: 0, logBytes: 0, events: [{ type: "status" }], }); }); test("fences a stale reconciler when its lease is taken over before persistence", async () => { const store = new TakeoverAtFencedWriteStore(); const { cas, kubernetes, request, root } = await fixture(1024, store); const submitter = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", }); await submitter.submitBuild(request); kubernetes.logs = "stale output\n"; kubernetes.observation = { phase: "succeeded" }; const stale = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileHolder: "stale-replica", resolveDigest: async () => `sha256:${"a".repeat(64)}`, }); const current = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", reconcileHolder: "current-replica", resolveDigest: async () => `sha256:${"b".repeat(64)}`, }); let currentReconciliation: Promise | undefined; store.onTakeover = () => { kubernetes.logs = "current output\n"; currentReconciliation = current.reconcileBuild(request.id); }; await expect(stale.reconcileBuild(request.id)).rejects.toThrow( "Build reconciliation lease ownership was lost", ); await currentReconciliation; const record = (await store.getBuild(request.id))!; expect(record.status).toMatchObject({ state: "succeeded", digest: `sha256:${"b".repeat(64)}`, logOffset: Buffer.byteLength("current output\n"), }); expect( record.status.events.filter((event) => event.type === "log"), ).toEqual([ expect.objectContaining({ message: "current output\n", sequence: 1 }), ]); expect(await store.ownsBuild(record.spec.imageKey, request.id)).toBe(false); }); test("cancels idempotently and cleans up only terminal build resources", async () => { const { controller, kubernetes, request, root } = await fixture(); await controller.submitBuild(request); await expect(controller.cleanupBuild(request.id)).rejects.toBeInstanceOf( BuildConflictError, ); expect(await controller.cancelBuild(request.id)).toMatchObject({ state: "failed", error: "Build cancelled", }); const deletes = kubernetes.deleted.length; expect(await controller.cancelBuild(request.id)).toMatchObject({ error: "Build cancelled", }); expect(kubernetes.deleted).toHaveLength(deletes); await controller.cleanupBuild(request.id); expect(kubernetes.deleted.length).toBe(deletes + 1); await expect( lstat(join(root, "workspaces", kubernetes.jobs[0]!.metadata.name)), ).rejects.toMatchObject({ code: "ENOENT" }); }); test("pushes to internal registry, records canonical result, and uses insecure flags", async () => { const { controller, kubernetes, request } = await internalFixture(); await controller.submitBuild(request); expect(kubernetes.jobs).toHaveLength(1); const container = (kubernetes.jobs[0]!.spec as any).template.spec .containers[0]; expect(container.args).toContain( "--output=type=image,name=cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/demo-web:latest,push=true,registry.insecure=true", ); expect(container.args).toContain( "--import-cache=type=registry,ref=cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/cache-demo-web,registry.insecure=true", ); expect(container.args).toContain( "--export-cache=type=registry,ref=cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/cache-demo-web,mode=max,registry.insecure=true", ); kubernetes.observation = { phase: "succeeded", finishedAt: "2026-09-02T00:00:04.000Z", }; const status = await controller.reconcileBuild(request.id); expect(status).toMatchObject({ state: "succeeded" }); const result = await controller.getBuildResult(request.id); expect(result.image).toBe("registry.neko-piranha.ts.net/kuber/demo-web"); expect(result.reference).toBe( `registry.neko-piranha.ts.net/kuber/demo-web@sha256:${"f".repeat(64)}`, ); }); test("does not add insecure flags when pushRegistryInsecure is false", async () => { const root = await mkdtemp(join(tmpdir(), "kuber-controller-secure-")); roots.push(root); const cas = new FilesystemCas(join(root, "cas")); const store = new MemoryBuildStore(); const kubernetes = new FakeKubernetes(); const source = Buffer.from("FROM scratch\n"); const sourceDigest = await cas.put(source); const manifest: WorkspaceManifest = { version: BUILD_PROTOCOL_VERSION, files: [ { path: "Dockerfile", type: "file", digest: sourceDigest, size: source.byteLength, mode: 0o644, }, ], }; const workspace = await cas.put(Buffer.from(JSON.stringify(manifest))); let now = 0; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "registry.test/cache/app", imageName: (request) => `registry.test/${request.project}-${request.service}:latest`, pushImage: (request) => `internal.registry:5000/${request.project}-${request.service}:latest`, maxLogBytes: 1024, now: () => new Date(Date.UTC(2026, 8, 2, 0, 0, now++)), resolveDigest: async () => `sha256:${"f".repeat(64)}`, }); const request: BuildRequest = { version: BUILD_PROTOCOL_VERSION, id: "request-secure", project: "demo", service: "web", spec: { architecture: "amd64", image: "registry.test/demo-web:latest", context: ".", buildArgs: [], workspace, }, }; await controller.submitBuild(request); const container = (kubernetes.jobs[0]!.spec as any).template.spec .containers[0]; for (const arg of container.args) { expect(arg).not.toContain("registry.insecure"); } }); test("pushImage from request cannot redirect to a different image in results", async () => { const root = await mkdtemp(join(tmpdir(), "kuber-controller-redirect-")); roots.push(root); const cas = new FilesystemCas(join(root, "cas")); const store = new MemoryBuildStore(); const kubernetes = new FakeKubernetes(); const source = Buffer.from("FROM scratch\n"); const sourceDigest = await cas.put(source); const manifest: WorkspaceManifest = { version: BUILD_PROTOCOL_VERSION, files: [ { path: "Dockerfile", type: "file", digest: sourceDigest, size: source.byteLength, mode: 0o644, }, ], }; const workspace = await cas.put(Buffer.from(JSON.stringify(manifest))); let now = 0; const controller = new BuildController({ cas, store, kubernetes, namespace: "builds", workspaceRoot: join(root, "workspaces"), workspaceClaimName: "workspaces", cacheImage: "canonical.test/cache/app", imageName: (request) => `canonical.test/${request.project}-${request.service}:latest`, pushImage: (request) => `internal.test/${request.project}-${request.service}:latest`, maxLogBytes: 1024, now: () => new Date(Date.UTC(2026, 8, 2, 0, 0, now++)), resolveDigest: async () => `sha256:${"f".repeat(64)}`, }); const request: BuildRequest = { version: BUILD_PROTOCOL_VERSION, id: "request-redirect", project: "demo", service: "web", spec: { architecture: "amd64", image: "canonical.test/demo-web:latest", context: ".", buildArgs: [], workspace, }, }; await controller.submitBuild(request); const container = (kubernetes.jobs[0]!.spec as any).template.spec .containers[0]; expect( container.args.some((arg: string) => arg.includes("name=internal.test/demo-web:latest"), ), ).toBe(true); kubernetes.observation = { phase: "succeeded", finishedAt: "2026-09-02T00:00:04.000Z", }; await controller.reconcileBuild(request.id); const result = await controller.getBuildResult(request.id); expect(result.image).toBe("canonical.test/demo-web"); expect(result.reference).toBe( `canonical.test/demo-web@sha256:${"f".repeat(64)}`, ); }); });