diff --git a/changelogs/2.6.2.md b/changelogs/2.6.2.md new file mode 100644 index 0000000..b0d30fc --- /dev/null +++ b/changelogs/2.6.2.md @@ -0,0 +1,11 @@ +# 2.6.2-rc1 + +## Fixed + +- Restore maintenance support for leading `*.domain` hostnames. +- Bound snapshot negotiation requests to five minutes and validate manifest path conflicts in linear time. + +## Changed + +- Emit formatted request log lines without the generic JSON fallback. +- Log build submission phase timings and materialization counts, sizes, CAS reads, and filesystem writes to investigate late Job creation. The observed 174-second delay is not yet resolved. diff --git a/command/main.ts b/command/main.ts index d79ffd4..6bfebce 100644 --- a/command/main.ts +++ b/command/main.ts @@ -37,7 +37,7 @@ function cloneCommand(command: T): T { export const main = defineCommand({ meta: { name: "kuber", - version: "2.6.1", + version: "2.6.2-rc1", description: "Docker Compose -> K8s translation layer", }, args: { diff --git a/lib/build.ts b/lib/build.ts index a8c2229..b875978 100644 --- a/lib/build.ts +++ b/lib/build.ts @@ -346,6 +346,7 @@ async function uploadBlob( request: ApiRequester, scheduler: TaskScheduler, project?: string, + signal?: AbortSignal, ): Promise { const uploadPath = `/blobs/${encodeURIComponent(digest)}/uploads`; const projectQuery = project ? `?project=${encodeURIComponent(project)}` : ""; @@ -354,6 +355,7 @@ async function uploadBlob( request<{ offset: number; complete: boolean }>(path, { method: "POST", json: { size: data.byteLength }, + signal, }), ); let offset = progress.offset; @@ -362,6 +364,7 @@ async function uploadBlob( const uploaded = await scheduler.run(() => request<{ offset: number }>(path, { method: "PATCH", + signal, headers: { "content-type": "application/octet-stream", "upload-offset": String(offset), @@ -378,6 +381,7 @@ async function uploadBlob( request(`${uploadPath}/complete${projectQuery}`, { method: "POST", json: {}, + signal, }), ); } @@ -389,17 +393,24 @@ export async function uploadWorkspaceSnapshot( reporter?: BuildReporter, scheduler?: TaskScheduler, project?: string, + signal?: AbortSignal, ): Promise { const blobs = new Map(snapshot.blobs.map((blob) => [blob.digest, blob.data])); blobs.set(snapshot.digest, serializeWorkspaceManifest(snapshot.manifest)); const requestScheduler = scheduler ?? new TaskScheduler(); for (;;) { + signal?.throwIfAborted(); const negotiation = await requestScheduler.run(() => - request("/snapshots/negotiate", { - method: "POST", - json: { workspace: snapshot.digest, ...(project && { project }) }, - }), + request( + "/snapshots/negotiate", + { + method: "POST", + json: { workspace: snapshot.digest, ...(project && { project }) }, + signal, + }, + { timeoutMs: 300_000 }, + ), ); if (negotiation.ready) return; if (negotiation.missing.length === 0) @@ -414,7 +425,14 @@ export async function uploadWorkspaceSnapshot( if (!data) throw new Error(`Server requested unknown workspace blob ${digest}`); await reporter?.progress?.(`Uploading ${digest}`); - await uploadBlob(digest, data, request, requestScheduler, project); + await uploadBlob( + digest, + data, + request, + requestScheduler, + project, + signal, + ); }, ); } @@ -712,6 +730,7 @@ export async function buildServices( reporter, options.scheduler, project, + options.signal, ); const references: string[] = new Array(plans.length); diff --git a/lib/request-log.ts b/lib/request-log.ts index 1291e3f..a92e93b 100644 --- a/lib/request-log.ts +++ b/lib/request-log.ts @@ -100,8 +100,7 @@ export const processLogger: ProcessLogger = { log(entry) { try { const line = requestLine(entry); - if (line === "") return; - console.log(line ?? `KUBER_REQUEST ${JSON.stringify(logValue(entry))}`); + if (line) console.log(line); } catch { // Process logging must never change request behavior. } diff --git a/package.json b/package.json index c6bf7ac..f912378 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.6.1", + "version": "2.6.2-rc1", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/server/build-controller.ts b/server/build-controller.ts index 333855f..d68647e 100644 --- a/server/build-controller.ts +++ b/server/build-controller.ts @@ -23,6 +23,7 @@ import { materializeWorkspace, parseWorkspaceManifest, type MaterializeCas, + type MaterializeTiming, } from "./materialize"; import { parseImageReference, resolveRegistryDigest } from "./registry"; @@ -39,12 +40,7 @@ export interface BuildCas extends MaterializeCas { } export type JobPhase = - | "queued" - | "creating" - | "starting" - | "running" - | "succeeded" - | "failed"; + "queued" | "creating" | "starting" | "running" | "succeeded" | "failed"; export interface BuildJobObservation { phase: JobPhase; @@ -451,9 +447,68 @@ export class BuildController { })); validateRequest(request); const imageKey = `${request.project}\0${request.service}\0${request.spec.image}`; + const started = performance.now(); const previous = this.imageSubmissionLocks.get(imageKey) ?? Promise.resolve(); - const current = previous.then(() => this.submitBuildInternal(request)); + const current = previous.then(async () => { + const queueWaitMs = performance.now() - started; + const phases: Record = {}; + const materialize: MaterializeTiming = { casReadMs: 0, fsWriteMs: 0 }; + let outcome = "failure"; + let jobCreated = false; + const timed = async ( + phase: string, + run: () => Promise, + ): Promise => { + const start = performance.now(); + try { + return await run(); + } finally { + phases[phase] = performance.now() - start; + } + }; + try { + const status = await this.submitBuildInternal( + request, + timed, + materialize, + () => { + jobCreated = true; + }, + ); + outcome = "success"; + return status; + } finally { + // Fixed-schema numeric telemetry only: never serialize a request or error. + try { + console.info( + JSON.stringify({ + event: "build_submission_timing", + buildRequestId: + /^[0-9a-f]{8}-[0-9a-f]{4}-[1-8][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test( + request.id, + ) + ? request.id + : null, + outcome, + jobCreated, + elapsedMs: performance.now() - started, + queueWaitMs, + phases, + materialize: { + fileCount: materialize.fileCount ?? null, + manifestBytes: materialize.manifestBytes ?? null, + fileBytes: materialize.fileBytes ?? null, + casReadMs: materialize.casReadMs, + fsWriteMs: materialize.fsWriteMs, + }, + }), + ); + } catch { + // Telemetry must not change submission behavior. + } + } + }); const entry = current.catch(() => undefined); this.imageSubmissionLocks.set(imageKey, entry); try { @@ -466,9 +521,14 @@ export class BuildController { private async submitBuildInternal( request: BuildRequest, + timed: (phase: string, run: () => Promise) => Promise, + materialize: MaterializeTiming, + markJobCreated: () => void, ): Promise { const imageKey = `${request.project}\0${request.service}\0${request.spec.image}`; - const snapshot = await this.negotiateSnapshot(request.spec.workspace); + const snapshot = await timed("snapshotMs", () => + this.negotiateSnapshot(request.spec.workspace), + ); if (!snapshot.ready) throw new BuildConflictError( `Workspace snapshot is incomplete: ${snapshot.missing.join(", ")}`, @@ -511,26 +571,33 @@ export class BuildController { }; let stored: BuildRecord; try { - const result = await this.options.store.createBuild(record); + const result = await timed("recordCreateMs", () => + this.options.store.createBuild(record), + ); stored = result.record; - for (const superseded of result.superseded ?? []) { - // Deletion is best effort and does not wait for the Job or its pods. - await this.options.kubernetes - .deleteJob(this.options.namespace, superseded.spec.jobName) - .catch(() => undefined); - await rm( - join(this.options.workspaceRoot, superseded.spec.workspaceSubPath), - { recursive: true, force: true }, - ).catch(() => undefined); - } + await timed("supersededCleanupMs", async () => { + for (const superseded of result.superseded ?? []) { + // Deletion is best effort and does not wait for the Job or its pods. + await this.options.kubernetes + .deleteJob(this.options.namespace, superseded.spec.jobName) + .catch(() => undefined); + await rm( + join(this.options.workspaceRoot, superseded.spec.workspaceSubPath), + { recursive: true, force: true }, + ).catch(() => undefined); + } + }); if (!result.created) return recordStatus(stored); - const current = await this.options.store.getBuild(request.id); - if ( - !current || - current.status.state !== "queued" || - (this.options.store.ownsBuild && - !(await this.options.store.ownsBuild(imageKey, request.id))) - ) + const owns = await timed("initialOwnershipMs", async () => { + const current = await this.options.store.getBuild(request.id); + return ( + !!current && + current.status.state === "queued" && + (!this.options.store.ownsBuild || + (await this.options.store.ownsBuild(imageKey, request.id))) + ); + }); + if (!owns) throw new BuildConflictError( `Build '${request.id}' was superseded before Job creation`, ); @@ -540,72 +607,84 @@ export class BuildController { throw error; } try { - await this.materializer( - this.options.cas, - request.spec.workspace, - join(this.options.workspaceRoot, stored.spec.workspaceSubPath), + await timed("materializeMs", () => + this.materializer( + this.options.cas, + request.spec.workspace, + join(this.options.workspaceRoot, stored.spec.workspaceSubPath), + materialize, + ), ); if ( this.options.store.ownsBuild && - !(await this.options.store.ownsBuild(imageKey, request.id)) + !(await timed("finalOwnershipMs", () => + this.options.store.ownsBuild!(imageKey, request.id), + )) ) throw new BuildConflictError( `Build '${request.id}' lost image lock before Job creation`, ); - const cacheImage = - typeof this.options.cacheImage === "function" - ? this.options.cacheImage(request) - : this.options.cacheImage; - const pushImage = this.options.pushImage - ? this.options.pushImage(request) - : request.spec.image; - const pushImages = request.destinations?.map((destination) => - this.options.pushImage - ? this.options.pushImage({ - ...request, - service: destination.service, - spec: { ...request.spec, image: destination.image }, - destinations: undefined, - }) - : destination.image, + const job = await timed("jobSpecMs", async () => { + const cacheImage = + typeof this.options.cacheImage === "function" + ? this.options.cacheImage(request) + : this.options.cacheImage; + const pushImage = this.options.pushImage + ? this.options.pushImage(request) + : request.spec.image; + const pushImages = request.destinations?.map((destination) => + this.options.pushImage + ? this.options.pushImage({ + ...request, + service: destination.service, + spec: { ...request.spec, image: destination.image }, + destinations: undefined, + }) + : destination.image, + ); + return createBuildJob({ + name: jobName, + namespace: this.options.namespace, + spec: request.spec, + workspaceClaimName: this.options.workspaceClaimName, + workspaceSubPath: jobWorkspaceSubPath( + this.options.workspaceRoot, + stored.spec.workspaceSubPath, + ), + cacheImage, + pushImage, + pushImages, + pushRegistryInsecure: this.options.pushRegistryInsecure, + cacheRegistryInsecure: + this.options.cacheRegistryInsecure ?? + this.options.pushRegistryInsecure, + buildkitImage: this.options.buildkitImage, + serviceAccountName: this.options.serviceAccountName, + registrySecretName: this.options.registrySecretName, + nodeSelector: this.options.nodeSelector, + tolerations: this.options.tolerations, + }); + }); + await timed("jobCreateMs", () => this.options.kubernetes.createJob(job)); + markJobCreated(); + await timed("recordUpdateMs", () => + this.updateBuild(stored, (next) => { + next.status.jobCreated = true; + }), ); - const job = createBuildJob({ - name: jobName, - namespace: this.options.namespace, - spec: request.spec, - workspaceClaimName: this.options.workspaceClaimName, - workspaceSubPath: jobWorkspaceSubPath( - this.options.workspaceRoot, - stored.spec.workspaceSubPath, - ), - cacheImage, - pushImage, - pushImages, - pushRegistryInsecure: this.options.pushRegistryInsecure, - cacheRegistryInsecure: - this.options.cacheRegistryInsecure ?? - this.options.pushRegistryInsecure, - buildkitImage: this.options.buildkitImage, - serviceAccountName: this.options.serviceAccountName, - registrySecretName: this.options.registrySecretName, - nodeSelector: this.options.nodeSelector, - tolerations: this.options.tolerations, - }); - await this.options.kubernetes.createJob(job); - await this.updateBuild(stored, (next) => { - next.status.jobCreated = true; - }); return initial; } catch (error) { - try { - await this.terminateJob(stored.spec.jobName); - } catch { - // Job cleanup is best effort; the record still releases its lock. - } - await this.failBuild( - stored.metadata.name, - error instanceof Error ? error.message : String(error), - ); + await timed("failureCleanupMs", async () => { + try { + await this.terminateJob(stored.spec.jobName); + } catch { + // Job cleanup is best effort; the record still releases its lock. + } + await this.failBuild( + stored.metadata.name, + error instanceof Error ? error.message : String(error), + ); + }); throw error; } } diff --git a/server/maintenance.ts b/server/maintenance.ts index 1b18d57..9ba7bcc 100644 --- a/server/maintenance.ts +++ b/server/maintenance.ts @@ -27,7 +27,7 @@ export interface MaintenanceLeaseProvider { ): Promise<{ release(): Promise } | undefined>; } -/** Accept DNS hostnames only: no URL syntax, port, address literals, or wildcards. */ +/** Accept DNS hostnames and an optional leading *. only; no URL syntax, port, or address literals. */ export function normalizeMaintenanceHost(value: unknown): string { if (typeof value !== "string") throw new Error("host is required"); const host = value.trim().toLowerCase().replace(/\.$/, ""); @@ -35,7 +35,7 @@ export function normalizeMaintenanceHost(value: unknown): string { host.length === 0 || host.length > 253 || !host.includes(".") || - !/^(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z]{2,63}$/.test(host) + !/^(?:\*\.)?(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z]{2,63}$/.test(host) ) throw new Error("host must be a DNS hostname"); return host; diff --git a/server/materialize.ts b/server/materialize.ts index ca5700b..f5e3240 100644 --- a/server/materialize.ts +++ b/server/materialize.ts @@ -29,6 +29,16 @@ export interface MaterializeCas { get(digest: Sha256Digest): Promise; } +export interface MaterializeTiming { + fileCount?: number; + manifestBytes?: number; + fileBytes?: number; + // Sum of operation durations; concurrent operations can exceed wall time. + // CAS get includes its content hash verification. + casReadMs: number; + fsWriteMs: number; +} + const MATERIALIZE_CONCURRENCY = 20; async function mapConcurrent( @@ -101,8 +111,18 @@ export function parseWorkspaceManifest(data: Uint8Array): WorkspaceManifest { } paths.add(file.path); } + const ancestors = new Set(); for (const path of paths) { - if ([...paths].some((other) => other.startsWith(`${path}/`))) { + for ( + let separator = path.indexOf("/"); + separator !== -1; + separator = path.indexOf("/", separator + 1) + ) { + ancestors.add(path.slice(0, separator)); + } + } + for (const path of paths) { + if (ancestors.has(path)) { throw new Error(`Workspace path conflicts with a directory: ${path}`); } } @@ -125,13 +145,43 @@ export async function materializeWorkspace( cas: MaterializeCas, manifestDigest: Sha256Digest, destination: string, + timing?: MaterializeTiming, ): Promise { + const read = timing + ? async (digest: Sha256Digest) => { + const start = performance.now(); + try { + return await cas.get(digest); + } finally { + timing.casReadMs += performance.now() - start; + } + } + : (digest: Sha256Digest) => cas.get(digest); + const write = timing + ? async (operation: () => Promise): Promise => { + const start = performance.now(); + try { + return await operation(); + } finally { + timing.fsWriteMs += performance.now() - start; + } + } + : (operation: () => Promise) => operation(); assertSha256Digest(manifestDigest); - const manifest = parseWorkspaceManifest(await cas.get(manifestDigest)); + const manifestData = await read(manifestDigest); + if (timing) timing.manifestBytes = manifestData.byteLength; + const manifest = parseWorkspaceManifest(manifestData); + if (timing) { + timing.fileCount = manifest.files.length; + timing.fileBytes = manifest.files.reduce( + (total, file) => total + file.size, + 0, + ); + } - await mkdir(dirname(destination), { recursive: true }); + await write(() => mkdir(dirname(destination), { recursive: true })); try { - await lstat(destination); + await write(() => lstat(destination)); throw new Error(`Workspace destination already exists: ${destination}`); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; @@ -140,14 +190,16 @@ export async function materializeWorkspace( dirname(destination), `.${basename(destination)}.${process.pid}.${randomUUID()}.tmp`, ); - await mkdir(temporary, { mode: 0o755 }); + await write(() => mkdir(temporary, { mode: 0o755 })); try { await mapConcurrent(manifest.files, async (file) => { const target = join(temporary, file.path); if (relative(temporary, target).startsWith("..")) throw new Error("Unsafe workspace path"); - await mkdir(dirname(target), { recursive: true, mode: 0o755 }); - const data = await cas.get(file.digest); + await write(() => + mkdir(dirname(target), { recursive: true, mode: 0o755 }), + ); + const data = await read(file.digest); if (data.byteLength !== file.size) { throw new Error(`Workspace blob size mismatch for ${file.path}`); } @@ -155,24 +207,26 @@ export async function materializeWorkspace( const link = new TextDecoder("utf-8", { fatal: true }).decode(data); if (!safeSymlinkTarget(file.path, link)) throw new Error(`Unsafe symlink target for ${file.path}`); - await symlink(link, target); + await write(() => symlink(link, target)); } else { - const handle = await open( - target, - constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, - file.mode, + const handle = await write(() => + open( + target, + constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, + file.mode, + ), ); try { - await handle.writeFile(data); + await write(() => handle.writeFile(data)); } finally { - await handle.close(); + await write(() => handle.close()); } - await chmod(target, file.mode); + await write(() => chmod(target, file.mode)); } }); - await rename(temporary, destination); + await write(() => rename(temporary, destination)); } catch (error) { - await rm(temporary, { recursive: true, force: true }); + await write(() => rm(temporary, { recursive: true, force: true })); throw error; } return manifest; diff --git a/tests/command/maintenance.test.ts b/tests/command/maintenance.test.ts index 442c37f..351788b 100644 --- a/tests/command/maintenance.test.ts +++ b/tests/command/maintenance.test.ts @@ -61,4 +61,46 @@ describe("maintenance command", () => { output.mockRestore(); } }); + + test("wildcard host requires a second call before posting the normalized host", async () => { + const calls: Array<{ path: string; init?: ApiRequestInit }> = []; + const confirmations: string[] = []; + const output = spyOn(console, "log").mockImplementation(() => {}); + try { + const request = async (path: string, init?: ApiRequestInit): Promise => { + calls.push({ path, init }); + return { + host: "*.example.com", + enabled: init?.method === "POST", + hosts: init?.method === "POST" ? ["*.example.com"] : [], + } as T; + }; + const confirm = async (host: string) => { + confirmations.push(host); + return confirmations.length === 2; + }; + + await runMaintenance("*.EXAMPLE.COM.", request, confirm); + expect(calls).toEqual([{ path: "/maintenance/*.EXAMPLE.COM.", init: undefined }]); + expect(output).toHaveBeenCalledWith( + expect.stringContaining("*.example.com maintenance is currently"), + ); + + await runMaintenance("*.EXAMPLE.COM.", request, confirm); + expect(confirmations).toEqual(["*.example.com", "*.example.com"]); + expect(calls).toEqual([ + { path: "/maintenance/*.EXAMPLE.COM.", init: undefined }, + { path: "/maintenance/*.EXAMPLE.COM.", init: undefined }, + { + path: "/maintenance/*.example.com", + init: { method: "POST", json: { enabled: true } }, + }, + ]); + expect(output).toHaveBeenCalledWith( + "*.example.com maintenance is now \x1b[31mON\x1b[0m", + ); + } finally { + output.mockRestore(); + } + }); }); diff --git a/tests/lib/build-api.test.ts b/tests/lib/build-api.test.ts index 6aa4d90..aa2ed29 100644 --- a/tests/lib/build-api.test.ts +++ b/tests/lib/build-api.test.ts @@ -74,6 +74,7 @@ afterEach(async () => { describe("authenticated build API pipeline", () => { 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; @@ -105,7 +106,14 @@ describe("authenticated build API pipeline", () => { return { complete: true } as T; }; - await uploadWorkspaceSnapshot(snapshot, request); + await uploadWorkspaceSnapshot( + snapshot, + request, + undefined, + undefined, + undefined, + controller.signal, + ); expect( calls.map(({ path, init }) => [init?.method ?? "GET", path]), @@ -120,6 +128,21 @@ describe("authenticated build API pipeline", () => { ["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 () => { diff --git a/tests/lib/request-log.test.ts b/tests/lib/request-log.test.ts index 203b096..dc28516 100644 --- a/tests/lib/request-log.test.ts +++ b/tests/lib/request-log.test.ts @@ -106,6 +106,91 @@ describe("process request lines", () => { } }); + test("only emits compact access lines from the default logger", async () => { + const output = spyOn(console, "log").mockImplementation(() => {}); + const logs: Record[] = []; + const logger = { + log(entry: Record) { + logs.push(entry); + processLogger.log(entry); + }, + }; + try { + const request = new Request( + "https://kuber.astrxl.dev/items?token=query-private", + { + method: "POST", + headers: { + authorization: "Bearer header-private", + cookie: "session=cookie-private", + "x-api-key": "key-private", + }, + body: "body-private", + }, + ); + await logServerRequest( + request, + "request-private", + () => { + processLogger.log({ + event: "kuber.server.request.body", + headers: Object.fromEntries(request.headers), + body: { encoding: "base64", data: "Ym9keS1wcml2YXRl" }, + }); + return new Response(null, { + status: 201, + headers: { "set-cookie": "session=response-private" }, + }); + }, + logger, + ); + await expect( + logServerRequest( + new Request("https://kuber.astrxl.dev/fail", { + headers: { authorization: "Bearer failure-private" }, + }), + "failure-request", + () => { + throw new Error("failed"); + }, + logger, + ), + ).rejects.toThrow("failed"); + await logServerRequest( + new Request("https://kuber.astrxl.dev/items", { + headers: { cookie: "session=get-private" }, + }), + "get-request", + () => new Response(null, { status: 200 }), + logger, + ); + + expect(output.mock.calls.map(([line]) => line)).toEqual([ + "PST /items 201", + "GET /fail ERR", + "GET /items 200", + ]); + expect(logs).toContainEqual( + expect.objectContaining({ + event: "kuber.server.request.start", + headers: expect.objectContaining({ + authorization: "Bearer header-private", + }), + }), + ); + expect(logs).toContainEqual( + expect.objectContaining({ + event: "kuber.server.request.end", + responseHeaders: expect.objectContaining({ + "set-cookie": "session=response-private", + }), + }), + ); + } finally { + output.mockRestore(); + } + }); + test("indents Kubernetes lines for incoming requests without a CLI marker", async () => { const output = spyOn(console, "log").mockImplementation(() => {}); try { diff --git a/tests/server/build-controller.test.ts b/tests/server/build-controller.test.ts index 0a9fb99..76afad5 100644 --- a/tests/server/build-controller.test.ts +++ b/tests/server/build-controller.test.ts @@ -236,6 +236,155 @@ async function internalFixture( } describe("build controller", () => { + 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((resolve) => { + release = resolve; + }); + const started = new Promise((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((resolve) => { + release = resolve; + }); + const started = new Promise((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( diff --git a/tests/server/maintenance.test.ts b/tests/server/maintenance.test.ts index a760353..12b1821 100644 --- a/tests/server/maintenance.test.ts +++ b/tests/server/maintenance.test.ts @@ -40,16 +40,25 @@ class MemoryPersistence implements MaintenancePersistence { const lease = { acquire: async () => ({ release: async () => {} }) }; describe("maintenance override", () => { - test("normalizes only DNS hostnames", () => { + test("normalizes DNS hostnames and leading wildcards only", () => { expect(normalizeMaintenanceHost(" Sub.Domain.COM. ")).toBe( "sub.domain.com", ); + expect(normalizeMaintenanceHost(" *.EXAMPLE.COM. ")).toBe("*.example.com"); + expect(normalizeMaintenanceHost(" *.Sub.Example.COM. ")).toBe("*.sub.example.com"); for (const host of [ "http://example.com", "example.com:443", "127.0.0.1", "[::1]", - "*.example.com", + "*", + "*.", + "example.*.com", + "a.*.example.com", + "**.example.com", + "*.*.example.com", + "*.example.com:443", + "*.127.0.0.1", "example", "a..com", ]) @@ -74,6 +83,28 @@ describe("maintenance override", () => { expect(persistence.state).toEqual({ hosts: [] }); }); + test("persists normalized wildcards and renders them as Host rules", async () => { + const persistence = new MemoryPersistence(); + const service = new MaintenanceService(persistence, lease); + expect(await service.set(" *.EXAMPLE.COM. ", true)).toMatchObject({ + host: "*.example.com", + enabled: true, + hosts: ["*.example.com"], + }); + expect(persistence.state).toEqual({ hosts: ["*.example.com"] }); + expect(persistence.resources.at(-1)).toMatchObject({ + spec: { routes: [{ match: "Host(`*.example.com`)" }] }, + }); + expect(await service.status("*.EXAMPLE.COM.")).toMatchObject({ + host: "*.example.com", + enabled: true, + }); + expect(await service.set("*.example.com", false)).toMatchObject({ + enabled: false, + hosts: [], + }); + }); + test("renders the single shared error route and rewrite middleware", () => { expect(maintenanceMiddleware()).toMatchObject({ metadata: { name: "maintenance-override", namespace: "routing" }, @@ -100,6 +131,9 @@ describe("maintenance override", () => { ], }, }); + expect(maintenanceRoute(["*.example.com", "a.example.com"])).toMatchObject({ + spec: { routes: [{ match: "Host(`*.example.com`) || Host(`a.example.com`)" }] }, + }); }); test("maps an unavailable global lease to a retryable API conflict", async () => { diff --git a/tests/server/materialize.test.ts b/tests/server/materialize.test.ts index 4223242..d5446e4 100644 --- a/tests/server/materialize.test.ts +++ b/tests/server/materialize.test.ts @@ -26,6 +26,59 @@ function digest(data: Uint8Array | string): Sha256Digest { } describe("source materialization", () => { + test("collects aggregate CAS and filesystem timing with manifest counts", async () => { + const root = await mkdtemp(join(tmpdir(), "kuber-materialize-")); + roots.push(root); + const data = Buffer.from("hello"); + const manifest: WorkspaceManifest = { + version: BUILD_PROTOCOL_VERSION, + files: [ + { + path: "example", + type: "file", + digest: digest(data), + size: data.byteLength, + mode: 0o644, + }, + ], + }; + const manifestData = Buffer.from(JSON.stringify(manifest)); + const workspace = digest(manifestData); + let release!: () => void; + let entered!: () => void; + const blocked = new Promise((resolve) => { + release = resolve; + }); + const started = new Promise((resolve) => { + entered = resolve; + }); + const timing = { casReadMs: 0, fsWriteMs: 0 }; + const submission = materializeWorkspace( + { + get: async (requested) => { + if (requested === workspace) return manifestData; + entered(); + await blocked; + return data; + }, + }, + workspace, + join(root, "workspace"), + timing, + ); + await started; + await Bun.sleep(25); + release(); + await submission; + expect(timing).toMatchObject({ + fileCount: 1, + manifestBytes: manifestData.byteLength, + fileBytes: data.byteLength, + }); + expect(timing.casReadMs).toBeGreaterThan(15); + expect(timing.fsWriteMs).toBeGreaterThan(0); + }); + test("bounds concurrent CAS reads while materializing many files", async () => { const root = await mkdtemp(join(tmpdir(), "kuber-materialize-")); roots.push(root); @@ -151,6 +204,28 @@ describe("source materialization", () => { ), ), ).toThrow("conflicts"); + for (const [files, conflict] of [ + [["a/b/c", "a"], "Workspace path conflicts with a directory: a"], + [["a/b/c", "a/b"], "Workspace path conflicts with a directory: a/b"], + [["a", "a/b/c"], "Workspace path conflicts with a file: a/b/c"], + ] as const) { + expect(() => + parseWorkspaceManifest( + Buffer.from( + JSON.stringify({ + version: BUILD_PROTOCOL_VERSION, + files: files.map((path) => ({ + path, + type: "file", + digest: `sha256:${"a".repeat(64)}`, + size: 0, + mode: 0o644, + })), + }), + ), + ), + ).toThrow(conflict); + } const root = await mkdtemp(join(tmpdir(), "kuber-materialize-")); roots.push(root); @@ -178,4 +253,20 @@ describe("source materialization", () => { ).rejects.toThrow("Unsafe symlink"); await expect(lstat(destination)).rejects.toMatchObject({ code: "ENOENT" }); }); + + test("accepts a large manifest with distinct nested paths", () => { + const manifest: WorkspaceManifest = { + version: BUILD_PROTOCOL_VERSION, + files: Array.from({ length: 3_000 }, (_, index) => ({ + path: `dir-${index}/nested/file`, + type: "file" as const, + digest: `sha256:${"a".repeat(64)}` as Sha256Digest, + size: 0, + mode: 0o644 as const, + })), + }; + expect( + parseWorkspaceManifest(Buffer.from(JSON.stringify(manifest))), + ).toEqual(manifest); + }); });