From e3df2f4d7675d262e7e32acf80c2cb828b9a710c Mon Sep 17 00:00:00 2001 From: dmgnr Date: Fri, 11 Sep 2026 10:51:57 +0000 Subject: [PATCH] feat: improve resource reconciliation and observability --- README.md | 13 + command/main.ts | 2 +- lib/apply.ts | 3 +- lib/convert.ts | 151 +++++++++ lib/k8s-http.ts | 64 +++- lib/kubernetes-object.ts | 26 ++ lib/request-log.ts | 226 ++++++++++++ package.json | 2 +- server/app.ts | 182 ++++++---- server/build-controller.ts | 317 ++++++++++++++--- server/build-job.ts | 21 +- server/build-kubernetes.ts | 103 +++++- server/build-reconciler.ts | 84 +++++ server/build-store.ts | 108 +++++- server/index.ts | 59 +++- server/kubernetes-exec.ts | 92 ++--- server/kubernetes-logs.ts | 52 ++- server/kubernetes-state.ts | 3 +- tests/lib/apply.test.ts | 63 +++- tests/lib/convert-deployment.test.ts | 142 ++++++++ tests/lib/database.test.ts | 55 +++ tests/lib/k8s-http.test.ts | 45 +++ tests/lib/request-log.test.ts | 65 ++++ tests/server/app.test.ts | 46 +++ tests/server/build-controller.test.ts | 471 +++++++++++++++++++++++++- tests/server/build-job.test.ts | 61 ++++ tests/server/build-kubernetes.test.ts | 65 ++++ tests/server/build-reconciler.test.ts | 103 ++++++ 28 files changed, 2407 insertions(+), 217 deletions(-) create mode 100644 lib/kubernetes-object.ts create mode 100644 lib/request-log.ts create mode 100644 server/build-reconciler.ts create mode 100644 tests/lib/request-log.test.ts create mode 100644 tests/server/build-reconciler.test.ts diff --git a/README.md b/README.md index db3a031..882801b 100644 --- a/README.md +++ b/README.md @@ -253,6 +253,19 @@ variables: registry access. - `KUBER_BUILDKIT_IMAGE`: override the BuildKit runner image used for builds. - `KUBER_BUILD_DATA_CLAIM`: the PVC claim backing builds. +- `KUBER_BUILD_RECONCILE_MS` (default `30000`): the background build + reconciliation interval in milliseconds. Values must be greater than zero and + no greater than `2147483647` (the JavaScript timer maximum); invalid values + use the default. +- `KUBER_BUILD_RECONCILE_TIMEOUT_MS` (default `20000`, maximum `25000`): the + deadline for one background build-reconciliation scan in milliseconds. The + maximum leaves time within the 30-second reconciliation lease for normal + renewal or release. The reconciler renews both its workspace and per-build + leases every 10 seconds for the lifetime of a generation, including while an + observation, log read, or store call is pending. On timeout its abort signal + is cancelled and logged, but Kubernetes requests may be unabortable; the + generation remains active and continues its lease heartbeat until that request + settles, so a later scan cannot overlap it. See [Server Deployment](#server-deployment) for how `compose.yml` wires the internal registry variables for the kuber-server pod. diff --git a/command/main.ts b/command/main.ts index 487434e..c9e630c 100644 --- a/command/main.ts +++ b/command/main.ts @@ -22,7 +22,7 @@ import { trust } from "./trust"; export const main = defineCommand({ meta: { name: "kuber", - version: "2.3.4", + version: "2.4.0", description: "Docker Compose -> K8s translation layer", }, args: { diff --git a/lib/apply.ts b/lib/apply.ts index ccdf40b..4889f64 100644 --- a/lib/apply.ts +++ b/lib/apply.ts @@ -2,6 +2,7 @@ import { PatchStrategy } from "@kubernetes/client-node"; import type { KubernetesObject } from "@kubernetes/client-node"; import { LABELS } from "../const"; import { objectApi } from "./k8s"; +import { sanitizeKubernetesObject } from "./kubernetes-object"; const FIELD_MANAGER = "kuber"; @@ -94,7 +95,7 @@ export async function applyResource( resource: T, ): Promise { return objectApi.patch( - resource, + sanitizeKubernetesObject(resource), undefined, undefined, FIELD_MANAGER, diff --git a/lib/convert.ts b/lib/convert.ts index 9b1437d..efebfa2 100644 --- a/lib/convert.ts +++ b/lib/convert.ts @@ -8,6 +8,7 @@ import type { V1Namespace, V1PersistentVolumeClaim, V1PodSpec, + V1ResourceRequirements, V1Secret, V1Service, V1ServicePort, @@ -228,6 +229,155 @@ function toHostname(value: string | undefined): string | undefined { return /[a-z]/i.test(value) ? value : undefined; } +const CPU_DECIMAL_PATTERN = /^(\d+)(?:\.(\d+))?$/; +const MEMORY_QUANTITY_PATTERN = /^(?:\d+(?:\.\d+)?|\.\d+)([a-zA-Z]*)$/; + +function resourceValue(value: unknown): string { + return typeof value === "string" ? JSON.stringify(value) : String(value); +} + +function toKubernetesCpu(value: number | string, source: string): string { + const text = String(value).trim(); + const match = CPU_DECIMAL_PATTERN.exec(text); + if (!match) { + throw new Error( + `Invalid ${source} value ${resourceValue(value)}. Expected a non-negative finite decimal number of CPU cores.`, + ); + } + + const whole = BigInt(match[1] ?? "0"); + const fraction = match[2] ?? ""; + const millicores = + whole * 1000n + + (fraction ? BigInt(fraction.padEnd(3, "0").slice(0, 3)) : 0n); + if (fraction.slice(3).replaceAll("0", "") !== "") { + throw new Error( + `Invalid ${source} value ${resourceValue(value)}. CPU cores must be representable in whole millicores.`, + ); + } + + return millicores % 1000n === 0n + ? String(millicores / 1000n) + : `${millicores}m`; +} + +function toKubernetesMemory(value: string, source: string): string { + const text = value.trim(); + const match = MEMORY_QUANTITY_PATTERN.exec(text); + if (!match) { + throw new Error( + `Invalid ${source} value ${resourceValue(value)}. Expected a non-negative memory quantity.`, + ); + } + + const suffix = match[1] ?? ""; + const normalizedSuffix = (() => { + switch (suffix.toLowerCase()) { + case "b": + return ""; + case "k": + case "kb": + return "k"; + case "m": + case "mb": + // Compose uses lowercase m for decimal megabytes; Kubernetes uses it for milli. + return "M"; + case "g": + case "gb": + return "G"; + case "t": + case "tb": + return "T"; + case "p": + case "pb": + return "P"; + case "e": + case "eb": + return "E"; + default: + return suffix; + } + })(); + + if ( + ![ + "", + "n", + "u", + "m", + "k", + "M", + "G", + "T", + "P", + "E", + "Ki", + "Mi", + "Gi", + "Ti", + "Pi", + "Ei", + ].includes(normalizedSuffix) + ) { + throw new Error( + `Invalid ${source} value ${resourceValue(value)}. Expected a valid Kubernetes or Compose memory quantity.`, + ); + } + + return `${text.slice(0, text.length - suffix.length)}${normalizedSuffix}`; +} + +function toContainerResources( + service: Service, +): V1ResourceRequirements | undefined { + const resources = service.deploy?.resources; + if (!resources) return; + + const limits = { + ...(resources.limits?.cpus !== undefined + ? { + cpu: toKubernetesCpu( + resources.limits.cpus, + "deploy.resources.limits.cpus", + ), + } + : {}), + ...(resources.limits?.memory !== undefined + ? { + memory: toKubernetesMemory( + resources.limits.memory, + "deploy.resources.limits.memory", + ), + } + : {}), + }; + const requests = { + ...(resources.reservations?.cpus !== undefined + ? { + cpu: toKubernetesCpu( + resources.reservations.cpus, + "deploy.resources.reservations.cpus", + ), + } + : {}), + ...(resources.reservations?.memory !== undefined + ? { + memory: toKubernetesMemory( + resources.reservations.memory, + "deploy.resources.reservations.memory", + ), + } + : {}), + }; + + if (Object.keys(limits).length === 0 && Object.keys(requests).length === 0) + return; + return { + ...(Object.keys(limits).length > 0 ? { limits } : {}), + ...(Object.keys(requests).length > 0 ? { requests } : {}), + }; +} + function toBoolean(value: boolean | string | undefined): boolean | undefined { if (typeof value === "boolean") return value; if (!value) return; @@ -1002,6 +1152,7 @@ export function serviceToDeployment( ? [{ secretRef: { name: `${name}-env` } }] : undefined, ports: toContainerPorts(ports), + resources: toContainerResources(service), volumeMounts: toVolumeMounts(mounts), } satisfies V1Container, service["x-container"] ?? {}, diff --git a/lib/k8s-http.ts b/lib/k8s-http.ts index ab2962c..02215ca 100644 --- a/lib/k8s-http.ts +++ b/lib/k8s-http.ts @@ -6,6 +6,12 @@ import { ResponseContext } from "@kubernetes/client-node/dist/gen/http/http.js"; import { from } from "@kubernetes/client-node/dist/gen/rxjsStub.js"; import http from "node:http"; import https from "node:https"; +import { + logValue, + processLogger, + safeLog, + type ProcessLogger, +} from "./request-log"; type TransportOptions = { maxConcurrent: number; @@ -14,6 +20,7 @@ type TransportOptions = { baseRetryMs: number; maxRetryMs: number; random: () => number; + logger: ProcessLogger; }; const DefaultTransportOptions: TransportOptions = { @@ -23,6 +30,7 @@ const DefaultTransportOptions: TransportOptions = { baseRetryMs: 250, maxRetryMs: 5000, random: Math.random, + logger: processLogger, }; function abortError(): Error { @@ -129,7 +137,10 @@ class RequestLimiter { } } -function sendOnce(request: RequestContext): Promise { +function sendOnce( + request: RequestContext, + onComplete: (response: ResponseContext, body: string) => void, +): Promise { return new Promise((resolve, reject) => { const url = new URL(request.getUrl()); const transport = url.protocol === "http:" ? http : https; @@ -158,12 +169,16 @@ function sendOnce(request: RequestContext): Promise { else if (value !== undefined) headers[key] = String(value); } - resolve( - new ResponseContext(response.statusCode ?? 0, headers, { + const context = new ResponseContext( + response.statusCode ?? 0, + headers, + { text: async () => buffer.toString("utf8"), binary: async () => buffer, - }), + }, ); + onComplete(context, buffer.toString("utf8")); + resolve(context); }); }, ); @@ -198,12 +213,47 @@ export function createKubernetesHttpLibrary( send(request) { const result = (async () => { for (let attempt = 0; ; attempt += 1) { - const release = await limiter.acquire(request.getSignal()); + const startedAt = Date.now(); + const method = request.getHttpMethod().toString(); + const url = request.getUrl(); + const base = { + type: "kubernetes.request", + method, + url, + pathname: new URL(url).pathname, + headers: logValue(request.getHeaders()), + body: logValue(request.getBody()), + attempt: attempt + 1, + maxRetries: options.maxRetries, + }; + safeLog(options.logger, { + event: "kuber.k8s.request.start", + ...base, + }); + let release: (() => void) | undefined; let response: ResponseContext; try { - response = await sendOnce(request); + release = await limiter.acquire(request.getSignal()); + response = await sendOnce(request, (completed, body) => { + safeLog(options.logger, { + event: "kuber.k8s.request.end", + ...base, + status: completed.httpStatusCode, + responseHeaders: logValue(completed.headers), + responseBody: body, + durationMs: Date.now() - startedAt, + }); + }); + } catch (error) { + safeLog(options.logger, { + event: "kuber.k8s.request.failed", + ...base, + durationMs: Date.now() - startedAt, + error: logValue(error), + }); + throw error; } finally { - release(); + release?.(); } if ( diff --git a/lib/kubernetes-object.ts b/lib/kubernetes-object.ts new file mode 100644 index 0000000..922da3f --- /dev/null +++ b/lib/kubernetes-object.ts @@ -0,0 +1,26 @@ +import type { KubernetesObject } from "@kubernetes/client-node"; + +type KubernetesObjectWithStatus = KubernetesObject & { status?: unknown }; + +/** Returns an apply-safe copy without fields owned by the API server. */ +export function sanitizeKubernetesObject( + resource: T, +): T { + const { status: _status, ...object } = resource as KubernetesObjectWithStatus; + const { + managedFields: _managedFields, + resourceVersion: _resourceVersion, + uid: _uid, + creationTimestamp: _creationTimestamp, + generation: _generation, + selfLink: _selfLink, + deletionTimestamp: _deletionTimestamp, + deletionGracePeriodSeconds: _deletionGracePeriodSeconds, + ...metadata + } = object.metadata ?? {}; + + return { + ...object, + ...(object.metadata && { metadata }), + } as T; +} diff --git a/lib/request-log.ts b/lib/request-log.ts new file mode 100644 index 0000000..a859f65 --- /dev/null +++ b/lib/request-log.ts @@ -0,0 +1,226 @@ +export type ProcessLogEntry = Record; + +export type ProcessLogger = { + log(entry: ProcessLogEntry): void; +}; + +const MAX_RESPONSE_LOG_BYTES = 64 * 1024; + +export function logValue( + value: unknown, + seen = new WeakSet(), +): unknown { + if ( + value === null || + typeof value === "string" || + typeof value === "number" || + typeof value === "boolean" + ) + return value; + if (typeof value === "undefined") return undefined; + if (typeof value === "bigint") return value.toString(); + if (typeof value === "symbol" || typeof value === "function") + return String(value); + if (value instanceof Uint8Array) + return { encoding: "base64", data: Buffer.from(value).toString("base64") }; + if (value instanceof Error) + return { + name: value.name, + message: value.message, + ...(value.stack && { stack: value.stack }), + }; + if (typeof value !== "object") return String(value); + if (seen.has(value)) return "[Circular]"; + seen.add(value); + if (Array.isArray(value)) return value.map((item) => logValue(item, seen)); + try { + return Object.fromEntries( + Object.entries(value).map(([key, item]) => [key, logValue(item, seen)]), + ); + } catch (error) { + return { unserializable: String(error) }; + } +} + +export const processLogger: ProcessLogger = { + log(entry) { + try { + console.log(`KUBER_REQUEST ${JSON.stringify(logValue(entry))}`); + } catch { + // Process logging must never change request behavior. + } + }, +}; + +export function safeLog(logger: ProcessLogger, entry: ProcessLogEntry): void { + try { + logger.log(entry); + } catch { + // Logging is deliberately best-effort and process-only. + } +} + +function headers(headers: Headers): Record { + return Object.fromEntries(headers); +} + +function isStreamingContentType(value: string | null): boolean { + const contentType = value?.split(";", 1)[0]?.trim().toLowerCase(); + return ( + contentType === "text/event-stream" || + contentType === "application/x-ndjson" || + contentType === "application/ndjson" || + contentType === "application/json-seq" + ); +} + +function responseBodySize(value: Response): number | undefined { + const contentLength = value.headers.get("content-length"); + if (!contentLength || !/^\d+$/.test(contentLength)) return; + const size = Number(contentLength); + return Number.isSafeInteger(size) ? size : undefined; +} + +async function readBoundedResponseBody( + body: ReadableStream, +): Promise { + const reader = body.getReader(); + const chunks: Uint8Array[] = []; + let size = 0; + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + size += value.byteLength; + if (size > MAX_RESPONSE_LOG_BYTES) { + void reader.cancel(); + throw new RangeError("Response log body exceeds limit"); + } + chunks.push(value); + } + const result = new Uint8Array(size); + let offset = 0; + for (const chunk of chunks) { + result.set(chunk, offset); + offset += chunk.byteLength; + } + return result; +} + +function logResponseBody( + value: Response, + logger: ProcessLogger, + entry: ProcessLogEntry, +): void { + if (!value.body) return; + + if (isStreamingContentType(value.headers.get("content-type"))) { + safeLog(logger, { ...entry, bodySkipped: true, streaming: true }); + return; + } + + const size = responseBodySize(value); + if (size === undefined) { + safeLog(logger, { ...entry, bodySkipped: true, bodySizeUnknown: true }); + return; + } + if (size > MAX_RESPONSE_LOG_BYTES) { + safeLog(logger, { + ...entry, + bodySkipped: true, + bodyTooLarge: true, + contentLength: size, + }); + return; + } + + try { + const body = value.clone().body; + if (!body) return; + void readBoundedResponseBody(body) + .then((body) => safeLog(logger, { ...entry, body: logValue(body) })) + .catch((error) => + safeLog( + logger, + error instanceof RangeError + ? { ...entry, bodySkipped: true, bodyTooLarge: true } + : { ...entry, error: logValue(error) }, + ), + ); + } catch (error) { + safeLog(logger, { ...entry, error: logValue(error) }); + } +} + +export async function logServerRequest( + request: Request, + requestId: string, + handler: () => Promise | T, + logger: ProcessLogger = processLogger, +): Promise { + const startedAt = Date.now(); + const url = new URL(request.url); + const base = { + type: "server.request", + requestId, + method: request.method, + url: request.url, + pathname: url.pathname, + headers: headers(request.headers), + }; + safeLog(logger, { event: "kuber.server.request.start", ...base }); + try { + const response = await handler(); + safeLog(logger, { + event: "kuber.server.request.end", + ...base, + status: response instanceof Response ? response.status : 101, + durationMs: Date.now() - startedAt, + ...(response instanceof Response && { + responseHeaders: headers(response.headers), + }), + }); + if (response instanceof Response) + logResponseBody(response, logger, { + event: "kuber.server.response.body", + ...base, + status: response.status, + }); + return response; + } catch (error) { + safeLog(logger, { + event: "kuber.server.request.failed", + ...base, + durationMs: Date.now() - startedAt, + error: logValue(error), + }); + throw error; + } +} + +export async function logKubernetesRequest( + request: Omit, + handler: () => Promise, + logger: ProcessLogger = processLogger, +): Promise { + const startedAt = Date.now(); + const base = { type: "kubernetes.request", ...request }; + safeLog(logger, { event: "kuber.k8s.request.start", ...base }); + try { + const result = await handler(); + safeLog(logger, { + event: "kuber.k8s.request.end", + ...base, + status: 200, + durationMs: Date.now() - startedAt, + }); + return result; + } catch (error) { + safeLog(logger, { + event: "kuber.k8s.request.failed", + ...base, + durationMs: Date.now() - startedAt, + error: logValue(error), + }); + throw error; + } +} diff --git a/package.json b/package.json index a9060bf..4db522f 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.3.4", + "version": "2.4.0", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/server/app.ts b/server/app.ts index 34fe039..a9e09d1 100644 --- a/server/app.ts +++ b/server/app.ts @@ -65,6 +65,13 @@ import { type WorkspaceStore, } from "./workspace-store"; import { redactString } from "./redact"; +import { + logValue, + logServerRequest, + processLogger, + safeLog, + type ProcessLogEntry, +} from "../lib/request-log"; const API_PREFIX = "/api/v2"; const RUNTIME_SESSION_MS = 24 * 60 * 60 * 1000; @@ -132,10 +139,12 @@ export type RequestErrorLog = Omit & { export interface AppLogger { error(entry: RequestErrorLog): void; + log?(entry: ProcessLogEntry): void; } const defaultAppLogger: AppLogger = { error: (entry) => console.error(JSON.stringify(entry)), + log: processLogger.log, }; type LoginFailures = { count: number; resetAt: number }; @@ -428,6 +437,22 @@ export function createApp( }); } + function logRequestBody(request: Request, body: Uint8Array): void { + safeLog( + { log: (entry) => (logger.log ?? processLogger.log)(entry) }, + { + event: "kuber.server.request.body", + type: "server.request", + requestId: requestIds.get(request) ?? makeRequestId(), + method: request.method, + url: request.url, + pathname: new URL(request.url).pathname, + headers: Object.fromEntries(request.headers), + body: logValue(body), + }, + ); + } + async function readJson( request: Request, limit = bodyLimit, @@ -478,6 +503,7 @@ export function createApp( offset += chunk.byteLength; } const text = new TextDecoder().decode(payload); + logRequestBody(request, payload); let value: unknown; try { value = text ? JSON.parse(text) : {}; @@ -537,6 +563,7 @@ export function createApp( result.set(chunk, offset); offset += chunk.byteLength; } + logRequestBody(request, result); return result; } @@ -2331,80 +2358,87 @@ export function createApp( suppliedRequestId && /^[\x21-\x7e]{1,128}$/.test(suppliedRequestId) ? suppliedRequestId : makeRequestId(); - requestIds.set(request, requestId); - const url = new URL(request.url); - const origin = request.headers.get("origin"); - const originAllowed = - !origin || - (hasOriginConfiguration - ? allowedOrigins.has(origin) - : origin === url.origin); - let result: Response; - try { - if (!originAllowed) - throw new HttpError( - 403, - "Forbidden", - "ORIGIN_NOT_ALLOWED", - "Request origin is not allowed", - ); - if (request.method === "OPTIONS") { - if (!origin) - throw new HttpError( - 400, - "Bad request", - "ORIGIN_REQUIRED", - "Origin header is required for preflight", - ); - result = new Response(null, { - status: 204, - headers: { - "access-control-allow-methods": - "GET, POST, PUT, PATCH, DELETE, OPTIONS", - "access-control-allow-headers": - "Authorization, Content-Type, Idempotency-Key, If-Match, Upload-Offset, X-Request-Id", - "access-control-max-age": "600", - }, - }); - } else if ( - request.method === "GET" && - url.pathname === `${API_PREFIX}/health` - ) { - result = response({ status: "ok" }); - } else if ( - request.method === "POST" && - url.pathname === `${API_PREFIX}/login` - ) { - result = await handleLogin(request); - } else { - const identity = await authenticate(request); - if (!identity) - throw new HttpError( - 401, - "Unauthorized", - "UNAUTHORIZED", - "A valid kuber login is required", - { "www-authenticate": "Bearer" }, - ); - result = await handleAuthenticated(request, url, identity); - } - } catch (error) { - const normalized = normalizeError(error); - if (normalized.code === "INTERNAL_ERROR") - logRequestError(request, error, { - event: "request.failed", - code: normalized.code, - message: normalized.message, - workspaceId: workspaceIdFromPath(url.pathname), - }); - result = problem(normalized, requestId); - } - result.headers.set("x-request-id", requestId); - if (origin && originAllowed) { - result.headers.set("access-control-allow-origin", origin); - result.headers.set("vary", "Origin"); - } - return result; + return logServerRequest( + request, + requestId, + async () => { + requestIds.set(request, requestId); + const url = new URL(request.url); + const origin = request.headers.get("origin"); + const originAllowed = + !origin || + (hasOriginConfiguration + ? allowedOrigins.has(origin) + : origin === url.origin); + let result: Response; + try { + if (!originAllowed) + throw new HttpError( + 403, + "Forbidden", + "ORIGIN_NOT_ALLOWED", + "Request origin is not allowed", + ); + if (request.method === "OPTIONS") { + if (!origin) + throw new HttpError( + 400, + "Bad request", + "ORIGIN_REQUIRED", + "Origin header is required for preflight", + ); + result = new Response(null, { + status: 204, + headers: { + "access-control-allow-methods": + "GET, POST, PUT, PATCH, DELETE, OPTIONS", + "access-control-allow-headers": + "Authorization, Content-Type, Idempotency-Key, If-Match, Upload-Offset, X-Request-Id", + "access-control-max-age": "600", + }, + }); + } else if ( + request.method === "GET" && + url.pathname === `${API_PREFIX}/health` + ) { + result = response({ status: "ok" }); + } else if ( + request.method === "POST" && + url.pathname === `${API_PREFIX}/login` + ) { + result = await handleLogin(request); + } else { + const identity = await authenticate(request); + if (!identity) + throw new HttpError( + 401, + "Unauthorized", + "UNAUTHORIZED", + "A valid kuber login is required", + { "www-authenticate": "Bearer" }, + ); + result = await handleAuthenticated(request, url, identity); + } + } catch (error) { + const normalized = normalizeError(error); + if (normalized.code === "INTERNAL_ERROR") + logRequestError(request, error, { + event: "request.failed", + code: normalized.code, + message: normalized.message, + workspaceId: workspaceIdFromPath(url.pathname), + }); + result = problem(normalized, requestId); + } + result.headers.set("x-request-id", requestId); + if (origin && originAllowed) { + result.headers.set("access-control-allow-origin", origin); + result.headers.set("vary", "Origin"); + } + return result; + }, + { log: (entry) => (logger.log ?? processLogger.log)(entry) }, + ); }; } diff --git a/server/build-controller.ts b/server/build-controller.ts index 45e3a9c..76820f5 100644 --- a/server/build-controller.ts +++ b/server/build-controller.ts @@ -1,4 +1,4 @@ -import { createHash } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import { rm } from "node:fs/promises"; import { join } from "node:path"; import { @@ -14,9 +14,11 @@ import { BUILD_RECORD_API_VERSION, BuildStoreConflictError, type BuildRecord, + type BuildReconciliationLease, type BuildStore, type UploadRecord, } from "./build-store"; +import type { WorkspaceLease, WorkspaceLeaseProvider } from "./operation-store"; import { materializeWorkspace, parseWorkspaceManifest, @@ -28,6 +30,8 @@ export const DEFAULT_MAX_BLOB_BYTES = 1024 * 1024 * 1024; export const DEFAULT_MAX_UPLOAD_CHUNK_BYTES = 8 * 1024 * 1024; const SNAPSHOT_CHECK_CONCURRENCY = 20; export const DEFAULT_MAX_LOG_BYTES = 1024 * 1024; +const BUILD_RECONCILE_LEASE_TTL_MS = 30_000; +const BUILD_RECONCILE_HEARTBEAT_MS = BUILD_RECONCILE_LEASE_TTL_MS / 3; export interface BuildCas extends MaterializeCas { has(digest: Sha256Digest): Promise; @@ -53,6 +57,13 @@ export interface BuildKubernetesOperations { deleteJob(namespace: string, name: string): Promise; } +type ReconcileOptions = { signal?: AbortSignal }; + +function throwIfAborted(signal?: AbortSignal): void { + if (signal?.aborted) + throw new Error("Build reconciliation scan was cancelled"); +} + export interface BuildControllerOptions { cas: BuildCas; store: BuildStore; @@ -76,6 +87,12 @@ export interface BuildControllerOptions { now?: () => Date; materialize?: typeof materializeWorkspace; resolveDigest?: typeof resolveRegistryDigest; + onReconcileFailure?: (id: string, error: unknown) => void; + reconcileLeases?: WorkspaceLeaseProvider; + reconcileHolder?: string; + reconcileHeartbeatMs?: number; + reconcileSetTimeout?: typeof setTimeout; + reconcileClearTimeout?: typeof clearTimeout; } export class BuildValidationError extends Error { @@ -204,6 +221,8 @@ export class BuildController { private readonly materializer: typeof materializeWorkspace; private readonly digestResolver: typeof resolveRegistryDigest; private readonly imageSubmissionLocks = new Map>(); + private readonly reconcileLocks = new Map>(); + private readonly reconcileHolder: string; constructor(private readonly options: BuildControllerOptions) { this.now = options.now ?? (() => new Date()); @@ -213,6 +232,7 @@ export class BuildController { this.maxLogBytes = options.maxLogBytes ?? DEFAULT_MAX_LOG_BYTES; this.materializer = options.materialize ?? materializeWorkspace; this.digestResolver = options.resolveDigest ?? resolveRegistryDigest; + this.reconcileHolder = options.reconcileHolder ?? randomUUID(); } async negotiateSnapshot( @@ -534,47 +554,215 @@ export class BuildController { ); } - async reconcileBuild(id: string): Promise { + async reconcileBuild( + id: string, + options: ReconcileOptions = {}, + ): Promise { + const current = this.reconcileLocks.get(id); + if (current) return current; + const reconciliation = this.reconcileBuildWithLease(id, options.signal); + this.reconcileLocks.set(id, reconciliation); + try { + return await reconciliation; + } finally { + if (this.reconcileLocks.get(id) === reconciliation) + this.reconcileLocks.delete(id); + } + } + + private async reconcileBuildWithLease( + id: string, + signal?: AbortSignal, + ): Promise { + throwIfAborted(signal); + let lease: WorkspaceLease | undefined; + let reconcileLease: BuildReconciliationLease | undefined; + let heartbeat: ReturnType | undefined; + let heartbeatActive = false; + let leaseOwnershipLost = false; + try { + lease = this.options.reconcileLeases + ? await this.options.reconcileLeases.acquire( + `build-reconcile:${id}`, + this.reconcileHolder, + BUILD_RECONCILE_LEASE_TTL_MS, + ) + : undefined; + throwIfAborted(signal); + if (this.options.reconcileLeases && !lease) + return this.getBuildStatus(id); + reconcileLease = await this.options.store.acquireBuildReconciliationLease( + id, + this.reconcileHolder, + BUILD_RECONCILE_LEASE_TTL_MS, + ); + throwIfAborted(signal); + if (!reconcileLease) return this.getBuildStatus(id); + const reconcileLeaseToken = reconcileLease.token; + const requireLeaseOwnership = async () => { + throwIfAborted(signal); + if (leaseOwnershipLost) + throw new BuildConflictError( + "Build reconciliation lease ownership was lost", + ); + const ownsLease = lease + ? await lease.renew(BUILD_RECONCILE_LEASE_TTL_MS) + : true; + const ownsReconcileLease = + await this.options.store.renewBuildReconciliationLease( + id, + reconcileLeaseToken, + BUILD_RECONCILE_LEASE_TTL_MS, + ); + throwIfAborted(signal); + if (!ownsLease || !ownsReconcileLease) { + leaseOwnershipLost = true; + throw new BuildConflictError( + "Build reconciliation lease ownership was lost", + ); + } + }; + const heartbeatMs = Math.min( + BUILD_RECONCILE_HEARTBEAT_MS, + this.options.reconcileHeartbeatMs && + Number.isFinite(this.options.reconcileHeartbeatMs) && + this.options.reconcileHeartbeatMs > 0 + ? this.options.reconcileHeartbeatMs + : BUILD_RECONCILE_HEARTBEAT_MS, + ); + const scheduleHeartbeat = () => { + if (!heartbeatActive || leaseOwnershipLost) return; + heartbeat = (this.options.reconcileSetTimeout ?? setTimeout)(() => { + void requireLeaseOwnership() + .catch(() => { + leaseOwnershipLost = true; + }) + .finally(scheduleHeartbeat); + }, heartbeatMs); + }; + heartbeatActive = true; + scheduleHeartbeat(); + await requireLeaseOwnership(); + return await this.reconcileBuildInternal( + id, + requireLeaseOwnership, + reconcileLease, + ); + } finally { + heartbeatActive = false; + if (heartbeat !== undefined) + (this.options.reconcileClearTimeout ?? clearTimeout)(heartbeat); + try { + if (reconcileLease) + await this.options.store.releaseBuildReconciliationLease( + id, + reconcileLease.token, + ); + } finally { + await lease?.release(); + } + } + } + + private async reconcileBuildInternal( + id: string, + requireLeaseOwnership: () => Promise, + reconcileLease: BuildReconciliationLease, + ): Promise { let record = await this.requireBuild(id); if (record.status.state === "succeeded" || record.status.state === "failed") return recordStatus(record); - if (record.status.jobCreated) record = await this.captureLogs(record); + if (record.status.jobCreated) + record = await this.captureLogs( + record, + requireLeaseOwnership, + reconcileLease, + ); const observation = await this.options.kubernetes.getJob( this.options.namespace, record.spec.jobName, ); if (!observation) return recordStatus(record); if (observation.phase === "running" && record.status.state === "queued") { - record = await this.setState(record, "running", { - startedAt: observation.startedAt ?? this.now().toISOString(), - }); + await requireLeaseOwnership(); + record = await this.setState( + record, + "running", + { + startedAt: observation.startedAt ?? this.now().toISOString(), + }, + reconcileLease, + ); } else if (observation.phase === "failed") { - record = await this.setState(record, "failed", { - startedAt: record.status.startedAt ?? observation.startedAt, - finishedAt: observation.finishedAt ?? this.now().toISOString(), - error: observation.error ?? "BuildKit Job failed", - }); + await requireLeaseOwnership(); + record = await this.setState( + record, + "failed", + { + startedAt: record.status.startedAt ?? observation.startedAt, + finishedAt: observation.finishedAt ?? this.now().toISOString(), + error: observation.error ?? "BuildKit Job failed", + }, + reconcileLease, + ); } else if (observation.phase === "succeeded") { + let digest: Sha256Digest; try { - const digest = await this.digestResolver( - record.spec.request.spec.image, - ); + digest = await this.digestResolver(record.spec.request.spec.image); assertSha256Digest(digest); - record = await this.setState(record, "succeeded", { + } catch (error) { + await requireLeaseOwnership(); + record = await this.setState( + record, + "failed", + { + finishedAt: this.now().toISOString(), + error: `Unable to resolve pushed image digest: ${error instanceof Error ? error.message : String(error)}`, + }, + reconcileLease, + ); + return recordStatus(record); + } + await requireLeaseOwnership(); + record = await this.setState( + record, + "succeeded", + { startedAt: record.status.startedAt ?? observation.startedAt, finishedAt: observation.finishedAt ?? this.now().toISOString(), digest, - }); - } catch (error) { - record = await this.setState(record, "failed", { - finishedAt: this.now().toISOString(), - error: `Unable to resolve pushed image digest: ${error instanceof Error ? error.message : String(error)}`, - }); - } + }, + reconcileLease, + ); } return recordStatus(record); } + async reconcilePendingBuilds(options: ReconcileOptions = {}): Promise { + throwIfAborted(options.signal); + const records = await this.options.store.listBuilds(); + throwIfAborted(options.signal); + for (const record of records) { + throwIfAborted(options.signal); + if ( + !record.status.jobCreated || + record.status.state === "succeeded" || + record.status.state === "failed" + ) + continue; + try { + await this.reconcileBuild(record.metadata.name, options); + } catch (error) { + if (options.signal?.aborted) throw error; + try { + this.options.onReconcileFailure?.(record.metadata.name, error); + } catch { + // Failure reporting must not interrupt reconciliation of other builds. + } + } + } + } + async cancelBuild(id: string): Promise { const record = await this.requireBuild(id); if (record.status.state === "succeeded" || record.status.state === "failed") @@ -637,15 +825,26 @@ export class BuildController { private async updateBuild( record: BuildRecord, change: (next: BuildRecord) => void, + reconcileLease?: BuildReconciliationLease, ): Promise { + // A heartbeat may have advanced this record's resourceVersion while an + // unabortable Kubernetes observation was pending. + if (reconcileLease) record = await this.requireBuild(record.metadata.name); const next = clone(record); next.metadata.resourceVersion = String( Number(record.metadata.resourceVersion) + 1, ); change(next); + if (reconcileLease) { + if (next.status.reconcileLease?.token !== reconcileLease.token) + throw new BuildConflictError( + "Build reconciliation lease ownership was lost", + ); + } await this.options.store.replaceBuild( next, record.metadata.resourceVersion, + reconcileLease?.token, ); return this.requireBuild(record.metadata.name); } @@ -654,12 +853,17 @@ export class BuildController { record: BuildRecord, state: BuildStatus["state"], values: Partial, + reconcileLease?: BuildReconciliationLease, ): Promise { if (record.status.state === state && state === "running") return record; - return this.updateBuild(record, (next) => { - Object.assign(next.status, values, { state }); - next.status.events.push({ type: "status", status: recordStatus(next) }); - }); + return this.updateBuild( + record, + (next) => { + Object.assign(next.status, values, { state }); + next.status.events.push({ type: "status", status: recordStatus(next) }); + }, + reconcileLease, + ); } private async failBuild(id: string, message: string): Promise { @@ -675,7 +879,11 @@ export class BuildController { }); } - private async captureLogs(record: BuildRecord): Promise { + private async captureLogs( + record: BuildRecord, + requireLeaseOwnership: () => Promise, + reconcileLease: BuildReconciliationLease, + ): Promise { let raw: string | Uint8Array; try { raw = await this.options.kubernetes.getJobLogs( @@ -689,31 +897,36 @@ export class BuildController { const offset = bytes.byteLength < record.status.logOffset ? 0 : record.status.logOffset; if (bytes.byteLength === offset) return record; + await requireLeaseOwnership(); const delta = bytes.subarray(offset); - return this.updateBuild(record, (next) => { - next.status.logOffset = bytes.byteLength; - for (const message of delta - .toString("utf8") - .split(/(?<=\n)/) - .filter(Boolean)) { - const event: BuildEvent = { - type: "log", - id: record.metadata.name, - sequence: next.status.nextSequence++, - message, - }; - next.status.events.push(event); - next.status.logBytes += Buffer.byteLength(message); - } - while (next.status.logBytes > this.maxLogBytes) { - const index = next.status.events.findIndex( - (event) => event.type === "log", - ); - if (index === -1) break; - const [removed] = next.status.events.splice(index, 1); - if (removed?.type === "log") - next.status.logBytes -= Buffer.byteLength(removed.message); - } - }); + return this.updateBuild( + record, + (next) => { + next.status.logOffset = bytes.byteLength; + for (const message of delta + .toString("utf8") + .split(/(?<=\n)/) + .filter(Boolean)) { + const event: BuildEvent = { + type: "log", + id: record.metadata.name, + sequence: next.status.nextSequence++, + message, + }; + next.status.events.push(event); + next.status.logBytes += Buffer.byteLength(message); + } + while (next.status.logBytes > this.maxLogBytes) { + const index = next.status.events.findIndex( + (event) => event.type === "log", + ); + if (index === -1) break; + const [removed] = next.status.events.splice(index, 1); + if (removed?.type === "log") + next.status.logBytes -= Buffer.byteLength(removed.message); + } + }, + reconcileLease, + ); } } diff --git a/server/build-job.ts b/server/build-job.ts index 4ca0ce2..fa29200 100644 --- a/server/build-job.ts +++ b/server/build-job.ts @@ -100,6 +100,25 @@ export function createBuildJob(options: BuildJobOptions): KubernetesJob { "kuber.astrxl.dev/build": options.name, ...options.labels, }; + const amd64Toleration = { + key: "arch", + operator: "Equal", + value: "amd64", + effect: "NoExecute", + }; + const tolerations = [...(options.tolerations ?? [])]; + if ( + options.spec.architecture === "amd64" && + !tolerations.some( + (toleration) => + toleration.key === amd64Toleration.key && + toleration.operator === amd64Toleration.operator && + toleration.value === amd64Toleration.value && + toleration.effect === amd64Toleration.effect && + Object.keys(toleration).length === Object.keys(amd64Toleration).length, + ) + ) + tolerations.push(amd64Toleration); return { apiVersion: "batch/v1", @@ -134,7 +153,7 @@ export function createBuildJob(options: BuildJobOptions): KubernetesJob { ...options.nodeSelector, "kubernetes.io/arch": options.spec.architecture, }, - ...(options.tolerations ? { tolerations: options.tolerations } : {}), + ...(tolerations.length ? { tolerations } : {}), securityContext: { runAsNonRoot: true, runAsUser: 1000, diff --git a/server/build-kubernetes.ts b/server/build-kubernetes.ts index 84ec4e0..57f29ea 100644 --- a/server/build-kubernetes.ts +++ b/server/build-kubernetes.ts @@ -5,7 +5,7 @@ import { type V1DeleteOptions, type V1Job, } from "@kubernetes/client-node"; -import { createHash } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import { mkdir, open, readFile, rm, writeFile } from "node:fs/promises"; import { join } from "node:path"; import type { Sha256Digest } from "../shared/build-protocol"; @@ -19,6 +19,7 @@ import { buildTimestamp, compareBuildRecords, type BuildRecord, + type BuildReconciliationLease, type BuildStore, type CreateBuildResult, type UploadRecord, @@ -541,6 +542,7 @@ export class KubernetesBuildStore implements BuildStore { async replaceBuild( record: BuildRecord, expectedResourceVersion: string, + reconcileLeaseToken?: string, ): Promise { const current = await this.getBuild(record.metadata.name); if ( @@ -552,6 +554,14 @@ export class KubernetesBuildStore implements BuildStore { ); if (!sameSpec(current, record)) throw new BuildStoreConflictError("Build specification is immutable"); + if ( + reconcileLeaseToken && + current.status.reconcileLease?.token !== reconcileLeaseToken + ) { + throw new BuildStoreConflictError( + "Build reconciliation lease ownership was lost", + ); + } if ( current.status.state === "succeeded" && (record.status.state !== "succeeded" || @@ -584,6 +594,97 @@ export class KubernetesBuildStore implements BuildStore { } } + async acquireBuildReconciliationLease( + id: string, + holder: string, + ttlMs: number, + ): Promise { + if (!holder || !Number.isFinite(ttlMs) || ttlMs <= 0) + throw new BuildStoreConflictError("Invalid build reconciliation lease"); + for (let attempt = 0; attempt < 8; attempt++) { + const current = await this.getBuild(id); + if (!current) return; + if ( + current.status.reconcileLease && + Date.parse(current.status.reconcileLease.expiresAt) > Date.now() + ) { + return; + } + const lease: BuildReconciliationLease = { + holder, + token: randomUUID(), + expiresAt: new Date(Date.now() + ttlMs).toISOString(), + }; + const next = structuredClone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + next.status.reconcileLease = lease; + try { + // The old fence token and resourceVersion identify one persisted record. + // This CAS installs the successor fence before it can reconcile. + await this.replaceBuild(next, current.metadata.resourceVersion); + return lease; + } catch (error) { + if (!(error instanceof BuildStoreConflictError) || attempt === 7) + throw error; + } + } + throw new BuildStoreConflictError( + "Build reconciliation lease acquisition timed out", + ); + } + + async renewBuildReconciliationLease( + id: string, + token: string, + ttlMs: number, + ): Promise { + if (!token || !Number.isFinite(ttlMs) || ttlMs <= 0) return false; + for (let attempt = 0; attempt < 8; attempt++) { + const current = await this.getBuild(id); + if ( + !current || + current.status.reconcileLease?.token !== token || + Date.parse(current.status.reconcileLease.expiresAt) <= Date.now() + ) { + return false; + } + const next = structuredClone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + next.status.reconcileLease!.expiresAt = new Date( + Date.now() + ttlMs, + ).toISOString(); + try { + await this.replaceBuild(next, current.metadata.resourceVersion, token); + return true; + } catch (error) { + if (!(error instanceof BuildStoreConflictError)) throw error; + } + } + return false; + } + + async releaseBuildReconciliationLease( + id: string, + token: string, + ): Promise { + const current = await this.getBuild(id); + if (current?.status.reconcileLease?.token !== token) return; + const next = structuredClone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + delete next.status.reconcileLease; + try { + await this.replaceBuild(next, current.metadata.resourceVersion, token); + } catch (error) { + if (!(error instanceof BuildStoreConflictError)) throw error; + } + } + async ownsBuild(imageKey: string, buildId: string): Promise { const lock = await this.read<{ buildId: string }>( map(this.namespace, this.lockName(imageKey), "build-lock"), diff --git a/server/build-reconciler.ts b/server/build-reconciler.ts new file mode 100644 index 0000000..e70a069 --- /dev/null +++ b/server/build-reconciler.ts @@ -0,0 +1,84 @@ +export const DEFAULT_BUILD_RECONCILE_MS = 30 * 1000; +export const DEFAULT_BUILD_RECONCILE_TIMEOUT_MS = 20 * 1000; +/** Leaves time for the 30-second reconciliation lease to be renewed or released. */ +export const MAX_BUILD_RECONCILE_TIMEOUT_MS = 25 * 1000; +/** Node clamps timer delays above this value to approximately one millisecond. */ +export const MAX_BUILD_RECONCILE_MS = 2_147_483_647; + +export function buildReconcileIntervalMs(value: string | undefined): number { + const parsed = Number(value ?? DEFAULT_BUILD_RECONCILE_MS); + return Number.isFinite(parsed) && + parsed > 0 && + parsed <= MAX_BUILD_RECONCILE_MS + ? parsed + : DEFAULT_BUILD_RECONCILE_MS; +} + +export function buildReconcileTimeoutMs(value: string | undefined): number { + const parsed = Number(value ?? DEFAULT_BUILD_RECONCILE_TIMEOUT_MS); + return Number.isFinite(parsed) && + parsed > 0 && + parsed <= MAX_BUILD_RECONCILE_TIMEOUT_MS + ? parsed + : DEFAULT_BUILD_RECONCILE_TIMEOUT_MS; +} + +export interface BuildReconcileRunner { + reconcilePendingBuilds(options?: { signal?: AbortSignal }): Promise; +} + +export interface BackgroundBuildReconcilerOptions { + runner: BuildReconcileRunner; + timeoutMs: number; + onFailure(error: unknown): void; + setTimeout?: typeof setTimeout; + clearTimeout?: typeof clearTimeout; +} + +/** + * Runs one scan at a time. A deadline cancels its generation and reports the + * timeout, but the generation remains active until its promise settles. Some + * Kubernetes adapters cannot abort an in-flight request, so clearing it at the + * deadline could otherwise overlap work for the same build. + */ +export class BackgroundBuildReconciler { + private current: + | { + controller: AbortController; + timeout: ReturnType; + } + | undefined; + private readonly setTimeout: typeof setTimeout; + private readonly clearTimeout: typeof clearTimeout; + + constructor(private readonly options: BackgroundBuildReconcilerOptions) { + this.setTimeout = options.setTimeout ?? setTimeout; + this.clearTimeout = options.clearTimeout ?? clearTimeout; + } + + tick(): void { + if (this.current) return; + const controller = new AbortController(); + const timeout = this.setTimeout(() => { + if (this.current?.controller !== controller) return; + controller.abort("Build reconciliation scan timed out"); + this.options.onFailure( + new Error( + `Build reconciliation scan timed out after ${this.options.timeoutMs}ms`, + ), + ); + }, this.options.timeoutMs); + this.current = { controller, timeout }; + + void this.options.runner + .reconcilePendingBuilds({ signal: controller.signal }) + .catch((error) => { + if (!controller.signal.aborted) this.options.onFailure(error); + }) + .finally(() => { + if (this.current?.controller !== controller) return; + this.clearTimeout(timeout); + this.current = undefined; + }); + } +} diff --git a/server/build-store.ts b/server/build-store.ts index 5d5dd18..46c48c2 100644 --- a/server/build-store.ts +++ b/server/build-store.ts @@ -1,3 +1,4 @@ +import { randomUUID } from "node:crypto"; import type { BuildEvent, BuildRequest, @@ -29,9 +30,17 @@ export interface BuildRecord { logOffset: number; nextSequence: number; events: BuildEvent[]; + reconcileLease?: BuildReconciliationLease; }; } +/** A record-local fencing token for distributed build reconciliation. */ +export interface BuildReconciliationLease { + holder: string; + token: string; + expiresAt: string; +} + export interface UploadRecord { apiVersion: typeof BUILD_RECORD_API_VERSION; kind: "BuildUpload"; @@ -58,7 +67,19 @@ export interface BuildStore { replaceBuild( record: BuildRecord, expectedResourceVersion: string, + reconcileLeaseToken?: string, ): Promise; + acquireBuildReconciliationLease( + id: string, + holder: string, + ttlMs: number, + ): Promise; + renewBuildReconciliationLease( + id: string, + token: string, + ttlMs: number, + ): Promise; + releaseBuildReconciliationLease(id: string, token: string): Promise; ownsBuild?(imageKey: string, buildId: string): Promise; getUpload(digest: Sha256Digest): Promise; createUpload(record: UploadRecord): Promise; @@ -237,7 +258,11 @@ export class MemoryBuildStore implements BuildStore { return [...this.builds.values()].sort(compareBuildRecords).map(clone); } - async replaceBuild(record: BuildRecord, expectedResourceVersion: string) { + async replaceBuild( + record: BuildRecord, + expectedResourceVersion: string, + reconcileLeaseToken?: string, + ) { const current = this.builds.get(record.metadata.name); if ( !current || @@ -250,6 +275,14 @@ export class MemoryBuildStore implements BuildStore { if (!sameSpec(current, record)) { throw new BuildStoreConflictError("Build specification is immutable"); } + if ( + reconcileLeaseToken && + current.status.reconcileLease?.token !== reconcileLeaseToken + ) { + throw new BuildStoreConflictError( + "Build reconciliation lease ownership was lost", + ); + } if ( current.status.state === "succeeded" && (record.status.state !== "succeeded" || @@ -268,6 +301,79 @@ export class MemoryBuildStore implements BuildStore { } } + async acquireBuildReconciliationLease( + id: string, + holder: string, + ttlMs: number, + ): Promise { + if (!holder || !Number.isFinite(ttlMs) || ttlMs <= 0) + throw new BuildStoreConflictError("Invalid build reconciliation lease"); + const current = this.builds.get(id); + if (!current) return; + if ( + current.status.reconcileLease && + Date.parse(current.status.reconcileLease.expiresAt) > Date.now() + ) { + return; + } + const lease: BuildReconciliationLease = { + holder, + token: randomUUID(), + expiresAt: new Date(Date.now() + ttlMs).toISOString(), + }; + const next = clone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + next.status.reconcileLease = lease; + await this.replaceBuild(next, current.metadata.resourceVersion); + return clone(lease); + } + + async renewBuildReconciliationLease( + id: string, + token: string, + ttlMs: number, + ): Promise { + if (!token || !Number.isFinite(ttlMs) || ttlMs <= 0) return false; + const current = this.builds.get(id); + if ( + !current || + current.status.reconcileLease?.token !== token || + Date.parse(current.status.reconcileLease.expiresAt) <= Date.now() + ) { + return false; + } + const next = clone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + next.status.reconcileLease!.expiresAt = new Date( + Date.now() + ttlMs, + ).toISOString(); + try { + await this.replaceBuild(next, current.metadata.resourceVersion, token); + return true; + } catch (error) { + if (error instanceof BuildStoreConflictError) return false; + throw error; + } + } + + async releaseBuildReconciliationLease( + id: string, + token: string, + ): Promise { + const current = this.builds.get(id); + if (current?.status.reconcileLease?.token !== token) return; + const next = clone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + delete next.status.reconcileLease; + await this.replaceBuild(next, current.metadata.resourceVersion, token); + } + async ownsBuild(imageKey: string, buildId: string): Promise { return this.active.get(imageKey) === buildId; } diff --git a/server/index.ts b/server/index.ts index 79a3e69..2b0b998 100644 --- a/server/index.ts +++ b/server/index.ts @@ -1,4 +1,5 @@ import { createApp, cleanupExpiredSessions } from "./app"; +import { logServerRequest, processLogger, safeLog } from "../lib/request-log"; import { RedactingAuditStore } from "./audit-store"; import { readFile } from "node:fs/promises"; import { IMAGE_REGISTRY } from "../const"; @@ -37,6 +38,11 @@ import { import { KubernetesLogs } from "./kubernetes-logs"; import { createLogService } from "./log-service"; import { resolveRegistryDigest, type RegistryCredentials } from "./registry"; +import { + BackgroundBuildReconciler, + buildReconcileIntervalMs, + buildReconcileTimeoutMs, +} from "./build-reconciler"; const clients = createKubernetesClients(); const store = new KubernetesAuthStore(clients.objects); @@ -102,6 +108,10 @@ const buildStore = new KubernetesBuildStore( namespace, `${dataRoot}/uploads`, ); +const leases = new KubernetesWorkspaceLeaseProvider( + createKubernetesLeaseObjects(clients.coordination), + namespace, +); const builds = new BuildController({ cas: new FilesystemCas(`${dataRoot}/cas`), store: buildStore, @@ -130,6 +140,13 @@ const builds = new BuildController({ origin: registryResolveOrigin, insecure: registryResolveOrigin?.startsWith("http://"), }), + onReconcileFailure: (buildId, error) => + safeLog(processLogger, { + event: "build.reconcile.failed", + buildId, + error, + }), + reconcileLeases: leases, }); const logs = createLogService( new KubernetesLogs(clients.config, clients.apps, clients.core), @@ -152,10 +169,6 @@ if (bootstrapUsername && bootstrapPassword) { } } -const leases = new KubernetesWorkspaceLeaseProvider( - createKubernetesLeaseObjects(clients.coordination), -); - try { const recovered = await recoverStaleOperations(operationStore); if (recovered > 0) @@ -181,6 +194,20 @@ const sessionCleanupTimer = setInterval( ); sessionCleanupTimer.unref?.(); +const buildReconciler = new BackgroundBuildReconciler({ + runner: builds, + timeoutMs: buildReconcileTimeoutMs( + process.env.KUBER_BUILD_RECONCILE_TIMEOUT_MS, + ), + onFailure: (error) => + safeLog(processLogger, { event: "build.reconcile.scan.failed", error }), +}); +const buildReconcileTimer = setInterval( + () => buildReconciler.tick(), + buildReconcileIntervalMs(process.env.KUBER_BUILD_RECONCILE_MS), +); +buildReconcileTimer.unref?.(); + const app = createApp({ store, workspaceStore, @@ -223,16 +250,20 @@ const server = Bun.serve({ const url = new URL(request.url); const workspaceId = execUpgradeMatch(url); if (workspaceId && request.method === "GET") { - return handleExecUpgrade( - { - store, - workspaceStore, - auditStore, - }, - request, - workspaceId, - (upgradeRequest, connection) => - server.upgrade(upgradeRequest, { data: { connection } }), + const requestId = + request.headers.get("x-request-id") ?? crypto.randomUUID(); + return logServerRequest(request, requestId, () => + handleExecUpgrade( + { + store, + workspaceStore, + auditStore, + }, + request, + workspaceId, + (upgradeRequest, connection) => + server.upgrade(upgradeRequest, { data: { connection } }), + ), ); } return app(request); diff --git a/server/kubernetes-exec.ts b/server/kubernetes-exec.ts index 1d6268c..4f9a984 100644 --- a/server/kubernetes-exec.ts +++ b/server/kubernetes-exec.ts @@ -7,6 +7,7 @@ import { type V1Status, } from "@kubernetes/client-node"; import { PassThrough } from "node:stream"; +import { logKubernetesRequest } from "../lib/request-log"; import { type ExecDeployment, type ExecPod, @@ -17,18 +18,14 @@ import { type KubernetesExecRequest, } from "./exec-service"; -function selectorString( - labels: Readonly>, -): string { +function selectorString(labels: Readonly>): string { return Object.entries(labels) .sort(([left], [right]) => left.localeCompare(right)) .map(([key, value]) => `${key}=${value}`) .join(","); } -function matchLabels( - deployment: V1Deployment, -): Record { +function matchLabels(deployment: V1Deployment): Record { const matchExpressions = deployment.spec?.selector?.matchExpressions; if (matchExpressions?.length) return {}; return deployment.spec?.selector?.matchLabels ?? {}; @@ -37,11 +34,13 @@ function matchLabels( async function ownerDeploymentUid( apps: AppsV1Api, namespace: string, - ownerReferences: Array<{ - kind?: string; - name?: string; - uid?: string; - }> | undefined, + ownerReferences: + | Array<{ + kind?: string; + name?: string; + uid?: string; + }> + | undefined, ): Promise { if (!ownerReferences?.length) return ""; for (const ref of ownerReferences) { @@ -153,27 +152,25 @@ export class KubernetesExec implements KubernetesExecBackend { pod.metadata?.ownerReferences, ); const containerStatuses = pod.status?.containerStatuses ?? []; - const containers: ExecPodContainer[] = ( - pod.spec?.containers ?? [] - ).map((container) => { - const status = containerStatuses.find( - (item) => item.name === container.name, - ); - return { - name: container.name ?? "", - running: status?.state?.running !== undefined, - ready: status?.ready, - }; - }); + const containers: ExecPodContainer[] = (pod.spec?.containers ?? []).map( + (container) => { + const status = containerStatuses.find( + (item) => item.name === container.name, + ); + return { + name: container.name ?? "", + running: status?.state?.running !== undefined, + ready: status?.ready, + }; + }, + ); pods.push({ name, uid, deploymentUid, phase: pod.status?.phase, deletionTimestamp: pod.metadata?.deletionTimestamp - ? new Date( - pod.metadata.deletionTimestamp, - ).toISOString() + ? new Date(pod.metadata.deletionTimestamp).toISOString() : undefined, containers, }); @@ -206,22 +203,35 @@ export class KubernetesExec implements KubernetesExecBackend { signal.addEventListener("abort", abort, { once: true }); try { - const ws = await this.execClient.exec( - request.namespace, - request.pod, - request.container, - [...request.command], - stdout, - stderr, - stdin, - request.tty, - (status) => { - completed({ - exitCode: exitCodeFromStatus(status) ?? 1, - reason: status.reason, - message: status.message, - }); + const ws = await logKubernetesRequest( + { + method: "GET", + pathname: `/api/v1/namespaces/${request.namespace}/pods/${request.pod}/exec`, + body: { + container: request.container, + command: request.command, + tty: request.tty, + }, + attempt: 1, }, + () => + this.execClient.exec( + request.namespace, + request.pod, + request.container, + [...request.command], + stdout, + stderr, + stdin, + request.tty, + (status) => { + completed({ + exitCode: exitCodeFromStatus(status) ?? 1, + reason: status.reason, + message: status.message, + }); + }, + ), ); close = () => { try { diff --git a/server/kubernetes-logs.ts b/server/kubernetes-logs.ts index fbed4fe..9a469c4 100644 --- a/server/kubernetes-logs.ts +++ b/server/kubernetes-logs.ts @@ -6,6 +6,7 @@ import { type V1LabelSelector, } from "@kubernetes/client-node"; import { PassThrough } from "node:stream"; +import { logKubernetesRequest } from "../lib/request-log"; import { KubernetesLogError, type ContainerLogRequest, @@ -19,7 +20,9 @@ function selector(labels: Readonly>): string { .join(","); } -function matchLabels(value: V1LabelSelector | undefined): Record { +function matchLabels( + value: V1LabelSelector | undefined, +): Record { if (value?.matchExpressions?.length) return {}; return value?.matchLabels ?? {}; } @@ -62,7 +65,9 @@ export class KubernetesLogs implements KubernetesLogsBackend { return deployments.items.flatMap((deployment) => { const name = deployment.metadata?.name; const labels = matchLabels(deployment.spec?.selector); - return name && Object.keys(labels).length ? [{ name, selector: labels }] : []; + return name && Object.keys(labels).length + ? [{ name, selector: labels }] + : []; }); } catch (error) { throw logError(error); @@ -71,7 +76,10 @@ export class KubernetesLogs implements KubernetesLogsBackend { async getService(namespace: string, name: string) { try { - const service = await this.core.readNamespacedService({ namespace, name }); + const service = await this.core.readNamespacedService({ + namespace, + name, + }); return { name, selector: service.spec?.selector ?? {} }; } catch (error) { if (statusCode(error) === 404) return; @@ -79,10 +87,7 @@ export class KubernetesLogs implements KubernetesLogsBackend { } } - async listPods( - namespace: string, - labels: Readonly>, - ) { + async listPods(namespace: string, labels: Readonly>) { try { const pods = await this.core.listNamespacedPod({ namespace, @@ -92,7 +97,9 @@ export class KubernetesLogs implements KubernetesLogsBackend { name: pod.metadata?.name ?? "", uid: pod.metadata?.uid ?? "", phase: pod.status?.phase, - containers: (pod.spec?.containers ?? []).map((container) => container.name), + containers: (pod.spec?.containers ?? []).map( + (container) => container.name, + ), })); } catch (error) { throw logError(error); @@ -126,18 +133,27 @@ export class KubernetesLogs implements KubernetesLogsBackend { }; signal.addEventListener("abort", abort, { once: true }); try { - upstream = await this.logger.log( - request.namespace, - request.pod, - request.container, - output, + upstream = await logKubernetesRequest( { - follow: true, - tailLines: request.tailLines, - sinceSeconds: request.sinceSeconds, - sinceTime: request.sinceTime, - timestamps: request.timestamps, + method: "GET", + pathname: `/api/v1/namespaces/${request.namespace}/pods/${request.pod}/log`, + body: undefined, + attempt: 1, }, + () => + this.logger.log( + request.namespace, + request.pod, + request.container, + output, + { + follow: true, + tailLines: request.tailLines, + sinceSeconds: request.sinceSeconds, + sinceTime: request.sinceTime, + timestamps: request.timestamps, + }, + ), ); if (signal.aborted) abort(); for await (const chunk of output) yield new Uint8Array(chunk); diff --git a/server/kubernetes-state.ts b/server/kubernetes-state.ts index 1bd53e8..73b0f09 100644 --- a/server/kubernetes-state.ts +++ b/server/kubernetes-state.ts @@ -14,6 +14,7 @@ import { } from "@kubernetes/client-node"; import { createHash } from "node:crypto"; import { createKubernetesHttpLibrary } from "../lib/k8s-http"; +import { sanitizeKubernetesObject } from "../lib/kubernetes-object"; import { type AuditEvent, type AuditPersistence } from "./audit-store"; import { managementDependencies, @@ -1128,7 +1129,7 @@ export function createKubernetesManagementDependencies( } } return objects.patch( - resource, + sanitizeKubernetesObject(resource), undefined, undefined, FIELD_MANAGER, diff --git a/tests/lib/apply.test.ts b/tests/lib/apply.test.ts index 6fac232..c29437c 100644 --- a/tests/lib/apply.test.ts +++ b/tests/lib/apply.test.ts @@ -1,6 +1,9 @@ -import { describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"; import type { KubernetesObject } from "@kubernetes/client-node"; -import { getResourceKey, sortResources } from "../../lib/apply"; +import { applyResource, getResourceKey, sortResources } from "../../lib/apply"; +import { objectApi } from "../../lib/k8s"; + +afterEach(() => mock.restore()); function resource(kind: string, name = kind): KubernetesObject { return { apiVersion: "v1", kind, metadata: { name, namespace: "project" } }; @@ -72,3 +75,59 @@ describe("resource ordering", () => { ).toBe("Namespace::app"); }); }); + +describe("resource apply", () => { + test("does not send server-managed fields from a resource read from Kubernetes", async () => { + const returned = { + apiVersion: "postgresql.cnpg.io/v1", + kind: "Database", + metadata: { + name: "orders", + namespace: "database", + labels: { "kuber.dev/project": "demo" }, + annotations: { "example.dev/setting": "enabled" }, + managedFields: [{ manager: "kube-controller-manager" }], + resourceVersion: "42", + uid: "database-uid", + creationTimestamp: new Date("2026-09-08T00:00:00Z"), + generation: 3, + selfLink: + "/apis/postgresql.cnpg.io/v1/namespaces/database/databases/orders", + }, + spec: { name: "orders" }, + status: { phase: "Ready" }, + } as KubernetesObject; + const patch = spyOn(objectApi, "patch").mockImplementation( + async (resource) => resource as never, + ); + + await applyResource(returned); + + expect(patch).toHaveBeenCalledWith( + expect.objectContaining({ + metadata: { + name: "orders", + namespace: "database", + labels: { "kuber.dev/project": "demo" }, + annotations: { "example.dev/setting": "enabled" }, + }, + spec: { name: "orders" }, + }), + undefined, + undefined, + "kuber", + true, + "application/apply-patch+yaml", + ); + const sent = patch.mock.calls[0]?.[0] as typeof returned; + expect(sent.metadata).not.toHaveProperty("managedFields"); + expect(sent.metadata).not.toHaveProperty("resourceVersion"); + expect(sent.metadata).not.toHaveProperty("uid"); + expect(sent.metadata).not.toHaveProperty("creationTimestamp"); + expect(sent.metadata).not.toHaveProperty("generation"); + expect(sent.metadata).not.toHaveProperty("selfLink"); + expect(sent).not.toHaveProperty("status"); + expect(returned.metadata).toHaveProperty("managedFields"); + expect(returned).toHaveProperty("status"); + }); +}); diff --git a/tests/lib/convert-deployment.test.ts b/tests/lib/convert-deployment.test.ts index 6681090..dde0597 100644 --- a/tests/lib/convert-deployment.test.ts +++ b/tests/lib/convert-deployment.test.ts @@ -170,6 +170,133 @@ describe("Deployment conversion", () => { ).spec?.template.spec?.containers[0]?.envFrom, ).toEqual([{ secretRef: { name: "app-env" } }]); }); + + test("maps deploy resource limits and reservations to container resources", () => { + const service = { + image: "app", + deploy: { + resources: { + limits: { cpus: "0.25", memory: "3072M" }, + reservations: { cpus: 0.5, memory: "1g" }, + }, + }, + } as Service; + const deployment = serviceToDeployment( + "project", + "app", + compose(service), + service, + ); + + expect(deployment.spec?.template.spec?.containers[0]?.resources).toEqual({ + limits: { cpu: "250m", memory: "3072M" }, + requests: { cpu: "500m", memory: "1G" }, + }); + }); + + test("normalizes Compose memory units to decimal Kubernetes quantities", () => { + const service = { + image: "app", + deploy: { + resources: { + limits: { cpus: 2, memory: "256M" }, + reservations: { cpus: "2.0", memory: "100m" }, + }, + }, + } as Service; + const deployment = serviceToDeployment( + "project", + "app", + compose(service), + service, + ); + + expect(deployment.spec?.template.spec?.containers[0]?.resources).toEqual({ + limits: { cpu: "2", memory: "256M" }, + requests: { cpu: "2", memory: "100M" }, + }); + }); + + test.each(["1Ki", "1Mi", "1Gi", "1Ti", "1Pi", "1Ei"])( + "preserves Kubernetes binary memory quantity %s", + (memory) => { + const service = { + image: "app", + deploy: { + resources: { + limits: { memory }, + }, + }, + } as Service; + const deployment = serviceToDeployment( + "project", + "app", + compose(service), + service, + ); + + expect(deployment.spec?.template.spec?.containers[0]?.resources).toEqual({ + limits: { memory }, + }); + }, + ); + + test.each(["", "-1M", "1e3", "1MBi", "1foo"])( + "rejects invalid memory resource quantity %s", + (memory) => { + const service = { + image: "app", + deploy: { resources: { limits: { memory } } }, + } as Service; + expect(() => + serviceToDeployment("project", "app", compose(service), service), + ).toThrow("deploy.resources.limits.memory"); + }, + ); + + test("rejects invalid CPU resource quantities", () => { + for (const cpus of [ + "NaN", + "Infinity", + Number.POSITIVE_INFINITY, + "-0.5", + "0.0005", + ]) { + const service = { + image: "app", + deploy: { resources: { limits: { cpus } } }, + } as Service; + expect(() => + serviceToDeployment("project", "app", compose(service), service), + ).toThrow("deploy.resources.limits.cpus"); + } + }); + + test("lets x-container resource fields override deploy resources", () => { + const service = { + image: "app", + deploy: { + resources: { + limits: { cpus: "0.5", memory: "256M" }, + reservations: { cpus: "0.25", memory: "128M" }, + }, + }, + "x-container": { + resources: { limits: { memory: "768Mi" }, requests: { cpu: "900m" } }, + }, + } as Service; + const deployment = serviceToDeployment( + "project", + "app", + compose(service), + service, + ); + + expect(deployment.spec?.template.spec?.containers[0]?.resources).toEqual({ + limits: { cpu: "500m", memory: "768Mi" }, + requests: { cpu: "900m", memory: "128M" }, + }); + }); }); describe("port and routing conversion", () => { @@ -439,6 +566,21 @@ describe("autoscaled (range) conversion", () => { ).toEqual({ cpu: "500m" }); }); + test("uses a deploy CPU reservation instead of the HPA fallback request", async () => { + const service = { + image: "app", + deploy: { + replicas: "1-4", + resources: { reservations: { cpus: "0.25" } }, + }, + } as Service; + const resources = await composeToKubernetes("project", compose(service)); + const deployment = findDeployment(resources); + expect( + deployment.spec?.template.spec?.containers?.[0]?.resources?.requests, + ).toEqual({ cpu: "250m" }); + }); + test("does not inject a CPU request on fixed-replica services", async () => { const service = { image: "app" } as Service; const resources = await composeToKubernetes("project", compose(service)); diff --git a/tests/lib/database.test.ts b/tests/lib/database.test.ts index 7696456..a04685b 100644 --- a/tests/lib/database.test.ts +++ b/tests/lib/database.test.ts @@ -151,4 +151,59 @@ describe("managed PostgreSQL claims", () => { "Database", ]); }); + + test("reconciles a CNPG cluster returned with managed fields without sending them back", async () => { + const claim = { + service: "app", + username: "app", + database: "app", + secretName: "postgres-app", + }; + spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Secret") { + return { + ...resource, + data: { + username: Buffer.from("app").toString("base64"), + password: Buffer.from("secret").toString("base64"), + }, + } as never; + } + return { + ...resource, + metadata: { + ...resource.metadata, + managedFields: [{ manager: "cloudnative-pg" }], + resourceVersion: "42", + uid: "cluster-uid", + creationTimestamp: "2026-09-08T00:00:00Z", + generation: 3, + }, + spec: { managed: { roles: [{ name: "existing" }] } }, + status: { phase: "Ready" }, + } as never; + }); + const patch = spyOn(objectApi, "patch").mockImplementation( + async (resource) => { + if (resource.metadata?.managedFields) + throw new Error("metadata.managedFields must be nil"); + return resource as never; + }, + ); + + await reconcilePostgresClaim("project", claim); + + const cluster = patch.mock.calls.find( + ([resource]) => resource.kind === "Cluster", + )?.[0]; + expect(cluster).toMatchObject({ + metadata: { name: "postgres", namespace: "database" }, + spec: { + managed: { roles: [{ name: "existing" }, { name: "app" }] }, + }, + }); + expect(cluster?.metadata).not.toHaveProperty("managedFields"); + expect(cluster?.metadata).not.toHaveProperty("resourceVersion"); + expect(cluster).not.toHaveProperty("status"); + }); }); diff --git a/tests/lib/k8s-http.test.ts b/tests/lib/k8s-http.test.ts index 8b1b538..5ee153e 100644 --- a/tests/lib/k8s-http.test.ts +++ b/tests/lib/k8s-http.test.ts @@ -74,6 +74,51 @@ describe("Kubernetes HTTP transport", () => { expect(requests).toBe(2); }); + test("logs Kubernetes request starts, completions, and failures", async () => { + const logs: Record[] = []; + const logger = { + log: (entry: Record) => logs.push(entry), + }; + const { url } = await serve((_request, response) => response.end("ok")); + const client = createKubernetesHttpLibrary({ minIntervalMs: 0, logger }); + + await client.send(new RequestContext(url, HttpMethod.GET)).toPromise(); + expect(logs).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + event: "kuber.k8s.request.start", + method: "GET", + attempt: 1, + }), + expect.objectContaining({ + event: "kuber.k8s.request.end", + status: 200, + responseBody: "ok", + }), + ]), + ); + + const failedLogs: Record[] = []; + const { url: failedUrl } = await serve((request) => + request.socket.destroy(), + ); + const failedClient = createKubernetesHttpLibrary({ + minIntervalMs: 0, + logger: { log: (entry) => failedLogs.push(entry) }, + }); + await expect( + failedClient + .send(new RequestContext(failedUrl, HttpMethod.GET)) + .toPromise(), + ).rejects.toBeDefined(); + expect(failedLogs).toEqual( + expect.arrayContaining([ + expect.objectContaining({ event: "kuber.k8s.request.start" }), + expect.objectContaining({ event: "kuber.k8s.request.failed" }), + ]), + ); + }); + test("paces concurrent request starts", async () => { const starts: number[] = []; const { url } = await serve((_request, response) => { diff --git a/tests/lib/request-log.test.ts b/tests/lib/request-log.test.ts new file mode 100644 index 0000000..609f645 --- /dev/null +++ b/tests/lib/request-log.test.ts @@ -0,0 +1,65 @@ +import { describe, expect, test } from "bun:test"; +import { logServerRequest } from "../../lib/request-log"; + +describe("server response logging", () => { + test("does not consume streaming response bodies", async () => { + let pulls = 0; + const body = new ReadableStream({ + pull(controller) { + pulls += 1; + controller.enqueue(new TextEncoder().encode('{"message":"live"}\n')); + }, + }); + const logs: Record[] = []; + const response = new Response(body, { + headers: { "content-type": "application/x-ndjson; charset=utf-8" }, + }); + await Bun.sleep(0); + const pullsBeforeLogging = pulls; + const loggedResponse = await logServerRequest( + new Request("https://kuber.astrxl.dev/logs"), + "streaming-response", + () => response, + { log: (entry) => logs.push(entry) }, + ); + + expect(pulls).toBe(pullsBeforeLogging); + expect(logs).toContainEqual( + expect.objectContaining({ + event: "kuber.server.response.body", + bodySkipped: true, + streaming: true, + }), + ); + + const reader = loggedResponse.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toBe( + '{"message":"live"}\n', + ); + await reader.cancel(); + }); + + test("logs bounded non-streaming response bodies", async () => { + const logs: Record[] = []; + await logServerRequest( + new Request("https://kuber.astrxl.dev/health"), + "bounded-response", + () => + new Response("healthy", { + headers: { + "content-length": "7", + "content-type": "text/plain; charset=utf-8", + }, + }), + { log: (entry) => logs.push(entry) }, + ); + + await Bun.sleep(0); + expect(logs).toContainEqual( + expect.objectContaining({ + event: "kuber.server.response.body", + body: { encoding: "base64", data: "aGVhbHRoeQ==" }, + }), + ); + }); +}); diff --git a/tests/server/app.test.ts b/tests/server/app.test.ts index f880a2b..c164d09 100644 --- a/tests/server/app.test.ts +++ b/tests/server/app.test.ts @@ -21,6 +21,52 @@ function request( } describe("kuber API authentication", () => { + test("logs every request start and completion without changing responses", async () => { + const logs: Record[] = []; + const app = createApp({ + store: new MemoryAuthStore(), + requestId: () => "request-log-id", + logger: { + error: () => {}, + log: (entry) => logs.push(entry), + }, + }); + + const response = await app(request("/api/v2/health")); + + expect(response.status).toBe(200); + expect(logs).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + event: "kuber.server.request.start", + requestId: "request-log-id", + method: "GET", + pathname: "/api/v2/health", + }), + expect.objectContaining({ + event: "kuber.server.request.end", + requestId: "request-log-id", + status: 200, + durationMs: expect.any(Number), + }), + ]), + ); + }); + + test("ignores request logging failures", async () => { + const app = createApp({ + store: new MemoryAuthStore(), + logger: { + error: () => {}, + log: () => { + throw new Error("logger unavailable"); + }, + }, + }); + + expect((await app(request("/api/v2/health"))).status).toBe(200); + }); + test("logs in, resolves identity, and revokes the session", async () => { const store = new MemoryAuthStore(); await store.putUser({ diff --git a/tests/server/build-controller.test.ts b/tests/server/build-controller.test.ts index 4daaf15..d9099f7 100644 --- a/tests/server/build-controller.test.ts +++ b/tests/server/build-controller.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, test } from "bun:test"; +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"; @@ -12,7 +12,11 @@ import { } from "../../server/build-controller"; import { createApp } from "../../server/app"; import { hashToken, MemoryAuthStore } from "../../server/auth"; -import { MemoryBuildStore, type BuildStore } from "../../server/build-store"; +import { + MemoryBuildStore, + type BuildRecord, + type BuildStore, +} from "../../server/build-store"; import { FilesystemCas } from "../../server/cas"; import type { KubernetesJob } from "../../server/build-job"; import { @@ -78,6 +82,30 @@ class SupersededBeforeJobStore extends MemoryBuildStore { } } +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(), @@ -277,6 +305,14 @@ describe("build controller", () => { ).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", @@ -449,6 +485,437 @@ describe("build controller", () => { ).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); diff --git a/tests/server/build-job.test.ts b/tests/server/build-job.test.ts index cddbb22..5c62c5b 100644 --- a/tests/server/build-job.test.ts +++ b/tests/server/build-job.test.ts @@ -33,6 +33,18 @@ describe("BuildKit Job generation", () => { const container = pod.containers[0]; expect(pod.nodeSelector["kubernetes.io/arch"]).toBe(architecture); expect(pod.nodeSelector.pool).toBe("builders"); + if (architecture === "amd64") { + expect(pod.tolerations).toEqual([ + { + key: "arch", + operator: "Equal", + value: "amd64", + effect: "NoExecute", + }, + ]); + } else { + expect(pod.tolerations).toBeUndefined(); + } expect(pod.automountServiceAccountToken).toBe(false); expect(pod.securityContext).toMatchObject({ runAsNonRoot: true, @@ -81,6 +93,55 @@ describe("BuildKit Job generation", () => { }, ); + test("merges caller tolerations with the amd64 placement toleration", () => { + const amd64Toleration = { + key: "arch", + operator: "Equal", + value: "amd64", + effect: "NoExecute", + }; + const customToleration = { + key: "workload", + operator: "Equal", + value: "build", + effect: "NoSchedule", + }; + const job = createBuildJob({ + name: "build-tolerations", + namespace: "default", + spec: spec("amd64"), + workspaceClaimName: "workspace", + cacheImage: "registry.example.com/cache/demo-web", + tolerations: [customToleration, amd64Toleration], + }); + + expect((job.spec as any).template.spec.tolerations).toEqual([ + customToleration, + amd64Toleration, + ]); + }); + + test("preserves caller tolerations for arm64 without adding architecture placement", () => { + const customToleration = { + key: "workload", + operator: "Equal", + value: "build", + effect: "NoSchedule", + }; + const job = createBuildJob({ + name: "build-arm64-tolerations", + namespace: "default", + spec: spec("arm64"), + workspaceClaimName: "workspace", + cacheImage: "registry.example.com/cache/demo-web", + tolerations: [customToleration], + }); + + expect((job.spec as any).template.spec.tolerations).toEqual([ + customToleration, + ]); + }); + test("uses pushImage for output when set", () => { const job = createBuildJob({ name: "build-push", diff --git a/tests/server/build-kubernetes.test.ts b/tests/server/build-kubernetes.test.ts index a3b1048..579b0f9 100644 --- a/tests/server/build-kubernetes.test.ts +++ b/tests/server/build-kubernetes.test.ts @@ -78,7 +78,9 @@ class FakeObjects implements BuildObjectApi { readonly selectors: string[] = []; readonly deletes: Array<{ name: string; options?: DeleteOptions }> = []; onLockRead?: () => void; + onReplace?: (value: ConfigMap) => Promise | void; private lockRead = false; + private replacing = false; async create(value: ConfigMap): Promise { if ( @@ -115,6 +117,14 @@ class FakeObjects implements BuildObjectApi { ) ) throw { code: 422 }; + if (!this.replacing && this.onReplace) { + this.replacing = true; + try { + await this.onReplace(value); + } finally { + this.replacing = false; + } + } const current = this.maps.get(value.metadata.name); if ( !current || @@ -346,4 +356,59 @@ describe("KubernetesBuildStore", () => { "1", ); }); + + test("fences a stale replace when an expired lease is taken over", async () => { + const fake = new FakeObjects(); + const kuber = store(fake); + const record = build("fenced", "2026-09-02T00:00:00.000Z"); + record.status.reconcileLease = { + holder: "stale-replica", + token: "stale-token", + expiresAt: new Date(0).toISOString(), + }; + fake.maps.set(name("build", record.metadata.name), configMap(record)); + + const stale = (await kuber.getBuild(record.metadata.name))!; + stale.metadata.resourceVersion = "2"; + stale.status.logBytes = 99; + stale.status.reconcileLease!.expiresAt = "2100-01-01T00:00:00.000Z"; + let successorToken = ""; + fake.onReplace = async () => { + const successor = await kuber.acquireBuildReconciliationLease( + record.metadata.name, + "current-replica", + 60_000, + ); + successorToken = successor!.token; + }; + + await expect(kuber.replaceBuild(stale, "1", "stale-token")).rejects.toThrow( + "concurrently modified", + ); + fake.onReplace = undefined; + + const takenOver = (await kuber.getBuild(record.metadata.name))!; + expect(takenOver.metadata.resourceVersion).toBe("2"); + expect(takenOver.status.logBytes).toBe(0); + expect(takenOver.status.reconcileLease).toMatchObject({ + holder: "current-replica", + token: successorToken, + }); + expect(takenOver.status.reconcileLease?.expiresAt).not.toBe( + "2100-01-01T00:00:00.000Z", + ); + + const current = structuredClone(takenOver); + current.metadata.resourceVersion = "3"; + current.status.logBytes = 7; + await kuber.replaceBuild( + current, + takenOver.metadata.resourceVersion, + successorToken, + ); + expect((await kuber.getBuild(record.metadata.name))?.status).toMatchObject({ + logBytes: 7, + reconcileLease: { token: successorToken }, + }); + }); }); diff --git a/tests/server/build-reconciler.test.ts b/tests/server/build-reconciler.test.ts new file mode 100644 index 0000000..2742f71 --- /dev/null +++ b/tests/server/build-reconciler.test.ts @@ -0,0 +1,103 @@ +import { describe, expect, test } from "bun:test"; +import { + BackgroundBuildReconciler, + DEFAULT_BUILD_RECONCILE_MS, + DEFAULT_BUILD_RECONCILE_TIMEOUT_MS, + MAX_BUILD_RECONCILE_TIMEOUT_MS, + MAX_BUILD_RECONCILE_MS, + buildReconcileIntervalMs, + buildReconcileTimeoutMs, +} from "../../server/build-reconciler"; + +describe("build reconciliation interval", () => { + test("uses positive finite configured intervals and defaults invalid values", () => { + expect(buildReconcileIntervalMs("1500")).toBe(1500); + expect(buildReconcileIntervalMs(String(MAX_BUILD_RECONCILE_MS))).toBe( + MAX_BUILD_RECONCILE_MS, + ); + for (const value of [ + undefined, + "", + "0", + "-1", + "Infinity", + "nope", + String(MAX_BUILD_RECONCILE_MS + 1), + ]) + expect(buildReconcileIntervalMs(value)).toBe(DEFAULT_BUILD_RECONCILE_MS); + }); +}); + +describe("background build reconciliation", () => { + test("does not overlap a timed-out scan and resumes after it settles", async () => { + const scans: AbortSignal[] = []; + let expire!: () => void; + let releaseFirst!: () => void; + let firstSettled!: () => void; + const firstReleased = new Promise((resolve) => { + releaseFirst = resolve; + }); + const settled = new Promise((resolve) => { + firstSettled = resolve; + }); + const failures: unknown[] = []; + const reconciler = new BackgroundBuildReconciler({ + runner: { + reconcilePendingBuilds: async ({ signal } = {}) => { + scans.push(signal!); + if (scans.length === 1) { + await firstReleased; + firstSettled(); + } + }, + }, + timeoutMs: 100, + onFailure: (error) => failures.push(error), + setTimeout: ((callback: () => void) => { + expire = callback; + return 0 as unknown as ReturnType; + }) as typeof setTimeout, + clearTimeout: (() => {}) as typeof clearTimeout, + }); + + reconciler.tick(); + reconciler.tick(); + expect(scans).toHaveLength(1); + + expire(); + expect(scans[0]!.aborted).toBe(true); + reconciler.tick(); + await Promise.resolve(); + + expect(scans).toHaveLength(1); + expect(failures).toHaveLength(1); + + releaseFirst(); + await settled; + await Promise.resolve(); + await Promise.resolve(); + reconciler.tick(); + await Promise.resolve(); + + expect(scans).toHaveLength(2); + }); + + test("uses positive finite configured timeouts and defaults invalid values", () => { + expect(buildReconcileTimeoutMs("1500")).toBe(1500); + expect( + buildReconcileTimeoutMs(String(MAX_BUILD_RECONCILE_TIMEOUT_MS)), + ).toBe(MAX_BUILD_RECONCILE_TIMEOUT_MS); + for (const value of [ + undefined, + "", + "0", + "-1", + "Infinity", + "nope", + String(MAX_BUILD_RECONCILE_TIMEOUT_MS + 1), + ]) + expect(buildReconcileTimeoutMs(value)).toBe( + DEFAULT_BUILD_RECONCILE_TIMEOUT_MS, + ); + }); +});