1584 lines
54 KiB
TypeScript
1584 lines
54 KiB
TypeScript
import { afterEach, describe, expect, spyOn, test } from "bun:test";
|
|
import { createHash } from "node:crypto";
|
|
import { lstat, mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
import {
|
|
BuildConflictError,
|
|
BuildController,
|
|
BuildValidationError,
|
|
type BuildJobObservation,
|
|
type BuildKubernetesOperations,
|
|
type BuildControllerOptions,
|
|
} 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<boolean> {
|
|
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<void> {
|
|
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(),
|
|
extra: Partial<BuildControllerOptions> = {},
|
|
) {
|
|
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)}`,
|
|
...extra,
|
|
});
|
|
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<Sha256Digest>,
|
|
) {
|
|
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("stages URI outside source, namespaces cache and cleans superseded/terminal packages", async () => {
|
|
const staged: string[] = [];
|
|
const fixtureData = await fixture(1024, new MemoryBuildStore(), {
|
|
stageBuildpack: async (_uri, arch, destination) => {
|
|
expect(arch).toBe("arm64");
|
|
staged.push(destination);
|
|
await mkdir(destination, { recursive: true });
|
|
await writeFile(join(destination, "order.toml"), "order");
|
|
},
|
|
});
|
|
const { controller, request, root, store, kubernetes } = fixtureData;
|
|
request.spec.builder = "buildpacks";
|
|
request.spec.architecture = "arm64";
|
|
request.spec.buildpackUri =
|
|
"https://github.com/example/bun/releases/download/v1/buildpack.cnb";
|
|
await controller.submitBuild(request);
|
|
const first = (await store.getBuild(request.id))!;
|
|
expect(staged[0]).toBe(
|
|
join(root, "workspaces", `${first.spec.workspaceSubPath}-buildpack`),
|
|
);
|
|
expect(
|
|
await Bun.file(
|
|
join(root, "workspaces", first.spec.workspaceSubPath, "Dockerfile"),
|
|
).text(),
|
|
).toBe("FROM scratch\n");
|
|
const firstArgs = (kubernetes.jobs[0]!.spec as any).template.spec
|
|
.containers[0].args as string[];
|
|
const cache1 = firstArgs[firstArgs.indexOf("--cache-image") + 1];
|
|
expect(cache1).toMatch(/:buildpack-[a-f0-9]{24}$/);
|
|
const newer = structuredClone(request);
|
|
newer.id = "request-two";
|
|
newer.spec.buildpackUri = newer.spec.buildpackUri!.replace("/v1/", "/v2/");
|
|
await controller.submitBuild(newer);
|
|
expect(await Bun.file(join(staged[0]!, "order.toml")).exists()).toBe(false);
|
|
const nextArgs = (kubernetes.jobs[1]!.spec as any).template.spec
|
|
.containers[0].args as string[];
|
|
expect(nextArgs[nextArgs.indexOf("--cache-image") + 1]).not.toBe(cache1);
|
|
kubernetes.observation = { phase: "failed", error: "test" };
|
|
await controller.reconcileBuild(newer.id);
|
|
await controller.cleanupBuild(newer.id);
|
|
expect(await Bun.file(join(staged[1]!, "order.toml")).exists()).toBe(false);
|
|
});
|
|
|
|
test("invalid packages fail before Job and clean source and partial staging", async () => {
|
|
let destination = "";
|
|
const { controller, request, store, kubernetes, root } = await fixture(
|
|
1024,
|
|
new MemoryBuildStore(),
|
|
{
|
|
stageBuildpack: async (_uri, _arch, path) => {
|
|
destination = path;
|
|
await mkdir(path, { recursive: true });
|
|
await writeFile(join(path, "partial"), "bytes");
|
|
throw new Error(
|
|
"package target linux/arm64 does not match linux/amd64",
|
|
);
|
|
},
|
|
},
|
|
);
|
|
request.spec.builder = "buildpacks";
|
|
request.spec.buildpackUri =
|
|
"https://github.com/example/bun/releases/download/v1/buildpack.cnb";
|
|
await expect(controller.submitBuild(request)).rejects.toThrow(/target/);
|
|
expect(kubernetes.jobs).toHaveLength(0);
|
|
const record = (await store.getBuild(request.id))!;
|
|
expect(record.status.state).toBe("failed");
|
|
expect(await Bun.file(join(destination, "partial")).exists()).toBe(false);
|
|
expect(
|
|
await Bun.file(
|
|
join(root, "workspaces", record.spec.workspaceSubPath, "Dockerfile"),
|
|
).exists(),
|
|
).toBe(false);
|
|
});
|
|
|
|
test("validates URI/builder pairing before creating a record", async () => {
|
|
const { controller, request, store, kubernetes } = await fixture();
|
|
request.spec.buildpackUri =
|
|
"https://github.com/example/bun/releases/download/v1/buildpack.cnb";
|
|
await expect(controller.submitBuild(request)).rejects.toThrow(
|
|
/requires the buildpacks/,
|
|
);
|
|
request.spec.builder = "buildpacks";
|
|
request.spec.buildpackUri = "https://localhost/buildpack.cnb";
|
|
await expect(controller.submitBuild(request)).rejects.toThrow(/URI/);
|
|
expect(await store.getBuild(request.id)).toBeUndefined();
|
|
expect(kubernetes.jobs).toHaveLength(0);
|
|
});
|
|
test("correlates a UUID build request with the stored build and emits only summary keys", async () => {
|
|
const { controller, request, store } = await fixture();
|
|
request.id = "123e4567-e89b-42d3-a456-426614174000";
|
|
const lines: string[] = [];
|
|
const output = spyOn(console, "info").mockImplementation((line) => {
|
|
lines.push(String(line));
|
|
});
|
|
try {
|
|
await controller.submitBuild(request);
|
|
const summary = JSON.parse(lines[0]!);
|
|
const record = await store.getBuild(request.id);
|
|
expect(summary.buildRequestId).toBe(request.id);
|
|
expect(record?.metadata.name).toBe(summary.buildRequestId);
|
|
expect(record?.spec.jobName).toMatch(/^kuber-build-/);
|
|
expect(Object.keys(summary).sort()).toEqual(
|
|
[
|
|
"buildRequestId",
|
|
"elapsedMs",
|
|
"event",
|
|
"jobCreated",
|
|
"materialize",
|
|
"outcome",
|
|
"phases",
|
|
"queueWaitMs",
|
|
].sort(),
|
|
);
|
|
expect(lines[0]).not.toContain(request.spec.image);
|
|
expect(lines[0]).not.toContain(request.spec.workspace);
|
|
} finally {
|
|
output.mockRestore();
|
|
}
|
|
});
|
|
|
|
test("ignores summary logger failures without changing build success", async () => {
|
|
const { controller, request } = await fixture();
|
|
const output = spyOn(console, "info").mockImplementation(() => {
|
|
throw new Error("logger unavailable");
|
|
});
|
|
try {
|
|
const status = await controller.submitBuild(request);
|
|
expect(status.state).toBe("queued");
|
|
} finally {
|
|
output.mockRestore();
|
|
}
|
|
});
|
|
|
|
test("attributes deferred store work and serialized queue wait without logging identifiers", async () => {
|
|
let release!: () => void;
|
|
let entered!: () => void;
|
|
const blocked = new Promise<void>((resolve) => {
|
|
release = resolve;
|
|
});
|
|
const started = new Promise<void>((resolve) => {
|
|
entered = resolve;
|
|
});
|
|
class SlowStore extends MemoryBuildStore {
|
|
override async createBuild(record: BuildRecord) {
|
|
entered();
|
|
await blocked;
|
|
return super.createBuild(record);
|
|
}
|
|
}
|
|
const { controller, request } = await fixture(1024, new SlowStore());
|
|
const lines: string[] = [];
|
|
const output = spyOn(console, "info").mockImplementation((line) => {
|
|
lines.push(String(line));
|
|
});
|
|
try {
|
|
const first = controller.submitBuild(request);
|
|
await started;
|
|
const second = controller.submitBuild(request);
|
|
await Bun.sleep(25);
|
|
release();
|
|
await Promise.all([first, second]);
|
|
expect(lines).toHaveLength(2);
|
|
const [created, existing] = lines.map((line) => JSON.parse(line));
|
|
expect(created.phases.recordCreateMs).toBeGreaterThan(15);
|
|
expect(created.phases.materializeMs).toBeGreaterThan(0);
|
|
expect(existing.queueWaitMs).toBeGreaterThan(15);
|
|
expect(existing.jobCreated).toBe(false);
|
|
expect(existing.phases.materializeMs).toBeUndefined();
|
|
for (const line of lines) {
|
|
expect(line).not.toContain(request.id);
|
|
expect(line).not.toContain(request.spec.workspace);
|
|
expect(line).not.toContain(request.spec.image);
|
|
expect(line).not.toContain("Dockerfile");
|
|
}
|
|
} finally {
|
|
release();
|
|
output.mockRestore();
|
|
}
|
|
});
|
|
|
|
test("attributes deferred materialization and logs one safe failure summary", async () => {
|
|
const { cas, store, kubernetes, request, root } = await fixture();
|
|
let release!: () => void;
|
|
let entered!: () => void;
|
|
const blocked = new Promise<void>((resolve) => {
|
|
release = resolve;
|
|
});
|
|
const started = new Promise<void>((resolve) => {
|
|
entered = resolve;
|
|
});
|
|
const secret = "private-path-and-credential";
|
|
const controller = new BuildController({
|
|
cas,
|
|
store,
|
|
kubernetes,
|
|
namespace: "builds",
|
|
workspaceRoot: join(root, "workspaces"),
|
|
workspaceClaimName: "workspaces",
|
|
cacheImage: "cache",
|
|
materialize: async () => {
|
|
entered();
|
|
await blocked;
|
|
throw new Error(secret);
|
|
},
|
|
});
|
|
const lines: string[] = [];
|
|
const output = spyOn(console, "info").mockImplementation((line) => {
|
|
lines.push(String(line));
|
|
});
|
|
try {
|
|
const submission = controller.submitBuild(request);
|
|
const failure = submission.catch((error: unknown) => error);
|
|
await started;
|
|
await Bun.sleep(25);
|
|
release();
|
|
const error = await failure;
|
|
expect(error).toBeInstanceOf(Error);
|
|
expect((error as Error).message).toBe(secret);
|
|
expect(lines).toHaveLength(1);
|
|
const summary = JSON.parse(lines[0]!);
|
|
expect(summary).toMatchObject({
|
|
event: "build_submission_timing",
|
|
outcome: "failure",
|
|
jobCreated: false,
|
|
});
|
|
expect(summary.phases.materializeMs).toBeGreaterThan(15);
|
|
expect(summary.phases.failureCleanupMs).toBeGreaterThanOrEqual(0);
|
|
expect(summary.phases.jobCreateMs).toBeUndefined();
|
|
expect(lines[0]).not.toContain(secret);
|
|
expect(lines[0]).not.toContain(request.id);
|
|
} finally {
|
|
release();
|
|
output.mockRestore();
|
|
}
|
|
});
|
|
|
|
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("submits a Buildpacks Job with mapped destinations and retains result and status semantics", async () => {
|
|
const resolved: string[] = [];
|
|
const { controller, request, kubernetes } = await internalFixture(
|
|
1024,
|
|
async (image) => {
|
|
resolved.push(image);
|
|
return `sha256:${"f".repeat(64)}`;
|
|
},
|
|
);
|
|
request.spec.builder = "buildpacks";
|
|
request.destinations = [
|
|
{ service: "worker", image: "external.example/other:latest" },
|
|
];
|
|
expect(await controller.submitBuild(request)).toMatchObject({
|
|
state: "queued",
|
|
});
|
|
const pod = (kubernetes.jobs[0]!.spec as any).template.spec;
|
|
expect(pod.containers[0].command).toEqual(["/cnb/lifecycle/creator"]);
|
|
expect(pod.containers[0].args).toContain("--tag");
|
|
expect(pod.containers[0].args).toContain(
|
|
"cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/demo-worker:latest",
|
|
);
|
|
expect(pod.containers[0].args.at(-1)).toBe(
|
|
"cncf-distribution-svc.registry.svc.cluster.local:5000/kuber/demo-web:latest",
|
|
);
|
|
kubernetes.observation = { phase: "succeeded" };
|
|
expect(await controller.reconcileBuild(request.id)).toMatchObject({
|
|
state: "succeeded",
|
|
digest: `sha256:${"f".repeat(64)}`,
|
|
});
|
|
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("rejects unsupported builders and Buildpacks Dockerfile options before Job creation", async () => {
|
|
const { controller, request, kubernetes } = await fixture();
|
|
for (const incompatible of [
|
|
{ dockerfile: "Dockerfile" },
|
|
{ target: "production" },
|
|
{ buildArgs: ["A=B"] },
|
|
]) {
|
|
await expect(
|
|
controller.submitBuild({
|
|
...request,
|
|
spec: {
|
|
...request.spec,
|
|
builder: "buildpacks",
|
|
...incompatible,
|
|
},
|
|
}),
|
|
).rejects.toBeInstanceOf(BuildValidationError);
|
|
}
|
|
await expect(
|
|
controller.submitBuild({
|
|
...request,
|
|
spec: {
|
|
...request.spec,
|
|
builder: "unknown" as any,
|
|
},
|
|
}),
|
|
).rejects.toBeInstanceOf(BuildValidationError);
|
|
expect(kubernetes.jobs).toHaveLength(0);
|
|
});
|
|
|
|
test("rejects an unpinned Buildpacks image before creating a record", async () => {
|
|
const { controller, request, kubernetes, store } = await fixture();
|
|
request.spec.builder = "buildpacks";
|
|
(controller as any).options.buildpacksImage = "heroku/builder:26";
|
|
await expect(controller.submitBuild(request)).rejects.toThrow(
|
|
"digest-pinned",
|
|
);
|
|
expect(kubernetes.jobs).toHaveLength(0);
|
|
expect(await store.getBuild(request.id)).toBeUndefined();
|
|
});
|
|
|
|
test("fails a Buildpacks build with destinations in different registries", async () => {
|
|
const { controller, request, kubernetes } = await fixture();
|
|
request.spec.builder = "buildpacks";
|
|
request.destinations = [
|
|
{ service: "worker", image: "another.example.com/worker:latest" },
|
|
];
|
|
await expect(controller.submitBuild(request)).rejects.toThrow(
|
|
"same registry",
|
|
);
|
|
expect(kubernetes.jobs).toHaveLength(0);
|
|
expect(await controller.getBuildStatus(request.id)).toMatchObject({
|
|
state: "failed",
|
|
error: expect.stringContaining("same registry"),
|
|
});
|
|
});
|
|
|
|
test("Buildpacks alias digest mismatch fails and cancellation deletes its Job", async () => {
|
|
const { controller, request, kubernetes } = await internalFixture(
|
|
1024,
|
|
async (image) =>
|
|
`sha256:${(image.includes("worker") ? "e" : "f").repeat(64)}`,
|
|
);
|
|
request.spec.builder = "buildpacks";
|
|
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"),
|
|
});
|
|
|
|
const pending = await fixture();
|
|
pending.request.spec.builder = "buildpacks";
|
|
await pending.controller.submitBuild(pending.request);
|
|
expect(
|
|
await pending.controller.cancelBuild(pending.request.id),
|
|
).toMatchObject({ state: "failed" });
|
|
expect(pending.kubernetes.deleted).toContain(
|
|
pending.kubernetes.jobs[0]!.metadata.name,
|
|
);
|
|
});
|
|
|
|
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<void>((resolve) => {
|
|
startLogRead = resolve;
|
|
});
|
|
const logReadReleased = new Promise<void>((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<void>((resolve) => {
|
|
acquired = resolve;
|
|
});
|
|
let firstAcquire = true;
|
|
let releaseFirstAcquire!: () => void;
|
|
const firstAcquireReleased = new Promise<void>((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<void>((resolve) => {
|
|
releaseObservation = resolve;
|
|
});
|
|
const observationStartedPromise = new Promise<void>((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<typeof setTimeout>;
|
|
}) 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<void>((resolve) => {
|
|
releaseLogs = resolve;
|
|
});
|
|
const logsStartedPromise = new Promise<void>((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<typeof setTimeout>;
|
|
}) 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<unknown> | 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)}`,
|
|
);
|
|
});
|
|
});
|