From 5f7509a9c523d31de5a89cf8d503cec4fc390e78 Mon Sep 17 00:00:00 2001 From: dmgnr Date: Sun, 27 Sep 2026 10:40:50 +0000 Subject: [PATCH] feat: release 2.5.1 --- changelogs/2.5.0.md | 16 + changelogs/2.5.1.md | 8 + command/exec.ts | 4 +- command/logs.ts | 22 ++ command/main.ts | 124 ++++--- command/up.ts | 87 +++-- index.ts | 12 +- lib/api.ts | 10 +- lib/apply.ts | 9 +- lib/build.ts | 173 +++++++--- lib/convert.ts | 2 +- lib/database.ts | 3 +- lib/format.ts | 2 +- lib/graph.ts | 2 +- lib/logger.ts | 20 +- lib/merge.ts | 32 ++ lib/rollback.ts | 3 +- lib/rollout.ts | 212 ++++++++++++ lib/session.ts | 35 +- lib/shared.ts | 34 +- lib/storage.ts | 3 +- lib/workspace.ts | 195 +++++++++-- lib/yaml.ts | 26 +- package.json | 2 +- server/app.ts | 85 +++-- server/audit-store.ts | 5 + server/auth.ts | 26 +- server/build-controller.ts | 2 + server/build-kubernetes.ts | 56 +++ server/build-store.ts | 38 +++ server/kubernetes-logs.ts | 115 ++++--- server/kubernetes-state.ts | 147 ++++---- server/log-service.ts | 7 + server/maintenance.ts | 3 +- server/management.ts | 67 ++-- tests/command/logs.test.ts | 31 ++ tests/command/module-graph.test.ts | 14 + tests/command/startup.test.ts | 81 +++++ tests/command/up-api.test.ts | 89 ++++- tests/lib/api.test.ts | 15 + tests/lib/build-api.test.ts | 268 ++++++++++++++- tests/lib/rollout.test.ts | 151 ++++++++ tests/lib/session.test.ts | 15 + tests/lib/workspace.test.ts | 99 ++++++ tests/server/app.test.ts | 145 +++++++- tests/server/audit-store.test.ts | 14 + tests/server/auth.test.ts | 47 +++ tests/server/build-kubernetes.test.ts | 52 +++ tests/server/build-store.test.ts | 93 +++++ tests/server/kubernetes-audit-store.test.ts | 75 ++++ tests/server/kubernetes-state.test.ts | 65 ++++ tests/server/log-service.test.ts | 361 +++++++++++++++++++- tests/server/management.test.ts | 39 +++ 53 files changed, 2871 insertions(+), 370 deletions(-) create mode 100644 changelogs/2.5.0.md create mode 100644 changelogs/2.5.1.md create mode 100644 lib/merge.ts create mode 100644 lib/rollout.ts create mode 100644 tests/command/module-graph.test.ts create mode 100644 tests/command/startup.test.ts create mode 100644 tests/lib/rollout.test.ts create mode 100644 tests/server/kubernetes-audit-store.test.ts diff --git a/changelogs/2.5.0.md b/changelogs/2.5.0.md new file mode 100644 index 0000000..465aa1c --- /dev/null +++ b/changelogs/2.5.0.md @@ -0,0 +1,16 @@ +# 2.5.0 + +## Added + +- Wait for Kubernetes Deployment rollouts and report rollout progress during `up` and rollback operations. +- Bound concurrent builds and improve build throughput and reconciliation behavior. +- Add `logs -n` to limit the number of returned log lines. +- Defer CLI command-module loading until the selected command starts, reducing startup work. +- Support `up` from Gitless workspaces. +- Add bounded, ring-buffered audit history. +- Retain a configurable history of successful builds and prune older build records. + +## Fixed + +- Correct Kubernetes log streaming and lifecycle handling. +- Repair TLS handling for Kubernetes HTTP connections. diff --git a/changelogs/2.5.1.md b/changelogs/2.5.1.md new file mode 100644 index 0000000..1efd409 --- /dev/null +++ b/changelogs/2.5.1.md @@ -0,0 +1,8 @@ +# 2.5.1 + +## Fixed + +- Harden build reconciliation and lease handling to prevent concurrent updates from overwriting one another. +- Improve workspace snapshot collection for Git repositories, including ignored environment files. +- Improve Kubernetes log streaming, audit persistence, and API request handling. +- Fix CLI module graph and startup behavior, and improve command error handling. diff --git a/command/exec.ts b/command/exec.ts index fea3634..12036c6 100644 --- a/command/exec.ts +++ b/command/exec.ts @@ -25,7 +25,9 @@ export async function runExec( ): Promise { if (!deployment) throw new Error("Deployment name is required"); if (command.length === 0) throw new Error("Command is required"); - const tty = Boolean(process.stdin.isTTY && process.stdout.isTTY); + const tty = Boolean( + process.stdin.isTTY && process.stdout.isTTY && process.stdin.isRaw, + ); const session: ExecApiSession = await opener( project, { deployment, command, tty, ...(tty ? terminalSize() : {}) }, diff --git a/command/logs.ts b/command/logs.ts index 7da43c3..2ab2f74 100644 --- a/command/logs.ts +++ b/command/logs.ts @@ -49,6 +49,13 @@ function defaultWriter(): LogEventWriter { }; } +export function parseLogTailLines(value: string): number { + const tailLines = Number(value); + if (!/^\d+$/.test(value) || !Number.isSafeInteger(tailLines)) + throw new Error("-n must be a non-negative integer"); + return tailLines; +} + export async function runLogs( project: string, service: string | undefined, @@ -56,10 +63,17 @@ export async function runLogs( signal?: AbortSignal, stream: LogsApiStream = apiStreamNdjson, write: LogEventWriter = defaultWriter(), + tailLines?: number, ): Promise { + if ( + tailLines !== undefined && + (!Number.isSafeInteger(tailLines) || tailLines < 0) + ) + throw new Error("-n must be a non-negative integer"); const query = new URLSearchParams(); if (service) query.set("service", service); if (follow) query.set("follow", "true"); + if (tailLines !== undefined) query.set("tailLines", String(tailLines)); const suffix = query.size ? `?${query}` : ""; for await (const event of stream( @@ -87,6 +101,11 @@ export const logs = defineCommand({ alias: "f", description: "Follow log output", }, + tail: { + type: "string", + alias: "n", + description: "Show the last N lines per container", + }, }, async run({ args }) { const service = args._[0]; @@ -100,6 +119,9 @@ export const logs = defineCommand({ service, Boolean(args.follow), controller.signal, + apiStreamNdjson, + defaultWriter(), + args.tail === undefined ? undefined : parseLogTailLines(args.tail), ); } catch (error) { if (!controller.signal.aborted) throw error; diff --git a/command/main.ts b/command/main.ts index d04df49..1abbc86 100644 --- a/command/main.ts +++ b/command/main.ts @@ -1,29 +1,43 @@ -import tab from "@bomb.sh/tab/citty"; import { defineCommand } from "citty"; -import { audit } from "./audit"; -import { ci } from "./ci"; -import { login, logout, whoami } from "./auth"; -import { db } from "./db"; -import { down } from "./down"; -import { exec } from "./exec"; -import { exportCommand } from "./export"; -import { logs } from "./logs"; -import { operations } from "./operations"; -import { ps } from "./ps"; -import { restart } from "./restart"; -import { rollback } from "./rollback"; -import { s3 } from "./s3"; -import { start } from "./start"; -import { stop } from "./stop"; -import { up } from "./up"; -import { users } from "./users"; -import { trust } from "./trust"; -import { maintenance } from "./maintenance"; +import { wrapCommandErrors } from "../lib/error"; + +function lazyCommand( + load: () => Promise, +): () => Promise { + let pending: Promise | undefined; + return () => + (pending ??= load() + .then((command) => { + // Keep the exported command objects pristine: tests and callers may + // import them directly while the CLI needs wrapped handlers. + const cliCommand = cloneCommand(command); + wrapCommandErrors(cliCommand); + return cliCommand; + }) + .catch((error) => { + pending = undefined; + throw error; + })); +} + +function cloneCommand(command: T): T { + const clone = { ...command } as T & { subCommands?: unknown }; + const subCommands = clone.subCommands; + if (subCommands && typeof subCommands === "object") { + clone.subCommands = Object.fromEntries( + Object.entries(subCommands).map(([name, child]) => [ + name, + child && typeof child === "object" ? cloneCommand(child) : child, + ]), + ); + } + return clone; +} export const main = defineCommand({ meta: { name: "kuber", - version: "2.4.0", + version: "2.5.1", description: "Docker Compose -> K8s translation layer", }, args: { @@ -33,33 +47,55 @@ export const main = defineCommand({ }, }, subCommands: { - audit, - ci, - db, - down, - export: exportCommand, - exec, - logs, - maintenance, - login, - logout, - operations, - ps, - restart, - rollback, - s3, - start, - stop, - up, - users, - trust, - whoami, + audit: lazyCommand(() => import("./audit").then((m) => m.audit)), + ci: lazyCommand(() => import("./ci").then((m) => m.ci)), + db: lazyCommand(() => import("./db").then((m) => m.db)), + down: lazyCommand(() => import("./down").then((m) => m.down)), + export: lazyCommand(() => import("./export").then((m) => m.exportCommand)), + exec: lazyCommand(() => import("./exec").then((m) => m.exec)), + logs: lazyCommand(() => import("./logs").then((m) => m.logs)), + maintenance: lazyCommand(() => + import("./maintenance").then((m) => m.maintenance), + ), + login: lazyCommand(() => import("./auth").then((m) => m.login)), + logout: lazyCommand(() => import("./auth").then((m) => m.logout)), + operations: lazyCommand(() => + import("./operations").then((m) => m.operations), + ), + ps: lazyCommand(() => import("./ps").then((m) => m.ps)), + restart: lazyCommand(() => import("./restart").then((m) => m.restart)), + rollback: lazyCommand(() => import("./rollback").then((m) => m.rollback)), + s3: lazyCommand(() => import("./s3").then((m) => m.s3)), + start: lazyCommand(() => import("./start").then((m) => m.start)), + stop: lazyCommand(() => import("./stop").then((m) => m.stop)), + up: lazyCommand(() => import("./up").then((m) => m.up)), + users: lazyCommand(() => import("./users").then((m) => m.users)), + trust: lazyCommand(() => import("./trust").then((m) => m.trust)), + whoami: lazyCommand(() => import("./auth").then((m) => m.whoami)), }, }); -let completion: ReturnType | undefined; +type Completion = Awaited< + ReturnType +>; +let completion: Promise | undefined; export function initializeCompletion() { - completion ??= tab(main); + completion ??= (async () => { + // tab inspects command objects directly, unlike citty's lazy resolver. + const commands = main.subCommands; + if (commands && typeof commands !== "function") { + main.subCommands = Object.fromEntries( + await Promise.all( + Object.entries(commands).map(async ([name, command]) => [ + name, + typeof command === "function" ? await command() : await command, + ]), + ), + ); + } + const { default: tab } = await import("@bomb.sh/tab/citty"); + return tab(main); + })(); return completion; } diff --git a/command/up.ts b/command/up.ts index 2b6efea..22a3a0e 100644 --- a/command/up.ts +++ b/command/up.ts @@ -101,6 +101,7 @@ export type LiveResourceOperationOptions = { type UpContext = { compose?: ComposeSpecification; snapshot?: WorkspaceSnapshot; + workspaceRoot?: string; buildImages?: Record; serviceEnv?: Record>; resources?: KubernetesResource[]; @@ -751,7 +752,14 @@ export async function runUp( { title: "Snapshot workspace", task: async (taskCtx, task) => { - taskCtx.snapshot = await enumerateWorkspace(await getRepoRoot(cwd)); + let workspaceRoot: string; + try { + workspaceRoot = await getRepoRoot(cwd); + } catch (error) { + workspaceRoot = cwd; + } + taskCtx.workspaceRoot = workspaceRoot; + taskCtx.snapshot = await enumerateWorkspace(workspaceRoot); task.output = `${taskCtx.snapshot.manifest.files.length} files`; }, }, @@ -774,7 +782,13 @@ export async function runUp( }, stream: task.stdout(), }, - { ...config, request, snapshot: taskCtx.snapshot }, + { + ...config, + request, + snapshot: taskCtx.snapshot, + workspaceRoot: taskCtx.workspaceRoot, + signal: task.signal, + }, ); taskCtx.buildImages = result.images; await config.postBuild?.(result, await getHookContext()); @@ -822,38 +836,49 @@ export async function runUp( }, }, { - title: "Reconcile databases", - enabled: async () => - getComposePostgresClaims(await compose()).length > 0, + title: "Reconcile backing services", task: async (taskCtx, task) => { - const response = await managementRequest>( - project, - request, - `${workspacePath}/databases`, - { method: "POST", json: { compose: taskCtx.compose } }, - ); + let databaseEnv: Record> = {}; + let storageEnv: Record> = {}; + await task.newListr( + [ + { + title: "Reconcile databases", + enabled: () => + getComposePostgresClaims(taskCtx.compose!).length > 0, + task: async (_ctx, child) => { + const response = await managementRequest>( + project, + request, + `${workspacePath}/databases`, + { method: "POST", json: { compose: taskCtx.compose } }, + ); + databaseEnv = operationEnvironment(response); + child.output = `${Object.keys(databaseEnv).length} services`; + }, + }, + { + title: "Reconcile S3 storage", + enabled: () => + getComposeS3Claims(taskCtx.compose!).length > 0, + task: async (_ctx, child) => { + const response = await managementRequest>( + project, + request, + `${workspacePath}/storage`, + { method: "POST", json: { compose: taskCtx.compose } }, + ); + storageEnv = operationEnvironment(response); + child.output = `${Object.keys(storageEnv).length} services`; + }, + }, + ], + { concurrent: true }, + ).run(); taskCtx.serviceEnv = mergeServiceEnv( - taskCtx.serviceEnv, - operationEnvironment(response), + mergeServiceEnv(taskCtx.serviceEnv, databaseEnv), + storageEnv, ); - task.output = `${Object.keys(taskCtx.serviceEnv).length} services`; - }, - }, - { - title: "Reconcile S3 storage", - enabled: async () => getComposeS3Claims(await compose()).length > 0, - task: async (taskCtx, task) => { - const response = await managementRequest>( - project, - request, - `${workspacePath}/storage`, - { method: "POST", json: { compose: taskCtx.compose } }, - ); - taskCtx.serviceEnv = mergeServiceEnv( - taskCtx.serviceEnv, - operationEnvironment(response), - ); - task.output = `${Object.keys(taskCtx.serviceEnv).length} services`; }, }, { diff --git a/index.ts b/index.ts index eca76e9..377b217 100644 --- a/index.ts +++ b/index.ts @@ -7,10 +7,18 @@ import { extractConfigArgument } from "./lib/config"; import { formatUnknownError, wrapCommandErrors } from "./lib/error"; async function run() { - await initializeCompletion(); + const { configPath, rawArgs } = extractConfigArgument(process.argv.slice(2)); + // citty resolves every preceding lazy command while looking up an alias. + if (rawArgs[0] === "fuck") rawArgs[0] = "rollback"; + if ( + rawArgs.length === 0 || + rawArgs[0] === "complete" || + rawArgs[0] === "--help" || + rawArgs[0] === "-h" + ) + await initializeCompletion(); wrapCommandErrors(main); const cli = createMain(main); - const { configPath, rawArgs } = extractConfigArgument(process.argv.slice(2)); const contextFree = new Set(["login", "logout", "whoami", "maintenance"]); if (rawArgs[0] && contextFree.has(rawArgs[0])) { await cli({ rawArgs }); diff --git a/lib/api.ts b/lib/api.ts index 4c6c92b..c522b31 100644 --- a/lib/api.ts +++ b/lib/api.ts @@ -263,13 +263,15 @@ export async function* apiStreamNdjson( for (;;) { const { done, value } = await reader.read(); buffer += decoder.decode(value, { stream: !done }); - let newline = buffer.indexOf("\n"); + let start = 0; + let newline = buffer.indexOf("\n", start); while (newline !== -1) { - const line = buffer.slice(0, newline).replace(/\r$/, "").trim(); - buffer = buffer.slice(newline + 1); + const line = buffer.slice(start, newline).replace(/\r$/, "").trim(); if (line) yield JSON.parse(line) as T; - newline = buffer.indexOf("\n"); + start = newline + 1; + newline = buffer.indexOf("\n", start); } + buffer = buffer.slice(start); if (done) break; } const finalLine = buffer.replace(/\r$/, "").trim(); diff --git a/lib/apply.ts b/lib/apply.ts index 4889f64..5106d65 100644 --- a/lib/apply.ts +++ b/lib/apply.ts @@ -1,7 +1,5 @@ -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"; @@ -57,6 +55,7 @@ export function getResourceKey(resource: KubernetesObject): string { export async function getNamespaceManagementStatus( namespace: string, ): Promise<"missing" | "managed" | "external"> { + const { objectApi } = await import("./k8s"); try { const existing = await objectApi.read({ apiVersion: "v1", @@ -94,6 +93,10 @@ export async function assertManagedNamespace(namespace: string): Promise { export async function applyResource( resource: T, ): Promise { + const [{ objectApi }, { PatchStrategy }] = await Promise.all([ + import("./k8s"), + import("@kubernetes/client-node"), + ]); return objectApi.patch( sanitizeKubernetesObject(resource), undefined, @@ -119,6 +122,7 @@ export async function applyResources( export async function deleteResource( resource: KubernetesObject, ): Promise { + const { objectApi } = await import("./k8s"); await objectApi.delete(resource); } @@ -133,6 +137,7 @@ export async function deleteResources( export async function listManagedResources( namespace: string, ): Promise { + const { objectApi } = await import("./k8s"); const resources = await Promise.all( ManagedResources.map(async ({ apiVersion, kind }) => { const list = await objectApi.list( diff --git a/lib/build.ts b/lib/build.ts index 7ef5489..e654427 100644 --- a/lib/build.ts +++ b/lib/build.ts @@ -16,6 +16,7 @@ import { apiRequest, type ApiRequestInit, type ApiRequestOptions } from "./api"; import { DEFAULT_REGISTRY } from "./config"; import { enumerateWorkspace, + gitAvailable, serializeWorkspaceManifest, type WorkspaceSnapshot, } from "./workspace"; @@ -25,6 +26,7 @@ const UPLOAD_CHUNK_BYTES = 8 * 1024 * 1024; const DEFAULT_POLL_INTERVAL_MS = 1_000; const BUILD_POLL_REQUEST_TIMEOUT_MS = 300_000; const MAX_BUILD_POLL_ATTEMPTS = 3; +const DEFAULT_CONCURRENT_BUILDS = 3; export const MAX_CONCURRENT_REQUESTS = 20; export const MAX_REQUESTS_PER_SECOND = 40; @@ -157,7 +159,11 @@ export type BuildOptions = { pollIntervalMs?: number; sleep?: (milliseconds: number) => Promise; snapshot?: WorkspaceSnapshot; + workspaceRoot?: string; scheduler?: TaskScheduler; + /** Maximum number of service images built at once. */ + buildConcurrency?: number; + signal?: AbortSignal; }; type BuildPlan = { @@ -301,13 +307,19 @@ export function imageDigestChanged( } export async function getRepoRoot(cwd: string): Promise { - const { stdout } = await execFileAsync("git", [ - "-C", - cwd, - "rev-parse", - "--show-toplevel", - ]); - return stdout.trim(); + if (!(await gitAvailable())) return cwd; + try { + const { stdout } = await execFileAsync("git", [ + "-C", + cwd, + "rev-parse", + "--show-toplevel", + ]); + return stdout.trim(); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return cwd; + throw error; + } } async function uploadBlob( @@ -394,16 +406,25 @@ async function reportBuildEvent( event: BuildEvent, reporter?: BuildReporter, reportedStates?: Set, + service?: string, ): Promise { if (event.type === "status") { if (!reportedStates?.has(event.status.state)) { reportedStates?.add(event.status.state); - await reporter?.progress?.(`Build ${event.status.state}`); + await reporter?.progress?.( + `${service ? `${service}: ` : ""}Build ${event.status.state}`, + ); } return 0; } - if (reporter?.stream) reporter.stream.write(event.message); - else await reporter?.progress?.(event.message.trimEnd()); + if (reporter?.stream) + reporter.stream.write( + service ? `[${service}] ${event.message}` : event.message, + ); + else + await reporter?.progress?.( + `${service ? `${service}: ` : ""}${event.message.trimEnd()}`, + ); return event.sequence; } @@ -434,12 +455,16 @@ async function requestBuildPoll( init: ApiRequestInit | undefined, pollIntervalMs: number, sleep: (milliseconds: number) => Promise, + signal?: AbortSignal, ): Promise { for (let attempt = 1; attempt <= MAX_BUILD_POLL_ATTEMPTS; attempt++) { try { - return await request(path, init, { - timeoutMs: BUILD_POLL_REQUEST_TIMEOUT_MS, - }); + signal?.throwIfAborted(); + return await request( + path, + { ...init, signal }, + { timeoutMs: BUILD_POLL_REQUEST_TIMEOUT_MS }, + ); } catch (error) { if ( !isTransientBuildPollError(error) || @@ -447,6 +472,7 @@ async function requestBuildPoll( ) { throw error; } + signal?.throwIfAborted(); await sleep(pollIntervalMs); } } @@ -460,6 +486,8 @@ async function waitForBuild( pollIntervalMs: number, sleep: (milliseconds: number) => Promise, initial: BuildStatus, + signal?: AbortSignal, + service?: string, ): Promise { let status = initial; let sequence = 0; @@ -471,11 +499,12 @@ async function waitForBuild( undefined, pollIntervalMs, sleep, + signal, ); for (const event of events) sequence = Math.max( sequence, - await reportBuildEvent(event, reporter, reportedStates), + await reportBuildEvent(event, reporter, reportedStates, service), ); if (status.state === "succeeded" || status.state === "failed") return status; @@ -488,9 +517,12 @@ async function waitForBuild( }, pollIntervalMs, sleep, + signal, ); - if (status.state !== "succeeded" && status.state !== "failed") + if (status.state !== "succeeded" && status.state !== "failed") { + signal?.throwIfAborted(); await sleep(pollIntervalMs); + } } } @@ -500,24 +532,44 @@ export async function resolveBuildImages( options: BuildOptions = {}, ): Promise> { const request = options.request ?? apiRequest; - const images: Record = {}; - for (const [service, definition] of Object.entries(compose.services ?? {})) { - if (!definition.build) continue; - try { - images[service] = ( - await request("/images/resolve", { - method: "POST", - json: { project, service }, - }) - ).reference; - } catch (error) { - throw new Error( - `Cannot resolve a published image for service ${service}. Run kuber up to build it.`, - { cause: error }, - ); - } - } - return images; + const services = Object.entries(compose.services ?? {}) + .filter(([, definition]) => definition.build) + .map(([service]) => service); + const references: string[] = new Array(services.length); + let next = 0; + await Promise.all( + Array.from( + { length: Math.min(DEFAULT_CONCURRENT_BUILDS, services.length) }, + async () => { + for (;;) { + const index = next++; + if (index >= services.length) return; + const service = services[index]!; + options.signal?.throwIfAborted(); + try { + references[index] = ( + await request( + "/images/resolve", + { + method: "POST", + json: { project, service }, + signal: options.signal, + }, + ) + ).reference; + } catch (error) { + throw new Error( + `Cannot resolve a published image for service ${service}. Run kuber up to build it.`, + { cause: error }, + ); + } + } + }, + ), + ); + return Object.fromEntries( + services.map((service, index) => [service, references[index]!]), + ); } export async function buildServices( @@ -530,8 +582,20 @@ export async function buildServices( if (!Object.values(compose.services ?? {}).some((service) => service.build)) return { built: [], changed: [], images: {} }; + const concurrency = options.buildConcurrency ?? DEFAULT_CONCURRENT_BUILDS; + if (!Number.isSafeInteger(concurrency) || concurrency < 1) + throw new RangeError("buildConcurrency must be a positive integer"); + const request = options.request ?? apiRequest; - const repoRoot = await getRepoRoot(cwd); + let repoRoot = options.workspaceRoot; + if (!repoRoot && !options.snapshot) { + try { + repoRoot = await getRepoRoot(cwd); + } catch (error) { + repoRoot = cwd; + } + } + repoRoot ??= cwd; const snapshot = options.snapshot ?? (await enumerateWorkspace(repoRoot)); const plans = Object.entries(compose.services ?? {}).flatMap( ([name, service]) => { @@ -554,9 +618,25 @@ export async function buildServices( project, ); - const images: Record = {}; - for (const plan of plans) { - await reporter?.progress?.(`Building ${plan.name}`); + const references: string[] = new Array(plans.length); + let next = 0; + let failure: unknown; + let failed = false; + const worker = async () => { + while (!failed && next < plans.length) { + const index = next++; + const plan = plans[index]!; + try { + options.signal?.throwIfAborted(); + await reporter?.progress?.(`Building ${plan.name}`); + references[index] = await buildPlan(plan); + } catch (error) { + if (!failed) failure = error; + failed = true; + } + } + }; + const buildPlan = async (plan: BuildPlan): Promise => { const id = randomUUID(); const buildRequest: BuildRequest = { version: BUILD_PROTOCOL_VERSION, @@ -578,6 +658,7 @@ export async function buildServices( { method: "POST", json: buildRequest, + signal: options.signal, }, { timeoutMs: 300_000 }, ); @@ -588,15 +669,27 @@ export async function buildServices( options.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS, options.sleep ?? ((milliseconds) => Bun.sleep(milliseconds)), initial, + options.signal, + plan.name, ); if (status.state !== "succeeded") throw new Error( `Build failed for service ${plan.name}: ${status.error ?? "unknown error"}`, ); - images[plan.name] = ( - await request(`/builds/${encodeURIComponent(id)}/result`) + return ( + await request(`/builds/${encodeURIComponent(id)}/result`, { + signal: options.signal, + }) ).reference; - } + }; + + await Promise.all( + Array.from({ length: Math.min(concurrency, plans.length) }, worker), + ); + if (failed) throw failure; + const images = Object.fromEntries( + plans.map((plan, index) => [plan.name, references[index]!]), + ); return { built: plans.map((plan) => plan.name), diff --git a/lib/convert.ts b/lib/convert.ts index efebfa2..7871582 100644 --- a/lib/convert.ts +++ b/lib/convert.ts @@ -23,7 +23,7 @@ import { LABELS } from "../const"; import { getComposeArchPlacement } from "./arch"; import { isPostgresVolumeEntry } from "./database"; import { toEnvVars } from "./format"; -import { deepMerge } from "./shared"; +import { deepMerge } from "./merge"; import { isS3VolumeEntry } from "./storage"; import { LocalArtifactProvider, diff --git a/lib/database.ts b/lib/database.ts index 1965df1..d801f9f 100644 --- a/lib/database.ts +++ b/lib/database.ts @@ -3,7 +3,6 @@ import { randomUUID } from "node:crypto"; import type { ComposeSpecification, Service } from "../schema/docker.d"; import { LABELS } from "../const"; import { deleteResource, applyResource } from "./apply"; -import { objectApi } from "./k8s"; export const DATABASE_NAMESPACE = "database"; export const DATABASE_CLUSTER = "postgres"; @@ -171,6 +170,7 @@ async function readObject( resource: KubernetesObject, ): Promise { try { + const { objectApi } = await import("./k8s"); return (await objectApi.read(resource as never)) as T; } catch (error) { if ( @@ -407,6 +407,7 @@ export async function getRoleCredentials( export async function listManagedDatabaseResources( project: string, ): Promise { + const { objectApi } = await import("./k8s"); const result = await objectApi.list( "postgresql.cnpg.io/v1", "Database", diff --git a/lib/format.ts b/lib/format.ts index a07ceee..5464ea9 100644 --- a/lib/format.ts +++ b/lib/format.ts @@ -1,4 +1,3 @@ -import { Table } from "@cliffy/table"; import type { V1EnvVar } from "@kubernetes/client-node"; import type { Service } from "../schema/docker.d"; @@ -32,3 +31,4 @@ function toEnvVars( } export { toEnvVars, toTable }; +import { Table } from "@cliffy/table"; diff --git a/lib/graph.ts b/lib/graph.ts index f1c8698..7cb500e 100644 --- a/lib/graph.ts +++ b/lib/graph.ts @@ -1,5 +1,4 @@ import type { KubernetesObject } from "@kubernetes/client-node"; -import { objectApi } from "./k8s"; type GraphObject = KubernetesObject & { spec?: Record; @@ -434,6 +433,7 @@ export function buildNamespaceGraphs( export async function fetchNamespaceObjects( namespace: string, ): Promise { + const { objectApi } = await import("./k8s"); const lists = await Promise.all( resources.map(async ([apiVersion, kind]) => { const list = await objectApi.list(apiVersion, kind, namespace); diff --git a/lib/logger.ts b/lib/logger.ts index 2dcd384..1e7eadb 100644 --- a/lib/logger.ts +++ b/lib/logger.ts @@ -1,11 +1,12 @@ //reused logger from other projects but use for only deployment logging import { feature } from "bun:bundle"; -import type { InspectColor } from "node:util"; import { styleText } from "node:util"; import { inspect } from "bun"; +type InspectColor = Extract[0], string>; + const colors: InspectColor[] = [ "red", "green", @@ -42,6 +43,23 @@ export function createLogger(label?: string) { i = (i + 1) % colors.length; function mkLogFn(stream: Console | NodeJS.WriteStream, strip?: InspectColor) { return (strings: TemplateStringsArray, ...values: unknown[]) => { + const tty = stream === console ? process.stdout.isTTY : process.stderr.isTTY; + if (!tty) { + const prefix = `${(label ?? `${Date.now() - start}ms`).padStart(7)} ▌${strip ? "▌" : " "}`; + let message = strings.reduce( + (out, str, index) => + out + str + + (typeof values[index] === "string" + ? values[index] + : typeof values[index] !== "undefined" + ? inspect(values[index], { colors: false }) + : ""), + "", + ); + message = message.split("\n").join(`\n${" ".repeat(7)} ▌${strip ? "▌" : " "}`); + (stream.write as (data: string) => void)(`${prefix}${message}\n`); + return; + } const normalized = styleText( col, `${(strip ? styleText : (_: string, s: string) => s)("bold", (label ?? `${Date.now() - start}ms`).padStart(7))} ▌${strip ? styleText(strip, "▌") : " "}${strings diff --git a/lib/merge.ts b/lib/merge.ts new file mode 100644 index 0000000..9275f69 --- /dev/null +++ b/lib/merge.ts @@ -0,0 +1,32 @@ +export function deepMerge(target: any, source: any): any { + if ( + typeof target !== "object" || + target === null || + typeof source !== "object" || + source === null + ) { + return source; + } + + const result = Array.isArray(target) ? [...target] : { ...target }; + + for (const key of Object.keys(source)) { + const targetValue = result[key]; + const sourceValue = source[key]; + + if (Array.isArray(targetValue) && Array.isArray(sourceValue)) { + result[key] = [...targetValue, ...sourceValue]; + } else if ( + targetValue && + typeof targetValue === "object" && + sourceValue && + typeof sourceValue === "object" + ) { + result[key] = deepMerge(targetValue, sourceValue); + } else { + result[key] = sourceValue; + } + } + + return result; +} diff --git a/lib/rollback.ts b/lib/rollback.ts index 34c527d..f2e0f91 100644 --- a/lib/rollback.ts +++ b/lib/rollback.ts @@ -5,7 +5,6 @@ import type { } from "@kubernetes/client-node"; import { LABELS } from "../const"; import { ctx } from "./context"; -import { apps } from "./k8s"; import { listManagedDeployments, waitForDeploymentRollout } from "./shared"; const RevisionAnnotation = "deployment.kubernetes.io/revision"; @@ -51,6 +50,7 @@ export function templateMatches( export async function listDeploymentReplicaSets( deployment: V1Deployment, ): Promise { + const { apps } = await import("./k8s"); const name = deployment.metadata?.name; const ownerUid = deployment.metadata?.uid; if (!name || !ownerUid) return []; @@ -130,6 +130,7 @@ export async function planRollback( export async function rollbackDeployment( candidate: RollbackCandidate, ): Promise { + const { apps } = await import("./k8s"); const { name, previousRevision } = candidate; const deployment = await apps.readNamespacedDeployment({ namespace: ctx().project, diff --git a/lib/rollout.ts b/lib/rollout.ts new file mode 100644 index 0000000..219c505 --- /dev/null +++ b/lib/rollout.ts @@ -0,0 +1,212 @@ +import type { + V1Deployment, + V1Pod, + V1ReplicaSet, + V1ContainerStatus, +} from "@kubernetes/client-node"; + +const safe = (value: string | undefined) => + (value ?? "unknown").replace(/[^a-zA-Z0-9._/-]/g, "_").slice(0, 100); + +function containsTemplate(actual: unknown, desired: unknown): boolean { + if (Array.isArray(desired)) + return ( + Array.isArray(actual) && + actual.length === desired.length && + desired.every((value, index) => containsTemplate(actual[index], value)) + ); + if (desired && typeof desired === "object") + return ( + Object.keys(desired).length === 0 || + Boolean( + actual && + typeof actual === "object" && + Object.entries(desired).every(([key, value]) => + containsTemplate((actual as Record)[key], value), + ), + ) + ); + return actual === desired; +} + +export function rolloutReady(deployment: V1Deployment): boolean { + const desired = deployment.spec?.replicas ?? 1; + return ( + (deployment.status?.observedGeneration ?? 0) >= + (deployment.metadata?.generation ?? 0) && + (deployment.status?.updatedReplicas ?? 0) === desired && + (deployment.status?.availableReplicas ?? 0) === desired && + (deployment.status?.unavailableReplicas ?? 0) === 0 + ); +} + +export function rolloutConditionFailure( + deployment: V1Deployment, +): string | undefined { + if ( + (deployment.status?.observedGeneration ?? 0) < + (deployment.metadata?.generation ?? 0) + ) + return; + const condition = deployment.status?.conditions?.find( + (item) => + item.type === "Progressing" && + item.status === "False" && + item.reason === "ProgressDeadlineExceeded", + ); + if (condition) + return `Deployment ${safe(deployment.metadata?.name)} rollout failed: ${safe(condition.reason)}`; +} + +// An older ReplicaSet may still have crashing pods while its replacement is being +// created. Require both Deployment ownership and a match against the desired template. +export function currentReplicaSet( + deployment: V1Deployment, + sets: V1ReplicaSet[], +): V1ReplicaSet | undefined { + const uid = deployment.metadata?.uid; + const template = deployment.spec?.template; + if (!uid || !template) return; + const matches = sets.filter((set) => { + if ( + !set.metadata?.ownerReferences?.some( + (owner) => owner.kind === "Deployment" && owner.uid === uid, + ) + ) + return false; + const actual = set.spec?.template; + if (!actual) return false; + return ( + containsTemplate( + actual.metadata?.labels, + template.metadata?.labels ?? {}, + ) && + containsTemplate( + actual.metadata?.annotations, + template.metadata?.annotations ?? {}, + ) && + containsTemplate(actual.spec, template.spec) + ); + }); + return matches.sort( + (a, b) => + Number( + b.metadata?.annotations?.["deployment.kubernetes.io/revision"] ?? 0, + ) - + Number( + a.metadata?.annotations?.["deployment.kubernetes.io/revision"] ?? 0, + ), + )[0]; +} + +const failureReasons = new Set([ + "CrashLoopBackOff", + "ImagePullBackOff", + "ErrImagePull", + "InvalidImageName", + "CreateContainerConfigError", + "CreateContainerError", + "RunContainerError", + "CreateContainerRuntimeError", + "ErrImageNeverPull", + "OOMKilled", +]); + +function containerFailure(status: V1ContainerStatus): string | undefined { + const waiting = status.state?.waiting; + if (waiting?.reason && failureReasons.has(waiting.reason)) + return waiting.reason; + const terminated = status.state?.terminated ?? status.lastState?.terminated; + if (terminated?.reason === "OOMKilled") return "OOMKilled"; + if ( + terminated && + terminated.exitCode !== 0 && + (status.restartCount ?? 0) >= 2 + ) + return terminated.reason || "RepeatedNonzeroExit"; +} + +export function rolloutPodFailure( + deployment: V1Deployment, + replicaSet: V1ReplicaSet, + pods: V1Pod[], +): string | undefined { + const uid = replicaSet.metadata?.uid; + if (!uid) return; + for (const pod of pods) { + if ( + !pod.metadata?.ownerReferences?.some( + (owner) => owner.kind === "ReplicaSet" && owner.uid === uid, + ) + ) + continue; + for (const status of [ + ...(pod.status?.initContainerStatuses ?? []), + ...(pod.status?.containerStatuses ?? []), + ]) { + const reason = containerFailure(status); + if (reason) + return `Deployment ${safe(deployment.metadata?.name)} rollout failed: pod ${safe(pod.metadata?.name)}, container ${safe(status.name)}: ${safe(reason)} (exit ${status.state?.terminated?.exitCode ?? status.lastState?.terminated?.exitCode ?? "unknown"}, restarts ${status.restartCount ?? 0})`; + } + } +} + +export async function waitForRollout( + name: string, + deadline: number, + read: () => Promise, + listSets: () => Promise, + listPods: () => Promise, + signal?: AbortSignal, +): Promise { + const cancelled = () => + new Error("Workspace operation execution was cancelled"); + const timedOut = () => + new Error(`Timed out waiting for deployment ${safe(name)} rollout`); + async function bounded(operation: () => Promise): Promise { + if (signal?.aborted) throw cancelled(); + const remaining = deadline - Date.now(); + if (remaining <= 0) throw timedOut(); + return new Promise((resolve, reject) => { + const timer = setTimeout( + () => finish(() => reject(timedOut())), + remaining, + ); + const abort = () => finish(() => reject(cancelled())); + const finish = (settle: () => void) => { + clearTimeout(timer); + signal?.removeEventListener("abort", abort); + settle(); + }; + signal?.addEventListener("abort", abort, { once: true }); + if (signal?.aborted) return abort(); + Promise.resolve() + .then(operation) + .then( + (value) => finish(() => resolve(value)), + (error) => finish(() => reject(error)), + ); + }); + } + while (true) { + const deployment = await bounded(read); + const failure = rolloutConditionFailure(deployment); + if (failure) throw new Error(failure); + if (rolloutReady(deployment)) return; + const current = currentReplicaSet(deployment, await bounded(listSets)); + if (current) { + const podFailure = rolloutPodFailure( + deployment, + current, + await bounded(listPods), + ); + if (podFailure) throw new Error(podFailure); + } + await bounded( + () => + new Promise((resolve) => + setTimeout(resolve, Math.min(2_000, deadline - Date.now())), + ), + ); + } +} diff --git a/lib/session.ts b/lib/session.ts index 0d69865..8975d4b 100644 --- a/lib/session.ts +++ b/lib/session.ts @@ -19,6 +19,15 @@ export type KuberSession = { }; }; +let cachedSession: + | { + runtimePath: string; + persistentPath: string; + value: KuberSession | undefined; + expiresAt: number; + } + | undefined; + function runtimeSessionPath(): string { const directory = process.env.XDG_RUNTIME_DIR ?? @@ -76,10 +85,26 @@ async function readSessionFile( } export async function readSession(): Promise { - return ( - (await readSessionFile(runtimeSessionPath())) ?? - (await readSessionFile(persistentSessionPath())) - ); + const runtimePath = runtimeSessionPath(); + const persistentPath = persistentSessionPath(); + if ( + cachedSession?.runtimePath === runtimePath && + cachedSession.persistentPath === persistentPath && + (cachedSession.expiresAt === 0 || cachedSession.expiresAt > Date.now()) + ) + return cachedSession.value; + + const value = + (await readSessionFile(runtimePath)) ?? + (await readSessionFile(persistentPath)); + cachedSession = { + runtimePath, + persistentPath, + value, + // Expiration is checked on every lookup, even while the file is memoized. + expiresAt: value ? Date.parse(value.expiresAt) : 0, + }; + return value; } export async function writeSession( @@ -97,6 +122,7 @@ export async function writeSession( await rename(temporaryPath, path); await chmod(path, 0o600); await rm(getSessionPath(!persistent), { force: true }); + cachedSession = undefined; return path; } @@ -106,4 +132,5 @@ export async function removeSessions(): Promise { rm(path, { force: true }), ), ); + cachedSession = undefined; } diff --git a/lib/shared.ts b/lib/shared.ts index f38c89a..16b3a64 100644 --- a/lib/shared.ts +++ b/lib/shared.ts @@ -2,6 +2,7 @@ import type { V1Deployment, V1Pod } from "@kubernetes/client-node"; import { LABELS } from "../const"; import { ctx } from "./context"; import { apps, core, objectApi } from "./k8s"; +import { waitForRollout } from "./rollout"; const ManagedBySelector = `app.kubernetes.io/managed-by=${LABELS["app.kubernetes.io/managed-by"]}`; @@ -95,30 +96,15 @@ export async function waitForDeploymentRollout( name: string, timeoutMs = 300000, ) { - const startedAt = Date.now(); - - while (Date.now() - startedAt < timeoutMs) { - const deployment = await getDeployment(name); - const desiredReplicas = deployment.spec?.replicas ?? 1; - const observedGeneration = deployment.status?.observedGeneration ?? 0; - const generation = deployment.metadata?.generation ?? 0; - const updatedReplicas = deployment.status?.updatedReplicas ?? 0; - const availableReplicas = deployment.status?.availableReplicas ?? 0; - const unavailableReplicas = deployment.status?.unavailableReplicas ?? 0; - - if ( - observedGeneration >= generation && - updatedReplicas === desiredReplicas && - availableReplicas === desiredReplicas && - unavailableReplicas === 0 - ) { - return; - } - - await delay(2000); - } - - throw new Error(`Timed out waiting for deployment ${name} rollout`); + return waitForRollout( + name, + Date.now() + timeoutMs, + () => getDeployment(name), + async () => + (await apps.listNamespacedReplicaSet({ namespace: getProject() })).items, + async () => + (await core.listNamespacedPod({ namespace: getProject() })).items, + ); } export function throwWhen( diff --git a/lib/storage.ts b/lib/storage.ts index 89a2809..6a07675 100644 --- a/lib/storage.ts +++ b/lib/storage.ts @@ -2,7 +2,6 @@ import type { KubernetesObject, V1Secret } from "@kubernetes/client-node"; import type { ComposeSpecification, Service } from "../schema/docker.d"; import { LABELS } from "../const"; import { applyResource } from "./apply"; -import { objectApi } from "./k8s"; export const STORAGE_NAMESPACE = "garage-system"; export const STORAGE_CLUSTER = "garage"; @@ -37,6 +36,7 @@ async function readObject( resource: KubernetesObject, ): Promise { try { + const { objectApi } = await import("./k8s"); return (await objectApi.read(resource as never)) as T; } catch (error) { if ( @@ -342,6 +342,7 @@ export async function reconcileS3Claims( export async function listManagedStorageResources( project: string, ): Promise { + const { objectApi } = await import("./k8s"); const selector = `${Object.entries(LABELS) .map(([key, value]) => `${key}=${value}`) .join(",")},${STORAGE_PROJECT_LABEL}=${project}`; diff --git a/lib/workspace.ts b/lib/workspace.ts index c451eed..d0b2d82 100644 --- a/lib/workspace.ts +++ b/lib/workspace.ts @@ -1,6 +1,7 @@ -import { execFile } from "node:child_process"; +import { execFile, spawn } from "node:child_process"; import { createHash } from "node:crypto"; import { constants } from "node:fs"; +import { access } from "node:fs/promises"; import { chmod, lstat, @@ -10,9 +11,10 @@ import { readlink, realpath, symlink, + readFile, } from "node:fs/promises"; import type { Stats } from "node:fs"; -import { dirname, isAbsolute, relative, resolve, sep } from "node:path"; +import { delimiter, dirname, isAbsolute, relative, resolve, sep } from "node:path"; import { promisify } from "node:util"; import { BUILD_PROTOCOL_VERSION, @@ -24,6 +26,17 @@ import { const execFileAsync = promisify(execFile); +export async function gitAvailable(): Promise { + for (const directory of (process.env.PATH ?? "").split(delimiter)) { + if (!directory) continue; + try { + await access(resolve(directory, process.platform === "win32" ? "git.exe" : "git"), constants.X_OK); + return true; + } catch { /* Continue searching PATH. */ } + } + return false; +} + export type WorkspaceBlob = { digest: Sha256Digest; data: Uint8Array; @@ -103,31 +116,166 @@ async function selectedFiles(root: string): Promise { ); } -async function isIgnored(root: string, path: string): Promise { +async function filesystemFiles(root: string): Promise { + const files: string[] = []; + const ignoredDirectories = new Set([ + ".git", ".hg", ".svn", "node_modules", "vendor", "bower_components", + ".venv", "venv", "__pycache__", ".tox", ".mypy_cache", ".pytest_cache", + ".next", ".nuxt", ".svelte-kit", ".cache", ".turbo", "dist", "build", "coverage", + "target", "out", "tmp", "temp", + ]); + const sensitiveDirectory = /(?:^|[-_.])(?:secrets?|credentials?|configs?)(?:$|[-_.])/i; + const ignoreRules: Array<{ base: string; pattern: string; directory: boolean }> = []; + const loadIgnore = async (directory: string): Promise => { + try { + const content = await readFile(resolve(directory, ".gitignore"), "utf8"); + for (const raw of content.split(/\r?\n/)) { + const line = raw.trim(); + if (!line || line.startsWith("#")) continue; + // Negations are skipped: safely re-including descendants requires + // Git's parent-directory semantics, so fallback stays fail-closed. + if (line.startsWith("!")) continue; + const rule = line.replace(/^\//, ""); + if (rule) ignoreRules.push({ base: relative(root, directory).split(sep).join("/"), pattern: rule.replace(/\/$/, ""), directory: line.endsWith("/") }); + } + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + }; + const ignoredByRules = (path: string, isDirectory: boolean): boolean => { + let ignored = false; + for (const rule of ignoreRules) { + const prefix = rule.base ? `${rule.base}/` : ""; + if (rule.base && path !== rule.base && !path.startsWith(prefix)) continue; + const local = rule.base && path.startsWith(prefix) ? path.slice(prefix.length) : path; + const glob = rule.pattern.replace(/[.+^${}()|[\]\\]/g, "\\$&").replace(/\*\*/g, "__DOUBLESTAR__").replace(/\*/g, "[^/]*").replace(/\?/g, "[^/]").replace(/__DOUBLESTAR__/g, ".*"); + const matcher = new RegExp(`^(?:${glob})(?:/.*)?$`); + const basenameMatcher = new RegExp(`^(?:${glob})$`); + if ((matcher.test(local) || local.split("/").some((part) => basenameMatcher.test(part))) && (!rule.directory || isDirectory || local.includes("/"))) ignored = true; + } + return ignored; + }; + const visit = async (directory: string): Promise => { + await loadIgnore(directory); + for (const entry of await readdir(directory, { withFileTypes: true })) { + const source = resolve(directory, entry.name); + const path = relative(root, source).split(sep).join("/"); + if (entry.isDirectory()) { + if (ignoredDirectories.has(entry.name) || sensitiveDirectory.test(entry.name) || ignoredByRules(path, true)) continue; + await visit(source); + } + else if (entry.isFile() || entry.isSymbolicLink()) { + if (ignoredByRules(path, false)) continue; + if (/^\.env/i.test(entry.name) || /(?:secret|credential|password|token|private[-_.]?key)/i.test(entry.name) || /^(?:id_rsa|id_ed25519|known_hosts|config\.json|\.npmrc|\.pypirc|\.netrc)$/i.test(entry.name)) continue; + files.push(path); + } else { + throw new Error(`Special files are not allowed in workspaces: ${relative(root, source).split(sep).join("/")}`); + } + } + }; + await visit(root); + return files.sort((a, b) => Buffer.from(a).compare(Buffer.from(b))); +} + +async function isGitRepository(root: string): Promise { + if (!(await gitAvailable())) return false; try { - await execFileAsync("git", ["-C", root, "check-ignore", "-q", "--", path]); + await execFileAsync("git", ["-C", root, "rev-parse", "--show-toplevel"]); return true; } catch (error) { - if ((error as { code?: number }).code === 1) return false; + const code = (error as { code?: unknown }).code; + const exitCode = (error as { exitCode?: number }).exitCode; + const message = error instanceof Error ? error.message : String(error); + if (code === "ENOENT") return false; + if (exitCode === 128 || code === 128 || /not a git repository/i.test(message)) return false; throw error; } } -async function rejectSelectedSpecialFiles( - root: string, - directory = root, -): Promise { - for (const entry of await readdir(directory, { withFileTypes: true })) { - if (directory === root && entry.name === ".git") continue; - const source = resolve(directory, entry.name); - const path = relative(root, source).split(sep).join("/"); - if (entry.isDirectory()) { - if (!(await isIgnored(root, path))) - await rejectSelectedSpecialFiles(root, source); - continue; +async function rejectSelectedSpecialFiles(root: string): Promise { + // Let Git identify ignored directories using its own ignore engine. With + // --directory, ignored trees are returned as directory entries, so the + // filesystem walk below can prune them without visiting their contents. + const ignoredResult = await execFileAsync( + "git", + [ + "-C", + root, + "ls-files", + "--others", + "--ignored", + "--exclude-standard", + "--directory", + "-z", + ], + { encoding: "buffer", maxBuffer: 64 * 1024 * 1024 }, + ); + const decoder = new TextDecoder("utf-8", { fatal: true }); + const ignoredDirectories = new Set(); + let start = 0; + for ( + let end = ignoredResult.stdout.indexOf(0); + end !== -1; + end = ignoredResult.stdout.indexOf(0, start) + ) { + if (end > start) { + const path = decoder.decode(ignoredResult.stdout.subarray(start, end)); + if (path.endsWith("/")) ignoredDirectories.add(path.slice(0, -1)); } - if (entry.isFile() || entry.isSymbolicLink()) continue; - if (!(await isIgnored(root, path)) || entry.name.startsWith(".env")) + start = end + 1; + } + + const specialPaths: string[] = []; + const visit = async (directory: string): Promise => { + for (const entry of await readdir(directory, { withFileTypes: true })) { + if (directory === root && entry.name === ".git") continue; + const source = resolve(directory, entry.name); + const path = relative(root, source).split(sep).join("/"); + if (entry.isDirectory()) { + if (!ignoredDirectories.has(path)) await visit(source); + } else if (!entry.isFile() && !entry.isSymbolicLink()) { + specialPaths.push(path); + } + } + }; + await visit(root); + + if (specialPaths.length === 0) return; + const input = Buffer.from(`${specialPaths.join("\0")}\0`); + const child = spawn("git", ["-C", root, "check-ignore", "--stdin", "-z"]); + const output: Buffer[] = []; + const errors: Buffer[] = []; + child.stdout.on("data", (chunk: Buffer) => output.push(chunk)); + child.stderr.on("data", (chunk: Buffer) => errors.push(chunk)); + const completed = new Promise((resolveExit, rejectExit) => { + child.once("error", rejectExit); + child.once("close", (code) => { + if (code === 0 || code === 1) resolveExit(); + else + rejectExit( + new Error( + Buffer.concat(errors).toString("utf8") || + `git check-ignore exited with ${code}`, + ), + ); + }); + }); + child.stdin.end(input); + await completed; + const stdout = Buffer.concat(output); + const ignored = new Set(); + start = 0; + for ( + let end = stdout.indexOf(0); + end !== -1; + end = stdout.indexOf(0, start) + ) { + if (end > start) ignored.add(decoder.decode(stdout.subarray(start, end))); + start = end + 1; + } + for (const path of specialPaths) { + const name = path.slice(path.lastIndexOf("/") + 1); + if (!ignored.has(path) || name.startsWith(".env")) throw new Error(`Special files are not allowed in workspaces: ${path}`); } } @@ -201,8 +349,13 @@ export async function enumerateWorkspace( root: string, ): Promise { const repository = await realpath(root); - await rejectSelectedSpecialFiles(repository); - const paths = await selectedFiles(repository); + const gitBacked = await isGitRepository(repository); + const paths = gitBacked + ? await (async () => { + await rejectSelectedSpecialFiles(repository); + return selectedFiles(repository); + })() + : await filesystemFiles(repository); const files: WorkspaceFile[] = []; const blobs = new Map(); diff --git a/lib/yaml.ts b/lib/yaml.ts index 6b2a547..a416007 100644 --- a/lib/yaml.ts +++ b/lib/yaml.ts @@ -2,7 +2,6 @@ import z from "zod"; import { readdir } from "node:fs/promises"; import { join } from "node:path"; import type { ComposeSpecification } from "../schema/docker.d"; -import JsonSchema from "../schema/docker.ts"; import { file, YAML } from "bun"; const ComposeFileNames = new Set([ @@ -12,12 +11,29 @@ const ComposeFileNames = new Set([ "docker-compose.yaml", ]); -const ComposeSchema = z.fromJSONSchema( - JsonSchema as unknown as Parameters[0], -) as z.ZodType; +let composeSchema: Promise> | undefined; + +function getComposeSchema(): Promise> { + if (!composeSchema) { + composeSchema = import("../schema/docker.ts") + .then(({ default: schema }) => + z.fromJSONSchema( + schema as unknown as Parameters[0], + ) as z.ZodType, + ) + .catch((error) => { + composeSchema = undefined; + throw error; + }); + } + return composeSchema; +} export function readCompose(path: string): Promise { - return file(path).text().then(YAML.parse).then(ComposeSchema.parse); + return file(path) + .text() + .then(YAML.parse) + .then(async (value) => (await getComposeSchema()).parse(value)); } export async function resolveComposeFile( diff --git a/package.json b/package.json index 2bdb683..8cd4e10 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.4.2", + "version": "2.5.1", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/server/app.ts b/server/app.ts index 484352e..a1b3a2d 100644 --- a/server/app.ts +++ b/server/app.ts @@ -51,6 +51,7 @@ import { sanitizeOperationResult, type Operation, type OperationStore, + type WorkspaceLease, type WorkspaceLeaseProvider, } from "./operation-store"; import { @@ -982,13 +983,45 @@ export function createApp( location: `${API_PREFIX}/operations/${operation.metadata.name}`, }); - const lease = options.leases - ? await options.leases.acquire( - workspace.metadata.name, - operation.metadata.name, - ) - : undefined; - if (options.leases && !lease) { + // Waiting for readiness only observes Kubernetes state; it cannot mutate + // resources, so it must not occupy the workspace mutation lease. + const leaseScopes = action === "resources/wait" + ? [] + : action === "databases.reconcile" + ? [`${workspace.metadata.name}:scope:database`] + : action === "storage.reconcile" + ? [`${workspace.metadata.name}:scope:storage`] + : [ + workspace.metadata.name, + `${workspace.metadata.name}:scope:database`, + `${workspace.metadata.name}:scope:storage`, + ]; + const acquiredLeases: WorkspaceLease[] = []; + if (options.leases) { + try { + for (const scope of [...leaseScopes].sort()) { + const acquired = await options.leases.acquire( + scope, + operation.metadata.name, + ); + if (!acquired) { + throw new HttpError( + 409, + "Conflict", + "WORKSPACE_BUSY", + "Another workspace operation is running", + undefined, + operation.metadata.name, + ); + } + acquiredLeases.push(acquired); + } + } catch (error) { + await Promise.allSettled(acquiredLeases.map((lease) => lease.release())); + throw error; + } + } + if (leaseScopes.length > 0 && options.leases && acquiredLeases.length === 0) { throw new HttpError( 409, "Conflict", @@ -998,7 +1031,7 @@ export function createApp( operation.metadata.name, ); } - if (!lease && options.leases) throw new Error("Unreachable lease state"); + const leases = acquiredLeases; let operationStarted = false; let leaseRenewalTimer: ReturnType | undefined; @@ -1006,11 +1039,12 @@ export function createApp( let renewalInFlight: Promise | undefined; const executionController = new AbortController(); const renewLease = async (): Promise => { - if (!lease || leaseOwnershipLost) return false; + if (leases.length === 0 || leaseOwnershipLost) return false; if (renewalInFlight) return renewalInFlight; - renewalInFlight = lease - .renew(WORKSPACE_LEASE_TTL_MS) - .catch(() => false) + renewalInFlight = Promise.all( + leases.map((lease) => lease.renew(WORKSPACE_LEASE_TTL_MS).catch(() => false)), + ) + .then((renewals) => renewals.every(Boolean)) .then((renewed) => { if (!renewed) { leaseOwnershipLost = true; @@ -1024,13 +1058,13 @@ export function createApp( return renewalInFlight; }; const scheduleLeaseRenewal = () => { - if (!lease || leaseOwnershipLost) return; + if (leases.length === 0 || leaseOwnershipLost) return; leaseRenewalTimer = setTimeout(() => { void renewLease().finally(scheduleLeaseRenewal); }, WORKSPACE_LEASE_RENEW_INTERVAL_MS); }; const requireLeaseOwnership = async () => { - if (lease && !(await renewLease())) + if (leases.length > 0 && !(await renewLease())) throw new HttpError( 409, "Conflict", @@ -1150,16 +1184,19 @@ export function createApp( ); } finally { if (leaseRenewalTimer !== undefined) clearTimeout(leaseRenewalTimer); - try { - await lease?.release(); - } catch (error) { - logRequestError(request, error, { - event: "operation.lease_release.failed", - code: "LEASE_RELEASE_FAILED", - message: "Workspace operation lease could not be released", - workspaceId: workspace.metadata.name, - operationId: operation.metadata.name, - }); + const releases = await Promise.allSettled( + leases.map((lease) => lease.release()), + ); + for (const release of releases) { + if (release.status === "rejected") { + logRequestError(request, release.reason, { + event: "operation.lease_release.failed", + code: "LEASE_RELEASE_FAILED", + message: "Workspace operation lease could not be released", + workspaceId: workspace.metadata.name, + operationId: operation.metadata.name, + }); + } } } }; diff --git a/server/audit-store.ts b/server/audit-store.ts index 5784d90..225a25e 100644 --- a/server/audit-store.ts +++ b/server/audit-store.ts @@ -4,6 +4,7 @@ import { redactString, REDACTED } from "./redact"; export const AUDIT_REDACTED = REDACTED; export const MAX_AUDIT_EVENT_BYTES = 256 * 1024; +export const MAX_AUDIT_EVENTS = 100; const SENSITIVE_KEY = /(?:authorization|cookie|credential|password|passwd|secret|token|api[-_]?key|private[-_]?key)/i; @@ -125,6 +126,10 @@ export class MemoryAuditPersistence implements AuditPersistence { async append(event: AuditEvent) { this.events.push(clone(event)); + this.events.sort((a, b) => + a.metadata.creationTimestamp.localeCompare(b.metadata.creationTimestamp), + ); + while (this.events.length > MAX_AUDIT_EVENTS) this.events.shift(); } async list(workspaceId?: string) { diff --git a/server/auth.ts b/server/auth.ts index 969b3f9..cba7756 100644 --- a/server/auth.ts +++ b/server/auth.ts @@ -166,6 +166,20 @@ export class MemoryAuthStore implements AuthStore { readonly users = new Map(); readonly sessions = new Map(); readonly apiKeys = new Map(); + readonly apiKeysByTokenHash = new Map(); + + private storeApiKey(key: ApiKeyRecord): void { + this.apiKeys.set(key.id, key); + this.apiKeysByTokenHash.set(key.tokenHash, key); + } + + private deleteApiKey(id: string): void { + const key = this.apiKeys.get(id); + if (!key) return; + this.apiKeys.delete(id); + if (this.apiKeysByTokenHash.get(key.tokenHash) === key) + this.apiKeysByTokenHash.delete(key.tokenHash); + } async getUser(username: string): Promise { return this.users.get(username); @@ -208,7 +222,7 @@ export class MemoryAuthStore implements AuthStore { async deleteUser(username: string): Promise { await this.revokeUserSessions(username); for (const [id, key] of this.apiKeys) - if (key.username === username) this.apiKeys.delete(id); + if (key.username === username) this.deleteApiKey(id); return this.users.delete(username); } @@ -264,9 +278,7 @@ export class MemoryAuthStore implements AuthStore { } async getApiKey(tokenHash: string): Promise { - const key = [...this.apiKeys.values()].find( - (item) => item.tokenHash === tokenHash, - ); + const key = this.apiKeysByTokenHash.get(tokenHash); if (!key || key.disabled || Date.parse(key.expiresAt) <= Date.now()) return; const user = this.users.get(key.username); if (!user || user.disabled) return; @@ -279,7 +291,7 @@ export class MemoryAuthStore implements AuthStore { if (!user || user.disabled) throw new Error("API key user is not active"); if (this.apiKeys.has(normalized.id)) throw new Error("API key already exists"); - this.apiKeys.set(normalized.id, normalized); + this.storeApiKey(normalized); } async listApiKeys(username: string): Promise { @@ -291,7 +303,7 @@ export class MemoryAuthStore implements AuthStore { async revokeApiKey(username: string, id: string): Promise { const key = this.apiKeys.get(id); if (!key || key.username !== username) return false; - this.apiKeys.delete(id); + this.deleteApiKey(id); return true; } @@ -299,7 +311,7 @@ export class MemoryAuthStore implements AuthStore { const expired = [...this.apiKeys.values()].filter( (key) => Date.parse(key.expiresAt) <= now, ); - for (const key of expired) this.apiKeys.delete(key.id); + for (const key of expired) this.deleteApiKey(key.id); return expired.length; } } diff --git a/server/build-controller.ts b/server/build-controller.ts index 76820f5..884e1c7 100644 --- a/server/build-controller.ts +++ b/server/build-controller.ts @@ -761,6 +761,8 @@ export class BuildController { } } } + throwIfAborted(options.signal); + await this.options.store.pruneBuilds?.(50); } async cancelBuild(id: string): Promise { diff --git a/server/build-kubernetes.ts b/server/build-kubernetes.ts index 57f29ea..42154bc 100644 --- a/server/build-kubernetes.ts +++ b/server/build-kubernetes.ts @@ -539,6 +539,62 @@ export class KubernetesBuildStore implements BuildStore { .sort(compareBuildRecords); } + /** Prune old terminal history using fresh reads and Kubernetes RV preconditions. */ + async pruneBuilds(limit = 50): Promise { + if (!Number.isSafeInteger(limit) || limit < 0) return; + const records = await this.listBuilds(); + let remaining = records.length; + for (const candidate of records) { + if (remaining <= limit) break; + if (!terminal(candidate)) continue; + const current = await this.readConfigMap( + map(this.namespace, this.buildName(candidate.metadata.name), "build"), + ); + const record = current && payload(current); + if ( + !current?.metadata.resourceVersion || + !record || + record.kind !== "BuildRecord" || + record.metadata.name !== candidate.metadata.name || + current.metadata.resourceVersion !== + candidate.metadata.resourceVersion || + !terminal(record) + ) + continue; + const expiresAt = record.status.reconcileLease?.expiresAt; + if (expiresAt && Date.parse(expiresAt) > Date.now()) continue; + // A build ID cannot reacquire an image lock after becoming terminal, so + // an absent/mismatched lock cannot newly begin referencing this record. + // Treat unreadable or malformed locks as protected rather than guessing. + let lock: ConfigMap | undefined; + try { + lock = await this.readConfigMap( + map(this.namespace, this.lockName(record.spec.imageKey), "build-lock"), + ); + } catch { + continue; + } + if (lock) { + const lockRecord = payload<{ buildId?: string }>(lock); + if (!lockRecord?.buildId || lockRecord.buildId === record.metadata.name) + continue; + } + await this.delete( + map( + this.namespace, + current.metadata.name, + "build", + undefined, + current.metadata.resourceVersion, + ), + current.metadata.resourceVersion, + ); + // Conflicts are swallowed by delete; confirm absence before adjusting the + // retained count so concurrent updates never make us over-prune. + if (!(await this.readConfigMap(current))) remaining--; + } + } + async replaceBuild( record: BuildRecord, expectedResourceVersion: string, diff --git a/server/build-store.ts b/server/build-store.ts index 46c48c2..41835d6 100644 --- a/server/build-store.ts +++ b/server/build-store.ts @@ -64,6 +64,7 @@ export interface BuildStore { createBuild(record: BuildRecord): Promise; getBuild(id: string): Promise; listBuilds(): Promise; + pruneBuilds?(limit?: number): Promise; replaceBuild( record: BuildRecord, expectedResourceVersion: string, @@ -258,6 +259,43 @@ export class MemoryBuildStore implements BuildStore { return [...this.builds.values()].sort(compareBuildRecords).map(clone); } + async pruneBuilds(limit = 50): Promise { + if (!Number.isSafeInteger(limit) || limit < 0) return; + const records = await this.listBuilds(); + let remaining = records.length; + for (const record of records) { + if (remaining <= limit) break; + const current = this.builds.get(record.metadata.name); + if ( + !current || + current.metadata.resourceVersion !== record.metadata.resourceVersion + ) + continue; + if (!terminal(current)) continue; + if ( + current.status.reconcileLease && + Date.parse(current.status.reconcileLease.expiresAt) > Date.now() + ) + continue; + if (await this.ownsBuild(current.spec.imageKey, current.metadata.name)) + continue; + // ownsBuild yields; revalidate after it so a concurrent lease/status + // update cannot slip between the initial check and this synchronous delete. + const latest = this.builds.get(record.metadata.name); + if ( + !latest || + latest.metadata.resourceVersion !== record.metadata.resourceVersion || + !terminal(latest) || + (latest.status.reconcileLease && + Date.parse(latest.status.reconcileLease.expiresAt) > Date.now()) + ) + continue; + this.builds.delete(latest.metadata.name); + this.byImage.get(latest.spec.imageKey)?.delete(latest.metadata.name); + remaining--; + } + } + async replaceBuild( record: BuildRecord, expectedResourceVersion: string, diff --git a/server/kubernetes-logs.ts b/server/kubernetes-logs.ts index 9a469c4..d399208 100644 --- a/server/kubernetes-logs.ts +++ b/server/kubernetes-logs.ts @@ -2,10 +2,10 @@ import { AppsV1Api, CoreV1Api, KubeConfig, - Log, type V1LabelSelector, } from "@kubernetes/client-node"; -import { PassThrough } from "node:stream"; +import http, { type IncomingMessage, type RequestOptions } from "node:http"; +import https from "node:https"; import { logKubernetesRequest } from "../lib/request-log"; import { KubernetesLogError, @@ -36,22 +36,46 @@ function statusCode(error: unknown): number | undefined { function logError(error: unknown): KubernetesLogError { const status = statusCode(error); + const code = + error && typeof error === "object" && "code" in error + ? error.code + : undefined; return new KubernetesLogError( error instanceof Error ? error.message : "Kubernetes log request failed", - status === undefined || status === 408 || status === 429 || status >= 500, + code === undefined && + (status === undefined || status === 408 || status === 429 || status >= 500), ); } -export class KubernetesLogs implements KubernetesLogsBackend { - private readonly logger: Log; +function requestLogs( + url: URL, + options: RequestOptions, + signal: AbortSignal, +): Promise { + return new Promise((resolve, reject) => { + const transport = url.protocol === "https:" ? https : http; + const request = transport.request( + url, + { ...options, method: "GET", signal }, + resolve, + ); + request.once("error", reject); + request.end(); + }); +} +async function responseBody(response: IncomingMessage): Promise { + const chunks: Buffer[] = []; + for await (const chunk of response) chunks.push(Buffer.from(chunk)); + return Buffer.concat(chunks).toString("utf8"); +} + +export class KubernetesLogs implements KubernetesLogsBackend { constructor( - config: KubeConfig, + private readonly config: KubeConfig, private readonly apps = config.makeApiClient(AppsV1Api), private readonly core = config.makeApiClient(CoreV1Api), - ) { - this.logger = new Log(config); - } + ) {} async listDeployments( namespace: string, @@ -125,44 +149,57 @@ export class KubernetesLogs implements KubernetesLogsBackend { request: ContainerLogRequest, signal: AbortSignal, ): AsyncIterable { - const output = new PassThrough(); - let upstream: AbortController | undefined; - const abort = () => { - upstream?.abort(); - output.destroy(); - }; - signal.addEventListener("abort", abort, { once: true }); try { - upstream = await logKubernetesRequest( + if (signal.aborted) return; + const cluster = this.config.getCurrentCluster(); + if (!cluster) throw new Error("No currently active cluster"); + const url = new URL( + `${cluster.server}/api/v1/namespaces/${encodeURIComponent(request.namespace)}/pods/${encodeURIComponent(request.pod)}/log`, + ); + url.searchParams.set("container", request.container); + url.searchParams.set("follow", "true"); + url.searchParams.set("timestamps", String(request.timestamps)); + if (request.tailLines !== undefined) + url.searchParams.set("tailLines", String(request.tailLines)); + if (request.sinceSeconds !== undefined) + url.searchParams.set("sinceSeconds", String(request.sinceSeconds)); + if (request.sinceTime !== undefined) + url.searchParams.set("sinceTime", request.sinceTime); + + const options: RequestOptions = {}; + await this.config.applyToHTTPSOptions(options); + const response = await logKubernetesRequest( { method: "GET", - pathname: `/api/v1/namespaces/${request.namespace}/pods/${request.pod}/log`, + pathname: url.pathname, 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, - }, - ), + async () => { + const response = await requestLogs(url, options, signal); + const status = response.statusCode ?? 0; + if (status < 200 || status >= 300) { + const body = await responseBody(response); + let detail = body.trim(); + try { + const status = JSON.parse(body) as { message?: unknown }; + if (typeof status.message === "string") detail = status.message; + } catch { + // Proxies can return plain text rather than a Kubernetes Status. + } + throw new KubernetesLogError( + `Kubernetes log request failed (HTTP ${status})${detail ? `: ${detail.slice(0, 512)}` : ""}`, + status === 408 || status === 429 || status >= 500, + ); + } + return response; + }, ); - if (signal.aborted) abort(); - for await (const chunk of output) yield new Uint8Array(chunk); + for await (const chunk of response) + yield new Uint8Array(chunk); } catch (error) { - if (!signal.aborted) throw logError(error); - } finally { - signal.removeEventListener("abort", abort); - upstream?.abort(); - output.destroy(); + if (!signal.aborted) + throw error instanceof KubernetesLogError ? error : logError(error); } } } diff --git a/server/kubernetes-state.ts b/server/kubernetes-state.ts index 73b0f09..353e1df 100644 --- a/server/kubernetes-state.ts +++ b/server/kubernetes-state.ts @@ -15,7 +15,8 @@ import { 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 { waitForRollout } from "../lib/rollout"; +import { MAX_AUDIT_EVENTS, type AuditEvent, type AuditPersistence } from "./audit-store"; import { managementDependencies, type ManagementDependencies, @@ -942,28 +943,75 @@ export class KubernetesOperationPersistence implements OperationPersistence { export class KubernetesAuditPersistence implements AuditPersistence { constructor(private readonly objects = createKubernetesClients().objects) {} - async append(event: AuditEvent): Promise { - await createObject( - this.objects, - stateObject( - "ConfigMap", - event.metadata.name, - "audit", - event, - event.spec.workspaceId, - ), + private readonly name = "audit-log"; + + private async readRing(): Promise { + try { + return (await this.objects.read({ + apiVersion: "v1", kind: "ConfigMap", + metadata: { name: this.name, namespace: KUBER_STATE_NAMESPACE }, + })) as DataObject; + } catch (error) { + if (isNotFound(error)) return; + throw error; + } + } + + private async legacy(): Promise { + const result = await this.objects.list( + "v1", "ConfigMap", KUBER_STATE_NAMESPACE, undefined, undefined, + undefined, undefined, `${TYPE_LABEL}=audit`, ); + return result.items.map((item) => ({ + apiVersion: "v1", kind: "ConfigMap", ...item, + })) as DataObject[]; + } + + private async reconcile(add?: AuditEvent): Promise { + for (let attempt = 0; attempt < 128; attempt++) { + const [ring, legacy] = await Promise.all([this.readRing(), this.legacy()]); + const current = ring && parsePayload<{ kind: string; items: AuditEvent[] }>(ring); + const events = new Map(); + for (const event of [...(current?.items ?? []), ...legacy.map(parsePayload).filter((event): event is AuditEvent => event?.kind === "AuditEvent"), ...(add ? [add] : [])]) { + if (event?.kind === "AuditEvent") events.set(event.metadata.uid || event.metadata.name, event); + } + const ordered = [...events.values()].sort((a, b) => + a.metadata.creationTimestamp.localeCompare(b.metadata.creationTimestamp) || + a.metadata.name.localeCompare(b.metadata.name), + ); + let retained = ordered.slice(-MAX_AUDIT_EVENTS); + while (retained.length && Buffer.byteLength(JSON.stringify({ kind: "AuditLog", items: retained })) > 768 * 1024) retained.shift(); + const value = stateObject("ConfigMap", this.name, "audit-log", { kind: "AuditLog", items: retained }); + try { + if (ring) { + value.metadata!.resourceVersion = ring.metadata?.resourceVersion; + await this.objects.replace(value); + } else { + await this.objects.create(value); + } + } catch (error) { + if (isConflict(error)) continue; + throw error; + } + // Never remove legacy records until the merged ring has been committed. + for (const item of legacy.slice(0, 100)) { + const name = item.metadata?.name; + if (name) await deleteObject(this.objects, "ConfigMap", name); + } + return retained; + } + throw new Error("Audit log was concurrently modified; retry limit reached"); + } + + async append(event: AuditEvent): Promise { + await this.reconcile(event); } async list(workspaceId?: string): Promise { - return (await list(this.objects, "ConfigMap", "audit", workspaceId)) - .map((item) => parsePayload(item)) - .filter((item): item is AuditEvent => item?.kind === "AuditEvent") - .sort((a, b) => - a.metadata.creationTimestamp.localeCompare( - b.metadata.creationTimestamp, - ), - ); + const events = await this.reconcile(); + return events.filter((event) => + workspaceId ? event.spec.workspaceId === workspaceId : true, + ); } } @@ -1032,34 +1080,10 @@ function resource(identity: ResourceIdentity): KubernetesObject { }; } -async function sleepUntilExecutionCancelled( - delayMs: number, - execution?: OperationExecution, -): Promise { - const signal = execution?.signal; - if (!signal) { - await Bun.sleep(delayMs); - return; - } - if (signal.aborted) - throw new Error("Workspace operation execution was cancelled"); - await new Promise((resolve, reject) => { - const timer = setTimeout(() => { - signal.removeEventListener("abort", cancelSleep); - resolve(); - }, delayMs); - const cancelSleep = () => { - clearTimeout(timer); - reject(new Error("Workspace operation execution was cancelled")); - }; - signal.addEventListener("abort", cancelSleep, { once: true }); - }); -} - export function createKubernetesManagementDependencies( clients = createKubernetesClients(), ): ManagementDependencies { - const { objects, apps } = clients; + const { objects, apps, core } = clients; const revision = (replicaSet: V1ReplicaSet): number | undefined => { const value = Number( replicaSet.metadata?.annotations?.["deployment.kubernetes.io/revision"], @@ -1171,30 +1195,17 @@ export function createKubernetesManagementDependencies( name, timeoutMs = 300_000, execution?: OperationExecution, - ) => { - const started = Date.now(); - while (Date.now() - started < timeoutMs) { - if (execution?.signal?.aborted) - throw new Error("Workspace operation execution was cancelled"); - const deployment = await apps.readNamespacedDeployment({ - namespace: project, - name, - }); - const desired = deployment.spec?.replicas ?? 1; - if ( - (deployment.status?.observedGeneration ?? 0) >= - (deployment.metadata?.generation ?? 0) && - (deployment.status?.updatedReplicas ?? 0) === desired && - (deployment.status?.availableReplicas ?? 0) === desired && - (deployment.status?.unavailableReplicas ?? 0) === 0 - ) - return; - if (execution?.signal?.aborted) - throw new Error("Workspace operation execution was cancelled"); - await sleepUntilExecutionCancelled(2_000, execution); - } - throw new Error(`Timed out waiting for deployment ${name} rollout`); - }, + ) => + waitForRollout( + name, + Date.now() + timeoutMs, + () => apps.readNamespacedDeployment({ namespace: project, name }), + async () => + (await apps.listNamespacedReplicaSet({ namespace: project })).items, + async () => + (await core.listNamespacedPod({ namespace: project })).items, + execution?.signal, + ), planRollback: async (project, names) => { const deployments = await overrides.listDeployments!(project); const selected = names diff --git a/server/log-service.ts b/server/log-service.ts index 9c27748..8ee1b32 100644 --- a/server/log-service.ts +++ b/server/log-service.ts @@ -222,6 +222,13 @@ function isAbort(signal: AbortSignal): boolean { } function isRetryable(error: unknown): boolean { + if ( + error && + typeof error === "object" && + "code" in error && + typeof error.code === "string" + ) + return false; return !(error instanceof KubernetesLogError) || error.retryable; } diff --git a/server/maintenance.ts b/server/maintenance.ts index 860a445..1b18d57 100644 --- a/server/maintenance.ts +++ b/server/maintenance.ts @@ -31,12 +31,11 @@ export interface MaintenanceLeaseProvider { export function normalizeMaintenanceHost(value: unknown): string { if (typeof value !== "string") throw new Error("host is required"); const host = value.trim().toLowerCase().replace(/\.$/, ""); - console.log(host) if ( host.length === 0 || host.length > 253 || !host.includes(".") || - !/^(?:\*\.)?(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z]{2,63}$/.test(host) + !/^(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z]{2,63}$/.test(host) ) throw new Error("host must be a DNS hostname"); return host; diff --git a/server/management.ts b/server/management.ts index 43df58e..623f0cb 100644 --- a/server/management.ts +++ b/server/management.ts @@ -796,33 +796,46 @@ export function createManagementService(dependencies: ManagementDependencies) { execution?: OperationExecution, ): Promise { const selected = await targets(workspace, deploymentTargets); - for (const name of selected) { - const event = { - resource: { - apiVersion: "apps/v1", - kind: "Deployment", - name, - namespace: workspace.project, - }, - phase: "wait" as const, - }; - try { - throwIfExecutionAborted(execution); - await execution?.emit?.({ ...event, state: "started" }); - await dependencies.waitForDeployment( - workspace.project, - name, - timeoutMs, - execution, - ); - await execution?.emit?.({ ...event, state: "succeeded" }); - } catch (error) { - await execution?.emit?.({ - ...event, - state: execution?.signal?.aborted ? "aborted" : "failed", - }); - throw error; - } + const deadline = Date.now() + (timeoutMs ?? 300_000); + const controller = new AbortController(); + const cancel = () => controller.abort(); + execution?.signal?.addEventListener("abort", cancel, { once: true }); + if (execution?.signal?.aborted) cancel(); + try { + await Promise.all( + selected.map(async (name) => { + const event = { + resource: { + apiVersion: "apps/v1", + kind: "Deployment", + name, + namespace: workspace.project, + }, + phase: "wait" as const, + }; + try { + throwIfExecutionAborted({ signal: controller.signal }); + await execution?.emit?.({ ...event, state: "started" }); + await dependencies.waitForDeployment( + workspace.project, + name, + Math.max(0, deadline - Date.now()), + { ...execution, signal: controller.signal }, + ); + await execution?.emit?.({ ...event, state: "succeeded" }); + } catch (error) { + await execution?.emit?.({ + ...event, + state: controller.signal.aborted ? "aborted" : "failed", + }); + controller.abort(); + throw error; + } + }), + ); + } finally { + controller.abort(); + execution?.signal?.removeEventListener("abort", cancel); } }, diff --git a/tests/command/logs.test.ts b/tests/command/logs.test.ts index d99e61a..bc6f290 100644 --- a/tests/command/logs.test.ts +++ b/tests/command/logs.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test } from "bun:test"; import { + parseLogTailLines, runLogs, type LogApiEvent, type LogEventWriter, @@ -81,4 +82,34 @@ describe("logs API command runner", () => { ); expect(receivedSignal).toBe(controller.signal); }); + + test.each([false, true])("sends -n with follow=%s", async (follow) => { + const paths: string[] = []; + const stream: LogsApiStream = async function* (path, _init, options) { + paths.push(path); + expect(options?.timeoutMs).toBe(follow ? 0 : undefined); + yield* [] as never[]; + }; + await runLogs("shop", "web", follow, undefined, stream, () => {}, 0); + expect(paths).toEqual([ + `/workspaces/shop/logs?service=web${follow ? "&follow=true" : ""}&tailLines=0`, + ]); + }); + + test("accepts only non-negative safe integer line counts", async () => { + expect(parseLogTailLines("25")).toBe(25); + expect(parseLogTailLines("0")).toBe(0); + for (const invalid of ["-1", "1.5", "2abc", "", "1e3", "9007199254740992"]) + expect(() => parseLogTailLines(invalid)).toThrow("non-negative integer"); + + let called = false; + const stream: LogsApiStream = async function* () { + called = true; + yield* [] as never[]; + }; + await expect( + runLogs("shop", undefined, false, undefined, stream, () => {}, -1), + ).rejects.toThrow("non-negative integer"); + expect(called).toBe(false); + }); }); diff --git a/tests/command/module-graph.test.ts b/tests/command/module-graph.test.ts new file mode 100644 index 0000000..3d22b78 --- /dev/null +++ b/tests/command/module-graph.test.ts @@ -0,0 +1,14 @@ +import { describe, expect, test } from "bun:test"; +import { join } from "node:path"; + +describe("CLI module graph", () => { + test("does not resolve the Kubernetes client during CLI startup", () => { + const entry = join(import.meta.dir, "../../index.ts"); + const probe = Bun.spawnSync([process.execPath, entry, "--help"], { + cwd: import.meta.dir, + env: process.env, + }); + expect(probe.exitCode).toBe(0); + expect(new TextDecoder().decode(probe.stdout)).toContain("USAGE"); + }); +}); diff --git a/tests/command/startup.test.ts b/tests/command/startup.test.ts new file mode 100644 index 0000000..0aac6ca --- /dev/null +++ b/tests/command/startup.test.ts @@ -0,0 +1,81 @@ +import { describe, expect, test } from "bun:test"; +import { join } from "node:path"; + +const entry = join(import.meta.dir, "../../index.ts"); + +function cli(...args: string[]) { + const result = Bun.spawnSync([process.execPath, entry, ...args], { + cwd: import.meta.dir, + env: { ...process.env, TERM: "dumb", NO_COLOR: "1" }, + }); + return { + status: result.exitCode, + output: new TextDecoder().decode(result.stdout), + error: new TextDecoder().decode(result.stderr), + }; +} + +describe("CLI startup", () => { + test("registers all original names as lazy commands", () => { + const probe = Bun.spawnSync([ + process.execPath, + "-e", + `import { main } from ${JSON.stringify(join(import.meta.dir, "../../command/main.ts"))}; console.log(JSON.stringify(Object.entries(main.subCommands).map(([name, command]) => [name, typeof command])));`, + ]); + expect(probe.exitCode).toBe(0); + const commands = JSON.parse(new TextDecoder().decode(probe.stdout)) as [ + string, + string, + ][]; + expect(commands.map(([name]) => name)).toEqual([ + "audit", + "ci", + "db", + "down", + "export", + "exec", + "logs", + "maintenance", + "login", + "logout", + "operations", + "ps", + "restart", + "rollback", + "s3", + "start", + "stop", + "up", + "users", + "trust", + "whoami", + ]); + expect(commands.every(([, type]) => type === "function")).toBe(true); + }); + + test("root help retains descriptions, alias, and completion command", () => { + const result = cli("--help"); + expect(result.status).toBe(0); + expect(result.output).toContain("rollback, fuck"); + expect(result.output).toContain("complete"); + expect(result.output).toContain("Create and start deployments"); + + const missingCommand = cli(); + expect(missingCommand.status).toBe(1); + expect(missingCommand.output).toContain("complete"); + }); + + test("alias and named subcommand have the same help", () => { + const alias = cli("fuck", "--help"); + const named = cli("rollback", "--help"); + expect(alias.status).toBe(0); + expect(alias.output).toBe(named.output); + expect(alias.output).toContain("DEPLOYMENT"); + }); + + test("completion still resolves flags from nested commands", () => { + const result = cli("complete", "--", "logs", "--"); + expect(result.status).toBe(0); + expect(result.output).toContain("--follow"); + }); +}); diff --git a/tests/command/up-api.test.ts b/tests/command/up-api.test.ts index 8d17634..b6dbb75 100644 --- a/tests/command/up-api.test.ts +++ b/tests/command/up-api.test.ts @@ -45,9 +45,6 @@ describe("up API pipeline", () => { join(root, ".kuberrc.ts"), 'export default { project: "shop" };\n', ); - const git = Bun.spawn(["git", "init", "-q", root]); - expect(await git.exited).toBe(0); - process.chdir(root); process.env.XDG_CONFIG_HOME = configHome; const identity = await resolveTrustIdentity("shop", root); @@ -106,6 +103,92 @@ describe("up API pipeline", () => { } }); + test("adopts before overlapping database and storage reconciliation, then renders merged env", async () => { + const root = await mkdtemp(join(tmpdir(), "kuber-up-api-")); + const previousCwd = process.cwd(); + const order: string[] = []; + const started: string[] = []; + let bothReady!: () => void; + const bothStarted = new Promise((resolve) => { + bothReady = resolve; + }); + let releaseDatabase!: () => void; + let releaseStorage!: () => void; + const database = new Promise((resolve) => { + releaseDatabase = resolve; + }); + const storage = new Promise((resolve) => { + releaseStorage = resolve; + }); + try { + await writeFile( + join(root, "compose.yml"), + "services:\n app:\n image: nginx\n volumes:\n - postgresql:app\n - s3:assets\n", + ); + await writeFile( + join(root, ".kuberrc.ts"), + 'export default { project: "shop" };\n', + ); + expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0); + process.chdir(root); + const trust = await resolveTrustIdentity("shop", root); + const request: ApiRequester = async (path: string) => { + order.push(path); + if (path === "/workspaces/shop") throw new KuberApiError("missing", 404); + if (path === "/workspaces") + return { + metadata: { name: "shop", uid: "workspace", resourceVersion: "1" }, + } as T; + if (path.endsWith("/adopt")) return { resourcesAdopted: 0 } as T; + if (path.endsWith("/databases")) { + started.push("database"); + if (started.length === 2) bothReady(); + await database; + return { app: { DATABASE_URL: "postgresql://app", SHARED: "database" } } as T; + } + if (path.endsWith("/storage")) { + started.push("storage"); + if (started.length === 2) bothReady(); + await storage; + return { app: { AWS_ACCESS_KEY_ID: "key", SHARED: "storage" } } as T; + } + if (path.endsWith("/plan")) return { desired: [], stale: [] } as T; + return {} as T; + }; + const run = provideContext(() => runUp(false, request, { trust })); + try { + await bothStarted; + expect(started).toEqual(["database", "storage"]); + expect(order.indexOf("/workspaces/shop/adopt")).toBeLessThan( + order.indexOf("/workspaces/shop/databases"), + ); + expect(order).not.toContain("/workspaces/shop/resources/plan"); + releaseStorage(); + await Promise.resolve(); + expect(order).not.toContain("/workspaces/shop/resources/plan"); + releaseDatabase(); + const result = await run; + expect(result.serviceEnv).toEqual({ + app: { + DATABASE_URL: "postgresql://app", + SHARED: "storage", + AWS_ACCESS_KEY_ID: "key", + }, + }); + expect(order.indexOf("/workspaces/shop/resources/plan")).toBeGreaterThan( + order.indexOf("/workspaces/shop/storage"), + ); + } finally { + releaseStorage(); + releaseDatabase(); + await run.catch(() => {}); + } + } finally { + process.chdir(previousCwd); + await rm(root, { recursive: true, force: true }); + } + }); + test("creates workspace metadata without embedding source blobs", async () => { const calls: Array<{ path: string; init?: ApiRequestInit }> = []; const request: ApiRequester = async ( diff --git a/tests/lib/api.test.ts b/tests/lib/api.test.ts index 201041d..2b5f8e8 100644 --- a/tests/lib/api.test.ts +++ b/tests/lib/api.test.ts @@ -267,4 +267,19 @@ describe("authenticated API transport", () => { }; await expect(consume()).rejects.toBeInstanceOf(SyntaxError); }); + + test("parses many records packed into one NDJSON chunk", async () => { + const expected = Array.from({ length: 2_000 }, (_, sequence) => ({ sequence })); + globalThis.fetch = mock(async () => + new Response(`${expected.map((record) => JSON.stringify(record)).join("\n")}\n`), + ) as unknown as typeof fetch; + + const actual: Array<{ sequence: number }> = []; + for await (const record of apiStreamNdjson<{ sequence: number }>( + "/events", + {}, + { authenticated: false, baseUrl: "https://api.test" }, + )) actual.push(record); + expect(actual).toEqual(expected); + }); }); diff --git a/tests/lib/build-api.test.ts b/tests/lib/build-api.test.ts index d491b42..f6cfaa4 100644 --- a/tests/lib/build-api.test.ts +++ b/tests/lib/build-api.test.ts @@ -2,6 +2,7 @@ import { afterEach, describe, expect, test } from "bun:test"; import { mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; +import { Writable } from "node:stream"; import type { ApiRequestInit, ApiRequestOptions } from "../../lib/api"; import { buildServices, @@ -563,7 +564,7 @@ describe("authenticated build API pipeline", () => { web: `registry.server/kuber/shop-web@sha256:${"a".repeat(64)}`, }, }); - expect(output).toContain("build log"); + expect(output).toContain("web: build log"); expect(calls.find(({ path }) => path === "/builds")?.options).toEqual({ timeoutMs: 300_000, }); @@ -581,6 +582,218 @@ describe("authenticated build API pipeline", () => { expect(calls.some(({ path }) => path.endsWith("/result"))).toBe(true); }); + test("bounds overlapping builds, attributes logs, and returns images in compose order", async () => { + const root = await mkdtemp(join(tmpdir(), "kuber-build-api-")); + directories.push(root); + expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0); + const snapshot = emptySnapshot(); + const started: string[] = []; + const startWaiters = new Map void>(); + const waitForStart = (count: number) => + new Promise((resolve) => { + startWaiters.set(count, resolve); + }); + const firstTwo = waitForStart(2); + const third = waitForStart(3); + const fourth = waitForStart(4); + const released = new Map void>(); + const gates = new Map>(); + const ids = new Map(); + const output: string[] = []; + const stream = new Writable({ + write(chunk, _encoding, done) { + output.push(String(chunk)); + done(); + }, + }); + const request: ApiRequester = async ( + path: string, + init?: ApiRequestInit, + ) => { + if (path === "/snapshots/negotiate") + return { workspace: snapshot.digest, missing: [], ready: true } as T; + if (path === "/builds") { + const build = init?.json as BuildRequest; + started.push(build.service); + startWaiters.get(started.length)?.(); + ids.set(build.id, build.service); + gates.set( + build.service, + new Promise((resolve) => released.set(build.service, resolve)), + ); + return { id: build.id, state: "queued" } as T; + } + const id = path.split("/")[2]!; + const service = ids.get(id)!; + if (path.includes("/events")) + return [{ type: "log", sequence: 1, message: `log ${service}\n` }] as T; + if (path.endsWith("/reconcile")) { + await gates.get(service); + return { id, state: "succeeded" } as T; + } + if (path.endsWith("/result")) return { reference: `image:${service}` } as T; + throw new Error(`Unexpected request ${path}`); + }; + const run = buildServices( + "shop", + { + services: Object.fromEntries( + ["one", "two", "three", "four"].map((name) => [ + name, + { build: "." }, + ]), + ), + }, + root, + { stream }, + { request, snapshot, buildConcurrency: 2, sleep: async () => {} }, + ); + try { + await firstTwo; + expect(started).toEqual(["one", "two"]); + released.get("two")!(); + await third; + expect(started).toEqual(["one", "two", "three"]); + released.get("three")!(); + await fourth; + released.get("four")!(); + released.get("one")!(); + expect(await run).toEqual({ + built: ["one", "two", "three", "four"], + changed: ["one", "two", "three", "four"], + images: { + one: "image:one", + two: "image:two", + three: "image:three", + four: "image:four", + }, + }); + expect(new Set(ids.keys()).size).toBe(4); + expect(output).toEqual( + expect.arrayContaining([ + "[one] log one\n", + "[two] log two\n", + "[three] log three\n", + "[four] log four\n", + ]), + ); + } finally { + for (const release of released.values()) release(); + await run.catch(() => {}); + } + }); + + test("drains active builds and does not start queued services after failure", async () => { + const root = await mkdtemp(join(tmpdir(), "kuber-build-api-")); + directories.push(root); + expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0); + const snapshot = emptySnapshot(); + const started: string[] = []; + let bothReady!: () => void; + const bothStarted = new Promise((resolve) => { + bothReady = resolve; + }); + let fail!: (error: Error) => void; + let finish!: () => void; + const failing = new Promise((_resolve, reject) => { + fail = reject; + }); + const active = new Promise((resolve) => { + finish = resolve; + }); + const request: ApiRequester = async ( + path: string, + init?: ApiRequestInit, + ) => { + if (path === "/snapshots/negotiate") return { ready: true } as T; + if (path === "/builds") { + const service = (init!.json as BuildRequest).service; + started.push(service); + if (started.length === 2) bothReady(); + if (service === "one") return await failing; + await active; + return { state: "succeeded" } as T; + } + if (path.includes("/events")) return [] as T; + if (path.endsWith("/result")) return { reference: "image:two" } as T; + throw new Error(`Unexpected request ${path}`); + }; + const run = buildServices( + "shop", + { + services: { + one: { build: "." }, + two: { build: "." }, + three: { build: "." }, + }, + }, + root, + undefined, + { request, snapshot, buildConcurrency: 2 }, + ); + const outcome = run.then(() => undefined, (error: unknown) => error); + await bothStarted; + fail(new Error("first build failed")); + await Promise.resolve(); + expect(started).toEqual(["one", "two"]); + finish(); + expect(await outcome).toMatchObject({ message: "first build failed" }); + expect(started).toEqual(["one", "two"]); + }); + + test("passes cancellation to active build requests and skips queued builds", async () => { + const root = await mkdtemp(join(tmpdir(), "kuber-build-api-")); + directories.push(root); + expect(await Bun.spawn(["git", "init", "-q", root]).exited).toBe(0); + const snapshot = emptySnapshot(); + const controller = new AbortController(); + const cancelled = new DOMException("Cancelled", "AbortError"); + const started: string[] = []; + let ready!: () => void; + const bothStarted = new Promise((resolve) => { + ready = resolve; + }); + const request: ApiRequester = async ( + path: string, + init?: ApiRequestInit, + ) => { + if (path === "/snapshots/negotiate") return { ready: true } as T; + if (path === "/builds") { + started.push((init!.json as BuildRequest).service); + expect(init?.signal).toBe(controller.signal); + if (started.length === 2) ready(); + return await new Promise((_resolve, reject) => { + init?.signal?.addEventListener( + "abort", + () => reject(init.signal!.reason), + { once: true }, + ); + }); + } + throw new Error(`Unexpected request ${path}`); + }; + const outcome = buildServices( + "shop", + { + services: { + one: { build: "." }, + two: { build: "." }, + three: { build: "." }, + }, + }, + root, + undefined, + { request, snapshot, buildConcurrency: 2, signal: controller.signal }, + ).then( + () => undefined, + (error: unknown) => error, + ); + await bothStarted; + controller.abort(cancelled); + expect(await outcome).toBe(cancelled); + expect(started).toEqual(["one", "two"]); + }); + test("retries a timed out poll request with the build request timeout", async () => { const root = await mkdtemp(join(tmpdir(), "kuber-build-api-")); directories.push(root); @@ -668,4 +881,57 @@ describe("authenticated build API pipeline", () => { 'POST:/images/resolve:{"project":"shop","service":"web"}', ]); }); + + test("resolves build images concurrently with bounded concurrency and stable order", async () => { + const services = ["one", "two", "three", "four", "five"]; + const started: string[] = []; + const pending: Array<() => void> = []; + let active = 0; + let maximumActive = 0; + let releaseFirstWave!: () => void; + const firstWave = new Promise((resolve) => { + releaseFirstWave = resolve; + }); + const run = resolveBuildImages( + "shop", + { + services: Object.fromEntries([ + ...services.map((service) => [service, { build: "." }]), + ["cache", { image: "redis" }], + ]), + }, + { + request: async (_path: string, init?: ApiRequestInit) => { + const service = ((init?.json ?? {}) as { service: string }).service; + started.push(service); + active++; + maximumActive = Math.max(maximumActive, active); + if (started.length === 3) releaseFirstWave(); + await new Promise((resolve) => pending.push(resolve)); + active--; + return { reference: `image:${service}` } as T; + }, + }, + ); + + try { + await firstWave; + expect(started).toEqual(services.slice(0, 3)); + expect(maximumActive).toBe(3); + for (let i = 0; i < services.length; i++) { + const release = pending.shift(); + if (!release) throw new Error("Expected a pending image request"); + release(); + await Promise.resolve(); + await Promise.resolve(); + } + const images = await run; + expect(Object.keys(images)).toEqual(services); + expect(images).toEqual(Object.fromEntries(services.map((service) => [service, `image:${service}`]))); + expect(started).toEqual(services); + } finally { + for (const release of pending) release(); + await run.catch(() => {}); + } + }); }); diff --git a/tests/lib/rollout.test.ts b/tests/lib/rollout.test.ts new file mode 100644 index 0000000..45bde66 --- /dev/null +++ b/tests/lib/rollout.test.ts @@ -0,0 +1,151 @@ +import { describe, expect, test } from "bun:test"; +import type { + V1Deployment, + V1Pod, + V1ReplicaSet, +} from "@kubernetes/client-node"; +import { + currentReplicaSet, + rolloutPodFailure, + waitForRollout, +} from "../../lib/rollout"; + +const deployment = { + metadata: { name: "web", uid: "deployment-1", generation: 2 }, + spec: { + replicas: 1, + template: { + metadata: { labels: { app: "web" } }, + spec: { containers: [{ name: "app", image: "new:v2" }] }, + }, + }, + status: { + observedGeneration: 2, + updatedReplicas: 1, + availableReplicas: 0, + unavailableReplicas: 1, + }, +} as unknown as V1Deployment; +const set = (uid: string, image: string) => + ({ + metadata: { + uid, + ownerReferences: [{ kind: "Deployment", uid: "deployment-1" }], + }, + spec: { + template: { + metadata: { labels: { app: "web", "pod-template-hash": uid } }, + spec: { containers: [{ name: "app", image }] }, + }, + }, + }) as unknown as V1ReplicaSet; +const pod = (uid: string, status: object, init = false) => + ({ + metadata: { + name: `${uid}-pod`, + ownerReferences: [{ kind: "ReplicaSet", uid }], + }, + status: { + [init ? "initContainerStatuses" : "containerStatuses"]: [ + { name: "app", restartCount: 3, ...status }, + ], + }, + }) as unknown as V1Pod; + +describe("rollout failure diagnosis", () => { + test("ignores crashing pods from the previous ReplicaSet", () => { + const current = currentReplicaSet(deployment, [ + set("old", "old:v1"), + set("new", "new:v2"), + ]); + expect(current?.metadata?.uid).toBe("new"); + expect( + rolloutPodFailure(deployment, current!, [ + pod("old", { state: { waiting: { reason: "CrashLoopBackOff" } } }), + ]), + ).toBeUndefined(); + expect( + currentReplicaSet(deployment, [set("old", "old:v1")]), + ).toBeUndefined(); + }); + + test("reports current init and application container failures with safe identifiers", () => { + const current = set("new", "new:v2"); + expect( + rolloutPodFailure(deployment, current, [ + pod( + "new", + { state: { waiting: { reason: "ImagePullBackOff" } } }, + true, + ), + ]), + ).toContain("pod new-pod, container app: ImagePullBackOff"); + expect( + rolloutPodFailure(deployment, current, [ + pod("new", { + lastState: { terminated: { exitCode: 137, reason: "OOMKilled" } }, + }), + ]), + ).toContain("OOMKilled (exit 137, restarts 3)"); + expect( + rolloutPodFailure(deployment, current, [ + pod("new", { + lastState: { terminated: { exitCode: 1, reason: "Error" } }, + }), + ]), + ).toContain("Error (exit 1, restarts 3)"); + }); + + test("detects deadline condition without waiting for pod listing", async () => { + const failed = { + ...deployment, + status: { + ...deployment.status, + conditions: [ + { + type: "Progressing", + status: "False", + reason: "ProgressDeadlineExceeded", + }, + ], + }, + } as V1Deployment; + await expect( + waitForRollout( + "web", + Date.now() + 100, + async () => failed, + async () => { + throw Error("should not list"); + }, + async () => [], + ), + ).rejects.toThrow( + "Deployment web rollout failed: ProgressDeadlineExceeded", + ); + }); + + test("bounds hung reads by deadline and responds to cancellation", async () => { + const never = async (): Promise => new Promise(() => {}); + await expect( + waitForRollout( + "web", + Date.now() + 20, + never, + async () => [], + async () => [], + ), + ).rejects.toThrow("Timed out waiting for deployment web rollout"); + const controller = new AbortController(); + const waiting = waitForRollout( + "web", + Date.now() + 10_000, + never, + async () => [], + async () => [], + controller.signal, + ); + controller.abort(); + await expect(waiting).rejects.toThrow("cancelled"); + }); +}); diff --git a/tests/lib/session.test.ts b/tests/lib/session.test.ts index c70f89c..9b5f7f0 100644 --- a/tests/lib/session.test.ts +++ b/tests/lib/session.test.ts @@ -68,4 +68,19 @@ describe("API sessions", () => { ); expect(await readSession()).toBeUndefined(); }); + + test("invalidates a memoized missing session after writing", async () => { + await setupDirectories(); + expect(await readSession()).toBeUndefined(); + await writeSession(session, false); + expect(await readSession()).toEqual(session); + }); + + test("invalidates a memoized session after removal", async () => { + await setupDirectories(); + await writeSession(session, false); + expect(await readSession()).toEqual(session); + await removeSessions(); + expect(await readSession()).toBeUndefined(); + }); }); diff --git a/tests/lib/workspace.test.ts b/tests/lib/workspace.test.ts index 03f0de5..f7ba9a7 100644 --- a/tests/lib/workspace.test.ts +++ b/tests/lib/workspace.test.ts @@ -2,6 +2,7 @@ import { afterEach, describe, expect, test } from "bun:test"; import { execFile } from "node:child_process"; import { mkdtemp, + mkdir, readFile, readlink, rm, @@ -50,6 +51,80 @@ afterEach(async () => { }); describe("workspace snapshots", () => { + test("snapshots Gitless directories deterministically without secrets", async () => { + const root = await temporaryDirectory("kuber-workspace-filesystem-"); + await writeFile(join(root, "Dockerfile"), "FROM scratch\n"); + await writeFile(join(root, "app.txt"), "application"); + await writeFile(join(root, ".env.local"), "SECRET=hidden"); + await writeFile(join(root, "db_credentials.json"), "hidden"); + + const first = await enumerateWorkspace(root); + const second = await enumerateWorkspace(root); + expect(first).toEqual(second); + expect(first.manifest.files.map((file) => file.path)).toEqual([ + "Dockerfile", + "app.txt", + ]); + }); + + test("works without Git and conservatively prunes ignored/generated and credential files", async () => { + const root = await temporaryDirectory("kuber-workspace-no-git-"); + await writeFile(join(root, ".gitignore"), "local-only/\n*.generated\nsecrets/\n!secrets/keep.txt\n"); + await writeFile(join(root, "app.ts"), "source"); + await mkdir(join(root, "local-only")); + await writeFile(join(root, "local-only", "hidden"), "secret"); + await writeFile(join(root, "cache.generated"), "generated"); + await writeFile(join(root, "credentials.json"), "credential"); + await writeFile(join(root, ".env.local"), "SECRET=hidden"); + await mkdir(join(root, "secrets")); + await writeFile(join(root, "secrets/keep.txt"), "must remain excluded"); + await mkdir(join(root, "services/secrets"), { recursive: true }); + await writeFile(join(root, "services/secrets/production.yaml"), "secret"); + await mkdir(join(root, "app_credentials")); + await writeFile(join(root, "app_credentials/key.json"), "credential"); + await mkdir(join(root, "config")); + await writeFile(join(root, "config/private.yaml"), "config secret"); + await mkdir(join(root, "secretary")); + await writeFile(join(root, "secretary/notes.txt"), "ordinary directory"); + await mkdir(join(root, ".next")); + await writeFile(join(root, ".next/generated"), "generated"); + const previousPath = process.env.PATH; + try { + process.env.PATH = ""; + const snapshot = await enumerateWorkspace(root); + expect(snapshot.manifest.files.map((file) => file.path)).toEqual([ + ".gitignore", + "app.ts", + "secretary/notes.txt", + ]); + } finally { + if (previousPath === undefined) delete process.env.PATH; + else process.env.PATH = previousPath; + } + }); + + test("uses conservative filesystem mode for an initialized repo when Git is absent from PATH", async () => { + const root = await repository(); + await writeFile(join(root, ".gitignore"), "secrets/\n!secrets/keep.txt\n"); + await writeFile(join(root, "source.ts"), "source"); + await writeFile(join(root, ".env.local"), "SECRET=hidden"); + await mkdir(join(root, "secrets")); + await writeFile(join(root, "secrets/keep.txt"), "hidden"); + const previousPath = process.env.PATH; + try { + process.env.PATH = ""; + const first = await enumerateWorkspace(root); + const second = await enumerateWorkspace(root); + expect(first).toEqual(second); + expect(first.manifest.files.map((file) => file.path)).toEqual([ + ".gitignore", "source.ts", + ]); + } finally { + if (previousPath === undefined) delete process.env.PATH; + else process.env.PATH = previousPath; + } + }); + test("captures working tracked, untracked, and ignored dotenv files deterministically", async () => { const root = await repository(); await writeFile(join(root, ".gitignore"), "ignored*\n.env*\nsub/.env*\n"); @@ -139,6 +214,30 @@ describe("workspace snapshots", () => { ); }); + test("allows ignored special files and prunes ignored directories", async () => { + const root = await repository(); + await writeFile(join(root, ".gitignore"), "ignored-pipe\nignored-dir/\n"); + await run("mkfifo", [join(root, "ignored-pipe")]); + await run("mkdir", [join(root, "ignored-dir")]); + await run("mkfifo", [join(root, "ignored-dir/pipe")]); + + await expect(enumerateWorkspace(root)).resolves.toMatchObject({ + manifest: { files: [{ path: ".gitignore" }] }, + }); + }); + + test("protects dotenv special files even when ignored", async () => { + const root = await repository(); + await writeFile(join(root, ".gitignore"), "*.pipe\n!keep.pipe\nsub/\n"); + await run("mkfifo", [join(root, "blocked.pipe")]); + await expect(enumerateWorkspace(root)).resolves.toBeDefined(); + + await run("mkfifo", [join(root, ".env.pipe")]); + await expect(enumerateWorkspace(root)).rejects.toThrow( + "Special files are not allowed in workspaces: .env.pipe", + ); + }); + test("rejects traversal, unsorted manifests, ancestor collisions, and corrupt blobs", async () => { expect(() => validateWorkspacePath("../secret")).toThrow( "Unsafe workspace path", diff --git a/tests/server/app.test.ts b/tests/server/app.test.ts index c164d09..976beb8 100644 --- a/tests/server/app.test.ts +++ b/tests/server/app.test.ts @@ -185,6 +185,58 @@ describe("kuber API authentication", () => { }); describe("operation response safety", () => { + test("database and storage reconcile concurrently while conflicting scopes fail fast", async () => { + const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid" }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const operationStore = new MemoryOperationStore(); + const leases = new MemoryWorkspaceLeaseProvider(); + const entered = new Set(); + let release!: () => void; + const blocked = new Promise((resolve) => (release = resolve)); + let bothEntered!: () => void; + const bothStarted = new Promise((resolve) => (bothEntered = resolve)); + const reconcile = (name: string) => async () => { + entered.add(name); + if (entered.size === 2) bothEntered(); + await blocked; + return {}; + }; + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + operationStore, + leases, + management: { + reconcileDatabases: reconcile("database"), + reconcileStorage: reconcile("storage"), + stop: async () => [], + } as unknown as ManagementService, + }); + const post = (path: string, key: string, body: unknown = { compose: { services: {} } }) => + app(request(`/api/v2/workspaces/demo/${path}`, { + method: "POST", + headers: { "idempotency-key": key }, + body: JSON.stringify(body), + }, "token")); + + const database = post("databases", "database-one"); + const storage = post("storage", "storage-one"); + await bothStarted; + expect(entered).toEqual(new Set(["database", "storage"])); + + const sameScope = await post("databases", "database-two"); + expect(sameScope.status).toBe(409); + const broad = await post("lifecycle", "lifecycle", { action: "stop" }); + expect(broad.status).toBe(409); + + release(); + expect((await database).status).toBe(200); + expect((await storage).status).toBe(200); + }); + test("redacts database and storage reconciliation results immediately", async () => { const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid", @@ -1025,6 +1077,73 @@ describe("kuber v2 HTTP routes", () => { ); }); + test("a stalled observation-only wait does not block mutation and remains idempotent", async () => { + const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid" }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const operationStore = new MemoryOperationStore(); + const leases = new MemoryWorkspaceLeaseProvider(); + let waitStarted!: () => void; + let finishWait!: () => void; + const started = new Promise((resolve) => (waitStarted = resolve)); + const blocked = new Promise((resolve) => (finishWait = resolve)); + let waits = 0; + let stops = 0; + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + operationStore, + leases, + management: { + waitForResources: async () => { + waits += 1; + waitStarted(); + await blocked; + }, + stop: async () => { + stops += 1; + return ["api"]; + }, + } as unknown as ManagementService, + }); + const waitRequest = () => + app( + request( + "/api/v2/workspaces/demo/resources/wait", + { + method: "POST", + headers: { "idempotency-key": "wait-once", prefer: "respond-async" }, + body: JSON.stringify({ deployments: ["api"] }), + }, + "token", + ), + ); + + const submitted = await waitRequest(); + expect(submitted.status).toBe(202); + await started; + const duplicate = await waitRequest(); + expect(duplicate.status).toBe(202); + expect(waits).toBe(1); + + const stop = await app( + request( + "/api/v2/workspaces/demo/lifecycle", + { + method: "POST", + headers: { "idempotency-key": "stop-during-wait" }, + body: JSON.stringify({ action: "stop" }), + }, + "token", + ), + ); + expect(stop.status).toBe(200); + expect(stops).toBe(1); + finishWait(); + }); + test("keeps a concurrent idempotent operation pending when its lease acquisition is denied", async () => { const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid", @@ -1048,10 +1167,12 @@ describe("kuber v2 HTTP routes", () => { acquire: async () => { acquireCount += 1; if (acquireCount === 2) return undefined; - signalFirstAcquire(); - await new Promise((resolve) => { - allowFirstAcquire = resolve; - }); + if (acquireCount === 1) { + signalFirstAcquire(); + await new Promise((resolve) => { + allowFirstAcquire = resolve; + }); + } return { workspaceId: "demo", holder: "operation-same-key", @@ -1096,8 +1217,8 @@ describe("kuber v2 HTTP routes", () => { expect((await operationStore.get("operation-same-key"))?.status.state).toBe( "succeeded", ); - expect(acquireCount).toBe(2); - expect(releaseCount).toBe(1); + expect(acquireCount).toBe(4); + expect(releaseCount).toBe(3); }); test("lets the lease winner claim execution even when a duplicate created the operation", async () => { @@ -1219,7 +1340,7 @@ describe("kuber v2 HTTP routes", () => { ) ).status, ).toBe(500); - expect(releases).toBe(1); + expect(releases).toBe(3); }); test("fails a claimed operation before execution when lease ownership is lost", async () => { @@ -1268,7 +1389,7 @@ describe("kuber v2 HTTP routes", () => { ); expect(result.status).toBe(409); expect(executions).toBe(0); - expect(releases).toBe(1); + expect(releases).toBe(3); expect(await result.json()).toMatchObject({ code: "WORKSPACE_LEASE_LOST", }); @@ -1314,7 +1435,7 @@ describe("kuber v2 HTTP routes", () => { workspaceId: "demo", holder: "operation-blocked-loss", expiresAt: new Date().toISOString(), - renew: async () => ++renewals === 1, + renew: async () => ++renewals <= 3, release: async () => {}, }), }, @@ -1426,14 +1547,14 @@ describe("kuber v2 HTTP routes", () => { ), ); await started; - expect(renewals).toBe(1); + expect(renewals).toBe(3); scheduledRenewal?.(); await Promise.resolve(); await Promise.resolve(); - expect(renewals).toBe(2); + expect(renewals).toBe(6); releaseOperation(); expect((await response).status).toBe(200); - expect(renewals).toBe(3); + expect(renewals).toBe(9); expect(clearedTimers).toEqual([0]); } finally { globalThis.setTimeout = originalSetTimeout; diff --git a/tests/server/audit-store.test.ts b/tests/server/audit-store.test.ts index 42a6212..5f1984c 100644 --- a/tests/server/audit-store.test.ts +++ b/tests/server/audit-store.test.ts @@ -44,6 +44,20 @@ describe("audit store", () => { }); }); + test("memory persistence retains only the newest 100 events", async () => { + const store = new MemoryAuditStore( + () => new Date("2026-01-01T00:00:00.000Z"), + (() => { let id = 0; return () => `event-${++id}`; })(), + ); + for (let index = 0; index < 105; index++) { + await store.append({ actor: { username: "admin" }, action: `action-${index}`, outcome: "success" }); + } + const events = await store.list(); + expect(events).toHaveLength(100); + expect(events[0]?.spec.action).toBe("action-5"); + expect(events.at(-1)?.spec.action).toBe("action-104"); + }); + test("redacts credential-bearing URL strings at value level", () => { expect( redactAuditValue("git clone https://alice:s3cret@github.com/org/repo.git"), diff --git a/tests/server/auth.test.ts b/tests/server/auth.test.ts index 83afe79..bbfad32 100644 --- a/tests/server/auth.test.ts +++ b/tests/server/auth.test.ts @@ -23,6 +23,53 @@ describe("authorization", () => { }); describe("memory auth store", () => { + test("maintains the API-key hash index across key mutations", async () => { + const store = new MemoryAuthStore(); + await store.putUser({ + username: "alice", + passwordHash: "hash", + roles: ["viewer"], + }); + const activeHash = hashToken("active-api-key"); + const expiredHash = hashToken("expired-api-key"); + const activeKey = { + id: "active-key-id-1234", + tokenHash: activeHash, + username: "alice", + capabilities: ["kubernetes:read" as (typeof CAPABILITIES)[number]], + expiresAt: new Date(Date.now() + 60_000).toISOString(), + }; + await store.createApiKey(activeKey); + expect(store.apiKeysByTokenHash.get(activeHash)).toBe( + store.apiKeys.get(activeKey.id), + ); + expect(await store.getApiKey(activeHash)).toEqual( + store.apiKeys.get(activeKey.id), + ); + + expect(await store.revokeApiKey("alice", activeKey.id)).toBe(true); + expect(store.apiKeysByTokenHash.has(activeHash)).toBe(false); + expect(await store.getApiKey(activeHash)).toBeUndefined(); + + const expiredKey = { + ...activeKey, + id: "expired-key-id-1234", + tokenHash: expiredHash, + expiresAt: "2020-01-01T00:00:00.000Z", + }; + await store.createApiKey(expiredKey); + expect(store.apiKeysByTokenHash.get(expiredHash)).toBe( + store.apiKeys.get(expiredKey.id), + ); + expect(await store.deleteExpiredApiKeys()).toBe(1); + expect(store.apiKeysByTokenHash.has(expiredHash)).toBe(false); + + await store.createApiKey(activeKey); + expect(await store.deleteUser("alice")).toBe(true); + expect(store.apiKeysByTokenHash.has(activeHash)).toBe(false); + expect(await store.getApiKey(activeHash)).toBeUndefined(); + }); + test("creates, lists, updates, and deletes users", async () => { const store = new MemoryAuthStore(); const bob = await store.createUser({ diff --git a/tests/server/build-kubernetes.test.ts b/tests/server/build-kubernetes.test.ts index 579b0f9..1a57b11 100644 --- a/tests/server/build-kubernetes.test.ts +++ b/tests/server/build-kubernetes.test.ts @@ -213,6 +213,58 @@ function configMap(record: BuildRecord, imageLabel = true): ConfigMap { } describe("KubernetesBuildStore", () => { + test("prunes only old terminal records with resourceVersion preconditions", async () => { + const fake = new FakeObjects(); + for (let i = 0; i < 52; i++) { + const record = build( + `history-${i}`, + new Date(Date.UTC(2026, 0, 1, 0, 0, i)).toISOString(), + ); + record.status.state = i === 0 ? "queued" : "failed"; + if (i === 1) record.status.reconcileLease = { + holder: "worker", + token: "active-lease", + expiresAt: new Date(Date.now() + 60_000).toISOString(), + }; + fake.maps.set(name("build", record.metadata.name), configMap(record)); + } + const old = build( + "protected-lock", + new Date(Date.UTC(2025, 0, 1)).toISOString(), + ); + old.status.state = "failed"; + fake.maps.set(name("build", old.metadata.name), configMap(old)); + const lockName = name("build-lock", old.spec.imageKey); + fake.maps.set(lockName, { + apiVersion: "v1", + kind: "ConfigMap", + metadata: { + name: lockName, + namespace, + resourceVersion: "1", + labels: { "kuber.astrxl.dev/type": "build-lock" }, + }, + data: { payload: JSON.stringify({ buildId: old.metadata.name }) }, + }); + + await store(fake).pruneBuilds(50); + expect(fake.maps.has(name("build", "history-0"))).toBe(true); + expect(fake.maps.has(name("build", "history-1"))).toBe(true); + expect(fake.maps.has(name("build", "protected-lock"))).toBe(true); + expect(fake.maps.has(name("build", "history-2"))).toBe(false); + expect( + [...fake.maps.values()].filter((value) => value.metadata.labels?.["kuber.astrxl.dev/type"] === "build"), + ).toHaveLength(50); + expect(fake.deletes.length).toBeGreaterThan(0); + expect( + fake.deletes.every( + ({ options }) => !!options?.preconditions?.resourceVersion, + ), + ).toBe(true); + await store(fake).pruneBuilds(50); + expect(fake.maps.has(name("build", "history-1"))).toBe(true); + }); + test("creates records with valid, deterministic image index labels", async () => { const fake = new FakeObjects(); const record = build("validated", "2026-09-02T00:00:00.000Z"); diff --git a/tests/server/build-store.test.ts b/tests/server/build-store.test.ts index dfd7d57..987c479 100644 --- a/tests/server/build-store.test.ts +++ b/tests/server/build-store.test.ts @@ -227,4 +227,97 @@ describe("build store", () => { (await store.listBuilds()).map((record) => record.metadata.name), ).toEqual(["a", "legacy", "replacement"]); }); + + test("prunes only old terminal builds and keeps protected or active records", async () => { + const store = new MemoryBuildStore(); + for (let i = 0; i < 55; i++) { + const record = build( + `build-${String(i).padStart(2, "0")}`, + `image-${i}`, + new Date(Date.UTC(2026, 0, 1, 0, 0, i)).toISOString(), + ); + await store.createBuild(record); + if (i < 53) { + const current = (await store.getBuild(record.metadata.name))!; + const next = structuredClone(current); + next.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + next.status.state = "failed"; + next.status.finishedAt = current.status.createdAt; + await store.replaceBuild(next, current.metadata.resourceVersion); + } + } + const leased = (await store.getBuild("build-04"))!; + const leasedNext = structuredClone(leased); + leasedNext.metadata.resourceVersion = String( + Number(leased.metadata.resourceVersion) + 1, + ); + leasedNext.status.reconcileLease = { + holder: "worker", + token: "token", + expiresAt: new Date(Date.now() + 60_000).toISOString(), + }; + await store.replaceBuild(leasedNext, leased.metadata.resourceVersion); + + await store.pruneBuilds(50); + const remaining = await store.listBuilds(); + expect(remaining.map((record) => record.metadata.name)).toContain( + "build-04", + ); + expect(remaining.map((record) => record.metadata.name)).toContain( + "build-53", + ); + expect(remaining.map((record) => record.metadata.name)).toContain( + "build-54", + ); + expect( + remaining.some((record) => record.metadata.name === "build-00"), + ).toBe(false); + expect(remaining.some((record) => record.status.state === "queued")).toBe( + true, + ); + expect(remaining).toHaveLength(50); + await store.pruneBuilds(50); + expect(await store.getBuild("build-00")).toBeUndefined(); + }); + + test("continues past protected oldest records to prune eligible terminals", async () => { + const store = new MemoryBuildStore(); + for (let i = 0; i < 52; i++) { + const record = build(`retention-${i}`, `image-${i}`, + new Date(Date.UTC(2026, 0, 1, 0, 0, i)).toISOString()); + await store.createBuild(record); + const current = (await store.getBuild(record.metadata.name))!; + const next = structuredClone(current); + next.metadata.resourceVersion = "2"; + next.status.state = i === 51 ? "queued" : "failed"; + if (i < 2) next.status.reconcileLease = { + holder: "worker", token: `lease-${i}`, + expiresAt: new Date(Date.now() + 60_000).toISOString(), + }; + await store.replaceBuild(next, "1"); + } + await store.pruneBuilds(50); + const remaining = await store.listBuilds(); + expect(remaining).toHaveLength(50); + expect(remaining.map((record) => record.metadata.name)).toContain("retention-0"); + expect(remaining.map((record) => record.metadata.name)).toContain("retention-1"); + expect(remaining.map((record) => record.metadata.name)).toContain("retention-51"); + expect(await store.getBuild("retention-2")).toBeUndefined(); + }); + + test("preserves more active records than the retention limit", async () => { + const store = new MemoryBuildStore(); + for (let i = 0; i < 52; i++) + await store.createBuild( + build( + `active-${i}`, + `image-${i}`, + new Date(Date.UTC(2026, 0, 1, 0, 0, i)).toISOString(), + ), + ); + await store.pruneBuilds(50); + expect(await store.listBuilds()).toHaveLength(52); + }); }); diff --git a/tests/server/kubernetes-audit-store.test.ts b/tests/server/kubernetes-audit-store.test.ts new file mode 100644 index 0000000..460c060 --- /dev/null +++ b/tests/server/kubernetes-audit-store.test.ts @@ -0,0 +1,75 @@ +import { expect, test } from "bun:test"; +import type { KubernetesObject } from "@kubernetes/client-node"; +import { KubernetesAuditPersistence } from "../../server/kubernetes-state"; +import type { AuditEvent } from "../../server/audit-store"; + +type Obj = KubernetesObject & { data?: Record }; +class Objects { + values = new Map(); + revision = 0; + deletes = 0; + failReplace = false; + failCreate = false; + key(v: KubernetesObject) { return `${v.kind}:${v.metadata?.name}`; } + async read(v: KubernetesObject) { + const found = this.values.get(this.key(v)); + if (!found) throw { code: 404 }; + return structuredClone(found); + } + async create(v: Obj) { + if (this.failCreate) throw new Error("write failed"); + const key = this.key(v); + if (this.values.has(key)) throw { code: 409 }; + const saved = structuredClone(v); + saved.metadata = { ...saved.metadata, resourceVersion: String(++this.revision) }; + this.values.set(key, saved); + return saved; + } + async replace(v: Obj) { + if (this.failReplace) throw new Error("write failed"); + const key = this.key(v), previous = this.values.get(key); + if (!previous || previous.metadata?.resourceVersion !== v.metadata?.resourceVersion) throw { code: 409 }; + const saved = structuredClone(v); + saved.metadata = { ...saved.metadata, resourceVersion: String(++this.revision) }; + this.values.set(key, saved); + return saved; + } + async delete(v: KubernetesObject) { this.deletes++; this.values.delete(this.key(v)); } + async list(_api: string, _kind: string, _ns: string, _p?: string, _e?: boolean, _x?: boolean, _f?: string, selector?: string) { + const [label, expected] = selector?.split("=") ?? []; + return { items: [...this.values.values()].filter((v) => !label || v.metadata?.labels?.[label] === expected) }; + } + api() { return this as never; } +} +function event(id: number, workspaceId?: string): AuditEvent { + return { apiVersion: "kuber.astrxl.dev/v2", kind: "AuditEvent", metadata: { name: `audit-${id}`, uid: `uid-${id}`, resourceVersion: "1", creationTimestamp: new Date(id * 1000).toISOString() }, spec: { actor: { username: "user" }, action: `action-${id}`, outcome: "success", ...(workspaceId && { workspaceId }) } }; +} + +test("Kubernetes audit persistence stores one bounded ring and filters workspaces", async () => { + const objects = new Objects(), persistence = new KubernetesAuditPersistence(objects.api()); + for (let i = 0; i < 105; i++) await persistence.append(event(i, i % 2 ? "one" : "two")); + expect([...objects.values.values()].filter((v) => v.metadata?.labels?.["kuber.astrxl.dev/type"] === "audit-log")).toHaveLength(1); + expect(await persistence.list()).toHaveLength(100); + expect((await persistence.list("one")).every((v) => v.spec.workspaceId === "one")).toBe(true); + expect((await persistence.list())[0]?.spec.action).toBe("action-5"); +}); + +test("conflicting replicas preserve both appended events", async () => { + const objects = new Objects(), a = new KubernetesAuditPersistence(objects.api()), b = new KubernetesAuditPersistence(objects.api()); + await Promise.all(Array.from({ length: 40 }, (_, i) => (i % 2 ? a : b).append(event(i)))); + expect((await a.list()).map((v) => v.metadata.uid)).toHaveLength(40); +}); + +test("legacy records are migrated before deletion and newest 100 survive", async () => { + const objects = new Objects(); + for (let i = 0; i < 110; i++) await objects.create({ apiVersion: "v1", kind: "ConfigMap", metadata: { name: `audit-${i}`, namespace: "kuber-system", labels: { "kuber.astrxl.dev/type": "audit" } }, data: { payload: JSON.stringify(event(i)) } }); + const persistence = new KubernetesAuditPersistence(objects.api()); + objects.failCreate = true; + await expect(persistence.list()).rejects.toThrow("write failed"); + expect(objects.deletes).toBe(0); + objects.failCreate = false; + const migrated = await persistence.list(); + expect(migrated).toHaveLength(100); + expect(migrated[0]?.spec.action).toBe("action-10"); + expect(objects.values.has("ConfigMap:audit-log")).toBe(true); +}); diff --git a/tests/server/kubernetes-state.test.ts b/tests/server/kubernetes-state.test.ts index 5b6ef05..9bc2c60 100644 --- a/tests/server/kubernetes-state.test.ts +++ b/tests/server/kubernetes-state.test.ts @@ -937,6 +937,71 @@ class FakeApps { } } +test("production rollout wait reports current pod crashes promptly", async () => { + const current = { + metadata: { + uid: "rs-current", + namespace: "shop", + ownerReferences: [{ kind: "Deployment", uid: "deployment-uid" }], + }, + spec: { + template: { + metadata: { labels: { app: "web", "pod-template-hash": "current" } }, + spec: { containers: [{ name: "web", image: "web:v2" }] }, + }, + }, + }; + const clients = managementClients([], []); + clients.apps.readNamespacedDeployment = async () => + ({ + metadata: { name: "web", uid: "deployment-uid", generation: 2 }, + spec: { + replicas: 1, + template: { + metadata: { labels: { app: "web" } }, + spec: { containers: [{ name: "web", image: "web:v2" }] }, + }, + }, + status: { + observedGeneration: 2, + updatedReplicas: 1, + availableReplicas: 0, + unavailableReplicas: 1, + }, + }) as never; + clients.apps.listNamespacedReplicaSet = async () => + ({ items: [current] }) as never; + clients.core.listNamespacedPod = async () => + ({ + items: [ + { + metadata: { + name: "web-pod", + ownerReferences: [{ kind: "ReplicaSet", uid: "rs-current" }], + }, + status: { + containerStatuses: [ + { + name: "web", + restartCount: 3, + state: { waiting: { reason: "CrashLoopBackOff" } }, + }, + ], + }, + }, + ], + }) as never; + await expect( + createKubernetesManagementDependencies(clients).waitForDeployment( + "shop", + "web", + 1000, + ), + ).rejects.toThrow( + "Deployment web rollout failed: pod web-pod, container web: CrashLoopBackOff", + ); +}); + type ManagementClientsType = NonNullable< Parameters[0] >; diff --git a/tests/server/log-service.test.ts b/tests/server/log-service.test.ts index 495511f..f1907e0 100644 --- a/tests/server/log-service.test.ts +++ b/tests/server/log-service.test.ts @@ -1,4 +1,8 @@ -import { describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, test } from "bun:test"; +import { KubeConfig } from "@kubernetes/client-node"; +import http from "node:http"; +import https from "node:https"; +import { KubernetesLogs } from "../../server/kubernetes-logs"; import { createLogService, encodeLogEvent, @@ -108,6 +112,138 @@ function setup() { return { backend, service }; } +const servers: Array = []; +const tlsFixture = { + key: Buffer.from(`-----BEGIN PRIVATE KEY----- +MIIEvAIBADANBgkqhkiG9w0BAQEFAASCBKYwggSiAgEAAoIBAQCUiNil5x0/YDDX +VWBT9hQQOldUzaRdCooc+Nx3CaZaKtjIlWDJQ24GXz+bJPi9kDhNTpoGAKy7EiMj +C+8pXS3ggnMEGIHF89eYzgsMlKUZeZ6x+q9blw04tdKRHRpA96Hps5MZ0FQFZ6p5 +WzPqlyymvP8lAbeYhugekQjzqWfpmecNE3in7WpJH4XkM9w30u9mvjzV+5iyyIuD +HI27Hqujbqm2zkEN07EU1RQMTqTXGzlz7oX5bmfhG04w/47Jz1rLJ80qISWCmVD9 +8hbRsg2j67tHWMOoBhOa7qh1Tme67LCNEFdbFuc5UD1mP+22K9D/zCC2RtJQu5Sq +dEJW3iovAgMBAAECggEAD2K3ckPk0y4/EOcOkdPhEyc/6ZBdkKepU8PxbkEpIpji +mLBkdKSP7ogKOiNTwqsAMf3M1YdXXQ9NZXF0hgfZWzKYAFobgyo1cGYTXeu9yExB +RHVPmcClRXUMCS0HDai49FC+EYPzWBX7YhOw5oFfRiw4j5hEcL+0ponmb/rhwSAg +AGGIM5nGBT3wpXcbOubdxe8/Ex7vg+refUrDGs4T7sbYijoWm5utz3bdc/JxLDf0 +KwMesf7RZ+eeO2C6qOB4gKNjZr5JxApcJZ2WoG9MdWkaYnU8Kyquf2vhJdDZmqtx +kV0y8KSWiPTbVqXZOYcg5AKcymdF7P04zVe4rkEuoQKBgQDRDz9jem5r9U25J6VT +Oi7oWZmgsDcchfeerRBhQmEPiASkOFu2VgiM3CeUeGnNmrRlRsKUEoS00Bph8AMk +z0SfnBTFjIIiB9gERnqrH1vSnCbLhvnhkd7XVtoIG3JjXnJgALoIkmp9DM82VTFu +kBaPpOx1i5vQyI68xpAUNhjVjwKBgQC14p5qxrwn5y0VQg/9MkgLUfbsROqRX8SP +O1RRhR8teA8lL+jnTzq7HauPMnP8kEjfMHXjV124eDLKUHLrOoNwTg7S0MmgYcRU +DWaOCA5u5MpHxCfz2h+OGDVp0FP/QCyLwu6vjkvK+dMhvzK92QHzTqLsia2Tpspx +3HFLEm1RYQKBgG44E7tmuQDB+5A6jrcqXcCyPISzYtru5nYJ2DDuxi1iENBjxjaD +dU6OY2+rbFyxy5n5jGx0tvJ9JOutlnq5q/xaVbkxMwquB/15CwNdLRQEr49uQh/i +wBHYAGt1zQEGslZbC7mpN+tl7Xk/wSgBX2OsF96BFE0m79om9Z8yRjWRAoGAAeBI +iglqv26fBG0eBRqTq6o4xc8gLEe0m1WdVQnufGWUommQGXKzxGJV9rAqihxi5Ap3 +7NRl3xU+UN/rj4mW+X2UoZANxF29zLAmsqhancI2Y+8eCmHhmXGee2zusN9Ulkx4 +cc8h8QIKr3ptZ4/peT0CaTYyWCeMRwhjEscp4YECgYBlQyBkjrhBLswIK+Avul0C +heCCZkDZstdkYTlAh6+tnf0ycMEKXPF9T1O67fSRmzQnTuvb6VWgtMh3zvOMy+Ce +UnVMOhKLS0wPPTZc75kKju1C2hHZ2eSgOFbNFrQvWJMXGNjEvdscXjiCoCBbhBt+ +Pc3eI1dwQpPwp363CGPLrg== +-----END PRIVATE KEY-----\n`), + cert: Buffer.from(`-----BEGIN CERTIFICATE----- +MIIC3DCCAcSgAwIBAgIIMGH60RwCdvswDQYJKoZIhvcNAQEMBQAwFDESMBAGA1UE +AxMJMTI3LjAuMC4xMB4XDTI2MDkyNzEwMjUyN1oXDTM2MDkyNDEwMjUyN1owFDES +MBAGA1UEAxMJMTI3LjAuMC4xMIIBIjANBgkqhkiG9w0BAQEFAAOCAQ8AMIIBCgKC +AQEAlIjYpecdP2Aw11VgU/YUEDpXVM2kXQqKHPjcdwmmWirYyJVgyUNuBl8/myT4 +vZA4TU6aBgCsuxIjIwvvKV0t4IJzBBiBxfPXmM4LDJSlGXmesfqvW5cNOLXSkR0a +QPeh6bOTGdBUBWeqeVsz6pcsprz/JQG3mIboHpEI86ln6ZnnDRN4p+1qSR+F5DPc +N9LvZr481fuYssiLgxyNux6ro26pts5BDdOxFNUUDE6k1xs5c+6F+W5n4RtOMP+O +yc9ayyfNKiElgplQ/fIW0bINo+u7R1jDqAYTmu6odU5nuuywjRBXWxbnOVA9Zj/t +tivQ/8wgtkbSULuUqnRCVt4qLwIDAQABozIwMDAdBgNVHQ4EFgQUeDmjpd05dSFB +yENQKuBjSFVItWQwDwYDVR0RBAgwBocEfwAAATANBgkqhkiG9w0BAQwFAAOCAQEA +SZLDBqpXGPrmKfz3qvPdda3wJUp5MAvy1ZEA2IAMHMutUGemIcEXoU+4qd8aNhzz +EVS2gKsaFZ59MXqC0xhRnNreqP3lP5lPVUV+EAOXvAF+vrg07KBRrxYuE2hzgKxl +FP6swGVsMGyWF266PFe4p3QhcxSUvLBtCuYyJTEvL5FDhByE5M+F7zh5RFiiHvx5 +txljZfJcP0wDZVUalNGbJR1nXJ7GRVPO92ec38ly6PIvlR/BAvNCoLdIF69C12qC +c+5Sp92YbZ6+Lq1p75M6mY44A8FSXE40x3HXM8HcpeatXOxZuM3YAkrO337AgEx0 +9CZrJPz0ekLRl4UW5LG8mA== +-----END CERTIFICATE-----\n`), +}; +const unrelatedCa = Buffer.from(`-----BEGIN CERTIFICATE----- +MIICwzCCAaugAwIBAgIISfrhnhbLiWQwDQYJKoZIhvcNAQEMBQAwEDEOMAwGA1UE +AxMFd3JvbmcwHhcNMjYwOTI3MTAyNzAyWhcNMzYwOTI0MTAyNzAyWjAQMQ4wDAYD +VQQDEwV3cm9uZzCCASIwDQYJKoZIhvcNAQEBBQADggEPADCCAQoCggEBAK4MZhty +U3AIPjb92r7F7tsFX3UUo9fdZUeLjli2a97DZTiO/hs7a/oce27MlXq0BcpP6Rm5 +d0kCIb6IU7SgP3sH8dovaLnC1Xv9KdyFnkRp0eXD6xExmZ9Y6ftDHRkEI+la41ai +GdRp/ebLnTaRLxRAQUpd2xCvaa/weenGDxKFNNunVncF9WzAWYhUZqFR1wXBtGo0 +0AAuzuxYoenyFLHoozeoBJ7Cr23aR/MP0zWQK3blz6TZ7fTcZtDkgsThga+6F8O+ +OJGM7di6FRx7O37DRK6yVphgI1nClbEmw/sAtqc3MYQAh6CzoRFbceGQ6QEdGjq3 +zrv7tJhjQ2x/CNECAwEAAaMhMB8wHQYDVR0OBBYEFKIQJVhaqJhy9cwoA1lCUPff +O//lMA0GCSqGSIb3DQEBDAUAA4IBAQCQHcUDcgLie6wFWpCUFgcT0UMBcZ7FL3me +3vmvOv7T/ydOJGbgaKqz+c3NdrApsL9Jk0wmh0YelYK7TuhZnsNJ0pbehe509vJ9 +10B8bKNmZqgVgdTppsaw9UDFcEQn1UMhkhtYs51+acVntNMXz8/qEhrto8mGdPKO +39ydfekUA6fAw9Q8F/sRtvG7bB6Cmj/8Tieq/dJjO7ZXgqY6HaYrX9qQxlapRu2l +LgdidPbFADDoUsSXBVqX3F4Il+qRQrrGmB/TPPbGCsW/A0TA2jOPHGoh3mczjEm7 +UDgBJ/wg7flQRRYcz278Hm+71JrSOEdpC4F1ssbttKE5Jhg+91uU +-----END CERTIFICATE-----\n`); + +async function kubernetesHttpsServer( + tls: typeof tlsFixture, + options: { ca?: Buffer; skipTLSVerify?: boolean }, + handler: http.RequestListener, +): Promise { + const server = https.createServer(tls, handler); + servers.push(server); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", resolve); + }); + const address = server.address(); + if (!address || typeof address === "string") throw new Error("Missing port"); + const config = new KubeConfig(); + config.loadFromOptions({ + clusters: [{ + name: "test", + server: `https://127.0.0.1:${address.port}`, + ...(options.ca && { caData: options.ca.toString("base64") }), + ...(options.skipTLSVerify && { skipTLSVerify: true }), + }], + users: [{ name: "test" }], + contexts: [{ name: "test", cluster: "test", user: "test" }], + currentContext: "test", + }); + return new KubernetesLogs(config); +} + +async function kubernetesServer( + handler: http.RequestListener, +): Promise { + const server = http.createServer(handler); + servers.push(server); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const address = server.address(); + if (!address || typeof address === "string") throw new Error("Missing port"); + const config = new KubeConfig(); + config.loadFromOptions({ + clusters: [ + { + name: "test", + server: `http://127.0.0.1:${address.port}`, + skipTLSVerify: true, + }, + ], + users: [{ name: "test" }], + contexts: [{ name: "test", cluster: "test", user: "test" }], + currentContext: "test", + }); + return new KubernetesLogs(config); +} + +afterEach(async () => { + await Promise.all( + servers.splice(0).map( + (server) => + new Promise((resolve, reject) => + server.listening + ? server.close((error) => (error ? reject(error) : resolve())) + : resolve(), + ), + ), + ); +}); + describe("LogService", () => { test("enumerates every container in managed deployment pods", async () => { const { backend, service } = setup(); @@ -204,6 +340,76 @@ describe("LogService", () => { expect(JSON.parse(encodeLogEvent(events[0]!))).toEqual(events[0]); }); + test("passes a line count through a followed Kubernetes log request", async () => { + let received: URL | undefined; + const kubernetes = await kubernetesServer((request, response) => { + received = new URL(request.url!, "http://localhost"); + response.writeHead(200); + response.end("latest\n"); + }); + const controller = new AbortController(); + const iterator = kubernetes.streamContainerLogs( + { + namespace: "demo", + pod: "web-abc", + container: "web", + tailLines: 12, + timestamps: false, + }, + controller.signal, + ); + const chunks: Uint8Array[] = []; + for await (const chunk of iterator) chunks.push(chunk); + expect(new TextDecoder().decode(Buffer.concat(chunks))).toBe("latest\n"); + expect(received?.searchParams.get("tailLines")).toBe("12"); + expect(received?.searchParams.get("follow")).toBe("true"); + }); + + test("streams logs from HTTPS trusted by the configured cluster CA", async () => { + const tls = tlsFixture; + const kubernetes = await kubernetesHttpsServer( + tls, + { ca: tls.cert }, + (_request, response) => response.end("trusted\n"), + ); + const controller = new AbortController(); + const chunks: Uint8Array[] = []; + for await (const chunk of kubernetes.streamContainerLogs({ + namespace: "demo", pod: "web-abc", container: "web", timestamps: false, + }, controller.signal)) chunks.push(chunk); + expect(Buffer.concat(chunks).toString()).toBe("trusted\n"); + }); + + test("rejects HTTPS logs with an untrusted cluster CA", async () => { + const tls = tlsFixture; + const kubernetes = await kubernetesHttpsServer( + tls, + { ca: unrelatedCa }, + (_request, response) => response.end("should not be reached\n"), + ); + const controller = new AbortController(); + await expect(async () => { + for await (const _chunk of kubernetes.streamContainerLogs({ + namespace: "demo", pod: "web-abc", container: "web", timestamps: false, + }, controller.signal)) { /* consume until TLS rejects */ } + }).toThrow(); + }); + + test("allows an untrusted HTTPS log certificate only with explicit opt-in", async () => { + const tls = tlsFixture; + const kubernetes = await kubernetesHttpsServer( + tls, + { skipTLSVerify: true }, + (_request, response) => response.end("opted in\n"), + ); + const controller = new AbortController(); + const chunks: Uint8Array[] = []; + for await (const chunk of kubernetes.streamContainerLogs({ + namespace: "demo", pod: "web-abc", container: "web", timestamps: false, + }, controller.signal)) chunks.push(chunk); + expect(Buffer.concat(chunks).toString()).toBe("opted in\n"); + }); + test("returns per-container safe errors and continues collection", async () => { const { backend, service } = setup(); backend.reads.set( @@ -364,6 +570,159 @@ describe("LogService", () => { expect(attempts).toBe(2); }); + test("emits a coded TLS failure once without retrying", async () => { + const { backend, service } = setup(); + backend.pods[0] = { ...backend.pods[0]!, containers: ["web"] }; + let attempts = 0; + backend.streams.set("web-abc/web", () => + (async function* () { + attempts += 1; + yield* [] as string[]; + const error = new Error("self signed certificate in certificate chain") as Error & { + code: string; + }; + error.code = "SELF_SIGNED_CERT_IN_CHAIN"; + throw error; + })(), + ); + const controller = new AbortController(); + const iterator = service.follow({ + namespace: "demo", + target: { kind: "managed-deployments" }, + signal: controller.signal, + discoveryIntervalMs: 5, + heartbeatIntervalMs: 1_000, + }); + expect(await nextOfType(iterator, "error")).toMatchObject({ + message: "self signed certificate in certificate chain", + retryable: false, + }); + expect(await nextOfType(iterator, "heartbeat")).toMatchObject({ + type: "heartbeat", + }); + expect(attempts).toBe(1); + controller.abort(); + expect((await iterator.next()).done).toBe(true); + expect(attempts).toBe(1); + }); + + test("reports a broken upstream log connection per container without losing other streams", async () => { + const urls: URL[] = []; + const kubernetes = await kubernetesServer((request, response) => { + const url = new URL(request.url!, "http://localhost"); + urls.push(url); + if (url.searchParams.get("container") === "web") { + request.socket.destroy(); + } else { + response.writeHead(200); + response.write("sidecar healthy\n"); + } + }); + const { backend, service } = setup(); + backend.streams.set("web-abc/web", (signal) => + kubernetes.streamContainerLogs( + { namespace: "demo", pod: "web-abc", container: "web", timestamps: false }, + signal, + ), + ); + backend.streams.set("web-abc/sidecar", (signal) => + kubernetes.streamContainerLogs( + { namespace: "demo", pod: "web-abc", container: "sidecar", timestamps: false }, + signal, + ), + ); + const controller = new AbortController(); + const iterator = service.follow({ + namespace: "demo", + target: { kind: "managed-deployments" }, + signal: controller.signal, + discoveryIntervalMs: 1_000, + heartbeatIntervalMs: 1_000, + retryIntervalMs: 1_000, + }); + try { + const events = await Promise.all([iterator.next(), iterator.next()]); + expect(events.map(({ value }) => value?.type).sort()).toEqual([ + "error", + "log", + ]); + const error = events.find(({ value }) => value?.type === "error")!.value; + expect(error).toMatchObject({ + pod: "web-abc", + container: "web", + retryable: false, + }); + expect(error?.message).toMatch( + /socket.*connection.*closed|socket hang up|fetch failed/i, + ); + expect(error?.message).not.toContain("Error occurred in log request"); + expect( + events.find(({ value }) => value?.type === "log")!.value, + ).toMatchObject({ container: "sidecar", message: "sidecar healthy" }); + expect(urls).toHaveLength(2); + expect(urls[0]!.searchParams.get("follow")).toBe("true"); + } finally { + controller.abort(); + await iterator.return(undefined); + } + }); + + test.each([ + [ + 403, + JSON.stringify({ kind: "Status", message: "logs forbidden" }), + "logs forbidden", + false, + ], + [500, "proxy unavailable", "proxy unavailable", true], + ])( + "keeps upstream HTTP %i status and body for follow errors", + async (status, body, detail, retryable) => { + const kubernetes = await kubernetesServer((_request, response) => { + response.writeHead(status); + response.end(body); + }); + const read = async () => { + for await (const _chunk of kubernetes.streamContainerLogs( + { + namespace: "demo", + pod: "web-abc", + container: "web", + timestamps: false, + }, + new AbortController().signal, + )) { + // The upstream response must fail before yielding log data. + } + }; + await expect(read()).rejects.toMatchObject({ + message: `Kubernetes log request failed (HTTP ${status}): ${detail}`, + retryable, + }); + }, + ); + + test("aborting a follow request closes the upstream stream", async () => { + let disconnected!: () => void; + const closed = new Promise((resolve) => (disconnected = resolve)); + const kubernetes = await kubernetesServer((_request, response) => { + response.on("close", disconnected); + response.writeHead(200); + response.write("first\n"); + }); + const controller = new AbortController(); + const iterator = kubernetes.streamContainerLogs( + { namespace: "demo", pod: "web-abc", container: "web", timestamps: false }, + controller.signal, + )[Symbol.asyncIterator](); + expect(new TextDecoder().decode((await iterator.next()).value)).toBe( + "first\n", + ); + controller.abort(); + expect((await iterator.next()).done).toBe(true); + await closed; + }); + test("emits one terminal target error and closes follow mode", async () => { const { service } = setup(); const iterator = service.follow({ diff --git a/tests/server/management.test.ts b/tests/server/management.test.ts index f05c38d..2362a53 100644 --- a/tests/server/management.test.ts +++ b/tests/server/management.test.ts @@ -67,6 +67,45 @@ function dependencies( } describe("server management service", () => { + test("waits on deployments together and cancels siblings after a crash", async () => { + const started: string[] = []; + let siblingCancelled = false; + const service = createManagementService( + dependencies({ + listDeployments: async () => [ + object("Deployment", "healthy", "healthy-uid") as V1Deployment, + object("Deployment", "broken", "broken-uid") as V1Deployment, + ], + waitForDeployment: async (_project, name, timeout, execution) => { + expect(timeout).toBeGreaterThan(0); + started.push(name); + if (name === "broken") + throw new Error( + "Deployment broken rollout failed: CrashLoopBackOff", + ); + await new Promise((resolve) => { + if (execution?.signal?.aborted) { + siblingCancelled = true; + return resolve(); + } + execution?.signal?.addEventListener( + "abort", + () => { + siblingCancelled = true; + resolve(); + }, + { once: true }, + ); + }); + }, + }), + ); + await expect( + service.waitForResources(workspace, ["healthy", "broken"], 1000), + ).rejects.toThrow("Deployment broken rollout failed: CrashLoopBackOff"); + expect(started).toEqual(["healthy", "broken"]); + expect(siblingCancelled).toBe(true); + }); test("contains no authorization policy and operates on an explicit workspace", async () => { const calls: string[] = []; const service = createManagementService(