From 828bf3a3286e62ccede84ef8df8f28737b252998 Mon Sep 17 00:00:00 2001 From: dmgnr Date: Mon, 5 Oct 2026 11:09:52 +0000 Subject: [PATCH] feat: prepare 2.6.1-rc5 shared databases and build SSE --- changelogs/2.6.1.md | 11 + command/main.ts | 2 +- command/up.ts | 88 +++--- lib/api.ts | 116 +++++++- lib/build.ts | 123 ++++++--- lib/database.ts | 188 ++++++++++--- lib/version-update.ts | 229 ++++++++++++++++ package.json | 2 +- server/app.ts | 122 ++++++++- server/build-event-stream.ts | 255 ++++++++++++++++++ server/management.ts | 74 ++++- shared/version.ts | 5 + tests/command/up-api.test.ts | 114 +++++--- tests/lib/api-stream.test.ts | 67 +++++ tests/lib/api.test.ts | 68 +++++ tests/lib/build-api.test.ts | 83 ++++++ tests/lib/database.test.ts | 161 ++++++++++- tests/lib/version-update-process.test.ts | 51 ++++ tests/lib/version-update.test.ts | 201 ++++++++++++++ tests/server/api-keys.test.ts | 44 +-- tests/server/app.test.ts | 223 ++++++++++++++- tests/server/build-event-stream.test.ts | 330 +++++++++++++++++++++++ tests/server/management.test.ts | 115 +++++++- tests/server/operation-store.test.ts | 7 + tests/server/version-header.test.ts | 74 +++++ 25 files changed, 2557 insertions(+), 196 deletions(-) create mode 100644 changelogs/2.6.1.md create mode 100644 lib/version-update.ts create mode 100644 server/build-event-stream.ts create mode 100644 shared/version.ts create mode 100644 tests/lib/api-stream.test.ts create mode 100644 tests/lib/version-update-process.test.ts create mode 100644 tests/lib/version-update.test.ts create mode 100644 tests/server/build-event-stream.test.ts create mode 100644 tests/server/version-header.test.ts diff --git a/changelogs/2.6.1.md b/changelogs/2.6.1.md new file mode 100644 index 0000000..5e15e57 --- /dev/null +++ b/changelogs/2.6.1.md @@ -0,0 +1,11 @@ +# 2.6.1 + +## Added + +- Stream build logs over SSE with a resumable cursor, and fall back to polling when streaming is unavailable. +- Return the Kuber version in the `x-kuber-version` response header and best-effort install a newer global CLI version when detected. The running process is not re-executed, so the update applies to later invocations. + +## Changed + +- Support safe shared PostgreSQL use without taking over database labels; retain Database custom resources when running `down --full`. +- Report sanitized database diagnostic phases without exposing sensitive provider details. diff --git a/command/main.ts b/command/main.ts index 4e1e2b1..a3fc1de 100644 --- a/command/main.ts +++ b/command/main.ts @@ -37,7 +37,7 @@ function cloneCommand(command: T): T { export const main = defineCommand({ meta: { name: "kuber", - version: "2.6.0", + version: "2.6.1-rc5", description: "Docker Compose -> K8s translation layer", }, args: { diff --git a/command/up.ts b/command/up.ts index 6a23e7e..c7211ad 100644 --- a/command/up.ts +++ b/command/up.ts @@ -1,5 +1,5 @@ import { defineCommand } from "citty"; -import { Listr, ListrErrorTypes, type ListrTaskWrapper } from "listr2"; +import { Listr, type ListrTaskWrapper } from "listr2"; import { randomUUID } from "node:crypto"; import { isDeepStrictEqual } from "node:util"; import type { ComposeSpecification } from "../schema/docker.d"; @@ -104,6 +104,8 @@ type UpContext = { workspaceRoot?: string; buildImages?: Record; serviceEnv?: Record>; + databaseEnv?: Record>; + storageEnv?: Record>; resources?: KubernetesResource[]; plan?: ResourcePlan; }; @@ -796,14 +798,20 @@ export async function runUp( cwd, { progress: (message) => { - if (task.output !== message) task.output = message; + if ( + !/^Build (?:queued|creating|starting|running|done)\b/.test(message) && + task.output !== message + ) task.output = message; }, service: (name) => { let stream: ReturnType["stdout"]> | undefined; return { progress: (message) => { const child = children.get(name); - if (child && child.output !== message) child.output = message; + if (!child) return; + const phase = /^Build (queued|creating|starting|running|done)\b/.exec(message)?.[1]; + if (phase) child.title = `Build ${name}: ${phase}`; + // A failed status includes the full error; the parent bottom bar owns it. }, get stream() { const child = children.get(name); @@ -833,7 +841,6 @@ export async function runUp( error instanceof Error ? error : new Error(String(error)); operationFailure = failure; task.output = failure.stack ?? failure.message; - task.report(failure, ListrErrorTypes.HAS_FAILED); for (const resolve of pending.values()) resolve(failure); pending.clear(); }, @@ -851,11 +858,14 @@ export async function runUp( child: ListrTaskWrapper, ) => { children.set(name, child); - child.output = "Build queued"; + child.title = `Build ${name}: queued`; if (children.size === services.length) start(); const error = await completed.get(name); - if (error) throw error; - child.output = "done"; + if (error) { + child.title = `Build ${name}: failed`; + throw error; + } + child.title = `Build ${name}: done`; }, })), { @@ -916,24 +926,29 @@ export async function runUp( }, { title: "Reconcile backing services", - task: async (taskCtx, task) => { - let databaseEnv: Record> = {}; - let storageEnv: Record> = {}; - await task.newListr( + rendererOptions: { bottomBar: 1, persistentOutput: true }, + task: (taskCtx, task) => + 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`; + let response: Record; + try { + response = await managementRequest>( + project, + request, + `${workspacePath}/databases`, + { method: "POST", json: { compose: taskCtx.compose } }, + ); + } catch (error) { + task.output = error instanceof Error ? error.message : String(error); + throw error; + } + taskCtx.databaseEnv = operationEnvironment(response); + child.output = `${Object.keys(taskCtx.databaseEnv).length} services`; }, }, { @@ -941,28 +956,33 @@ export async function runUp( 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`; + let response: Record; + try { + response = await managementRequest>( + project, + request, + `${workspacePath}/storage`, + { method: "POST", json: { compose: taskCtx.compose } }, + ); + } catch (error) { + task.output = error instanceof Error ? error.message : String(error); + throw error; + } + taskCtx.storageEnv = operationEnvironment(response); + child.output = `${Object.keys(taskCtx.storageEnv).length} services`; }, }, ], { concurrent: true }, - ).run(); - taskCtx.serviceEnv = mergeServiceEnv( - mergeServiceEnv(taskCtx.serviceEnv, databaseEnv), - storageEnv, - ); - }, + ), }, { title: "Render manifests", task: async (taskCtx, task) => { + taskCtx.serviceEnv = mergeServiceEnv( + mergeServiceEnv(taskCtx.serviceEnv, taskCtx.databaseEnv ?? {}), + taskCtx.storageEnv ?? {}, + ); taskCtx.resources = await composeToKubernetes( project, taskCtx.compose!, @@ -1085,7 +1105,7 @@ export async function runUp( renderer: "default", // Only redraw when a task changes; a queued build should not generate // another frame on each spinner tick. Keep the renderer's task tree. - rendererOptions: { collapseErrors: false, collapseSubtasks: false, lazy: true }, + rendererOptions: { collapseErrors: false, collapseSubtasks: false, showErrorMessage: false, lazy: true }, }, ).run(); diff --git a/lib/api.ts b/lib/api.ts index c522b31..9248042 100644 --- a/lib/api.ts +++ b/lib/api.ts @@ -1,6 +1,8 @@ import { KUBER_API_BASE_URL } from "../const"; import type { ApiProblemDetails } from "../shared/api"; +import { KUBER_VERSION_HEADER } from "../shared/version"; import { readSession, type KuberSession } from "./session"; +import { observeServerVersion } from "./version-update"; export type { ApiProblemDetails } from "../shared/api"; @@ -19,6 +21,8 @@ export type ApiRequestOptions = { timeoutMs?: number; /** Byte offset for resumable binary upload requests. */ uploadOffset?: number; + /** Overrides version observation for embedded clients and tests. */ + onServerVersion?: (version: string | null) => void; }; export type ApiUploadOptions = ApiRequestOptions & { @@ -181,12 +185,23 @@ async function sendRequest( const headers = await requestHeaders(init, options); const { json, ...requestInit } = init; const body = Object.hasOwn(init, "json") ? JSON.stringify(json) : init.body; - return fetch(`${options.baseUrl ?? KUBER_API_BASE_URL}${path}`, { - ...requestInit, - body, - headers, - signal, - }); + const response = await fetch( + `${options.baseUrl ?? KUBER_API_BASE_URL}${path}`, + { + ...requestInit, + body, + headers, + signal, + }, + ); + try { + (options.onServerVersion ?? observeServerVersion)( + response.headers.get(KUBER_VERSION_HEADER), + ); + } catch { + // Version discovery must never interrupt an API request. + } + return response; } export async function apiRequest( @@ -288,3 +303,92 @@ export async function* apiStreamNdjson( deadline.clear(); } } + +/** Raised only when a server does not offer the SSE build-events representation. */ +export class ApiStreamUnsupportedError extends Error {} + +/** One authenticated SSE connection. The caller owns reconnection and its cursor. */ +export async function* apiStreamEvents( + path: string, + after: number, + signal?: AbortSignal, + options: ApiRequestOptions = {}, +): AsyncGenerator<{ event: T; id?: number; type?: string }, void, void> { + const deadline = requestDeadline(signal, 0); + const headers = new Headers({ accept: "text/event-stream" }); + headers.set("last-event-id", String(after)); + try { + const response = await sendRequest( + path, + { headers }, + options, + deadline.signal, + ); + if ([404, 405, 406, 415, 501].includes(response.status)) + throw new ApiStreamUnsupportedError("Build SSE events are unavailable"); + await assertResponseOk(response); + if ( + !response.headers + .get("content-type") + ?.toLowerCase() + .startsWith("text/event-stream") + ) { + await response.body?.cancel(); + throw new ApiStreamUnsupportedError("Build SSE events are unavailable"); + } + if (!response.body) throw new Error("Empty build SSE response"); + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + let data: string[] = []; + let eventType = ""; + let id: number | undefined; + try { + for (;;) { + const { value, done } = await reader.read(); + buffer += decoder.decode(value, { stream: !done }); + while (buffer.includes("\n")) { + const newline = buffer.indexOf("\n"); + const line = buffer.slice(0, newline).replace(/\r$/, ""); + buffer = buffer.slice(newline + 1); + if (line === "") { + if ( + data.length && + (eventType === "log" || + eventType === "status" || + eventType === "gap") + ) + yield { + event: JSON.parse(data.join("\n")) as T, + id, + ...(eventType === "gap" && { type: eventType }), + }; + data = []; + id = undefined; + eventType = ""; + } else if (line.startsWith("data:")) + data.push(line.slice(5).replace(/^ /, "")); + else if (line.startsWith("event:")) eventType = line.slice(6).trim(); + else if (line.startsWith("id:")) { + const raw = line.slice(3).trim(); + if (/^(0|[1-9]\d*)$/.test(raw) && Number.isSafeInteger(Number(raw))) + id = Number(raw); + } + } + // Bound a malformed or non-SSE response that never terminates a line. + if (buffer.length > 2 * 1024 * 1024) + throw new Error("Build SSE line too long"); + if (done) break; + } + } finally { + try { + await reader.cancel(); + } catch { + /* Upstream aborted. */ + } + reader.releaseLock(); + } + } finally { + deadline.clear(); + } +} diff --git a/lib/build.ts b/lib/build.ts index c4ab13c..a8c2229 100644 --- a/lib/build.ts +++ b/lib/build.ts @@ -12,7 +12,13 @@ import { type Sha256Digest, } from "../shared/build-protocol"; import { resolveComposeArch } from "./arch"; -import { apiRequest, type ApiRequestInit, type ApiRequestOptions } from "./api"; +import { + apiRequest, + apiStreamEvents, + ApiStreamUnsupportedError, + type ApiRequestInit, + type ApiRequestOptions, +} from "./api"; import { DEFAULT_REGISTRY } from "./config"; import { enumerateWorkspace, @@ -153,9 +159,17 @@ export type ApiRequester = ( options?: ApiRequestOptions, ) => Promise; +export type BuildEventStreamer = ( + path: string, + after: number, + signal?: AbortSignal, +) => AsyncIterable<{ event: BuildEvent; id?: number; type?: string }>; + export type BuildOptions = { registry?: string; request?: ApiRequester; + /** For embedded clients; defaults to authenticated SSE with the standard API requester. */ + streamEvents?: BuildEventStreamer; pollIntervalMs?: number; sleep?: (milliseconds: number) => Promise; snapshot?: WorkspaceSnapshot; @@ -496,20 +510,78 @@ async function waitForBuild( initial: BuildStatus, signal?: AbortSignal, prefixService = false, + streamEvents?: BuildEventStreamer, ): Promise { let status = initial; let sequence = 0; let failureLog = ""; // Each physical build belongs only to its destination services. A shared // reporter (the caller's global sink) receives each event just once. - const owners = new Map; - }>(); + const owners = new Map< + BuildReporter | undefined, + { + name: string; + reporter?: BuildReporter; + states: Set; + } + >(); for (const { name, reporter } of reporters) if (!owners.has(reporter)) owners.set(reporter, { name, reporter, states: new Set() }); + const deliver = async (event: BuildEvent, updateStatus = false) => { + if (event.type === "log" && event.sequence <= sequence) return; + if (event.type === "log") + failureLog = (failureLog + event.message).slice(-8_192); + for (const { name, reporter, states } of owners.values()) + await reportBuildEvent( + event, + reporter, + states, + prefixService ? name : undefined, + ); + if (event.type === "log") sequence = event.sequence; + else if (updateStatus) status = event.status; + }; + if (streamEvents) { + const path = `/builds/${encodeURIComponent(id)}/events`; + let unsupported = false; + while (!unsupported) { + try { + signal?.throwIfAborted(); + for await (const { event, id: cursor, type: eventType } of streamEvents( + path, + sequence, + signal, + )) { + if (eventType === "gap") { + const gap = event as unknown as { message?: string }; + throw new Error( + gap.message ?? + "Build log history was trimmed; some log output is unavailable.", + ); + } + if ( + event.type === "log" && + cursor !== undefined && + cursor !== event.sequence + ) + throw new Error("Build SSE log cursor does not match its sequence"); + await deliver(event, true); + if (status.state === "succeeded" || status.state === "failed") break; + } + if (status.state === "succeeded" || status.state === "failed") break; + } catch (error) { + if (error instanceof ApiStreamUnsupportedError) unsupported = true; + else if (signal?.aborted) throw signal.reason; + else if (!isTransientBuildPollError(error)) throw error; + } + if (!unsupported) { + signal?.throwIfAborted(); + await sleep(Math.max(100, pollIntervalMs)); + } + } + if (!unsupported) return withFailureLog(status, failureLog); + } for (;;) { const events = await requestBuildPoll( request, @@ -519,31 +591,9 @@ async function waitForBuild( sleep, signal, ); - for (const event of events) { - if (event.type === "log" && event.sequence <= sequence) continue; - if (event.type === "log") - failureLog = (failureLog + event.message).slice(-8_192); - for (const { name, reporter, states } of owners.values()) - await reportBuildEvent( - event, - reporter, - states, - prefixService ? name : undefined, - ); - if (event.type === "log") sequence = event.sequence; - } + for (const event of events) await deliver(event); if (status.state === "succeeded" || status.state === "failed") { - const details = failureLog.trim(); - if ( - status.state === "failed" && - details && - !status.error?.includes(details) - ) - return { - ...status, - error: `${status.error ?? "BuildKit Job failed"}\n${details}`, - }; - return status; + return withFailureLog(status, failureLog); } status = await requestBuildPoll( request, @@ -563,6 +613,18 @@ async function waitForBuild( } } +function withFailureLog(status: BuildStatus, failureLog: string): BuildStatus { + const details = failureLog.trim(); + return status.state === "failed" && + details && + !status.error?.includes(details) + ? { + ...status, + error: `${status.error ?? "BuildKit Job failed"}\n${details}`, + } + : status; +} + export async function resolveBuildImages( project: string, compose: ComposeSpecification, @@ -753,6 +815,7 @@ export async function buildServices( initial, options.signal, !reporter?.service, + options.streamEvents ?? (options.request ? undefined : apiStreamEvents), ); if (status.state !== "succeeded") throw new Error( diff --git a/lib/database.ts b/lib/database.ts index bdaab42..0149cbc 100644 --- a/lib/database.ts +++ b/lib/database.ts @@ -4,24 +4,65 @@ import type { ComposeSpecification, Service } from "../schema/docker.d"; import { LABELS } from "../const"; import { deleteResource, applyResource } from "./apply"; +export const DATABASE_RECONCILE_PHASES = [ + "namespace precheck", "database dependency", "database resource listing", + "database resource ownership", "claim discovery", "credential preparation", + "role secret lookup", "role secret apply", "cluster lookup", + "managed role preparation", "managed role update", "database preparation", + "database lookup", "database ownership", "database apply", "credential lookup", + "environment assembly", "operation execution", +] as const; + export class DatabaseReconciliationError extends Error { + readonly phase: string; + constructor(phase: string, error: unknown, claim?: PostgresClaim) { - const context = claim - ? ` for database ${claim.database} (service ${claim.service}, role ${claim.username})` - : ""; - const reason = error instanceof Error ? error.message : String(error); - // Provider errors can contain credentials or entire request bodies. Only - // expose a short, recognisable operational reason, never a raw response. - const knownReason = /^(forbidden|not found|conflict|permission denied|connection refused|timed out|timeout|unauthorized|unprocessable entity|service unavailable)\b/i.exec(reason); - const status = error && typeof error === "object" && "code" in error && - typeof error.code === "number" && error.code >= 400 && error.code < 600 - ? ` (HTTP ${error.code})` - : ""; - const safeReason = `${knownReason ? knownReason[1] : "Kubernetes request failed"}${status}`; - super(`Database reconciliation failed during ${phase}${context}: ${safeReason}`, { + const safePhase: string = DATABASE_RECONCILE_PHASES.some((known) => known === phase) + ? phase + : "operation execution"; + // Claim values are user-controlled and may themselves be credentials. + // Keep the actionable phase, but never persist or expose claim identifiers. + const context = claim ? " for requested database claim" : ""; + const provider = + error && typeof error === "object" + ? (error as Record) + : {}; + const body = + provider.body && typeof provider.body === "object" + ? (provider.body as Record) + : {}; + const statusCode = [provider.statusCode, provider.code, body.code].find( + (value) => + typeof value === "number" && + Number.isInteger(value) && + value >= 400 && + value < 600, + ); + const status = statusCode === undefined ? "" : ` (HTTP ${statusCode})`; + // Provider messages and response bodies can contain secrets, connection + // URLs or the full request. Match only a fixed vocabulary, never echo them. + const reason = error instanceof Error ? error.message : ""; + const knownReason = + /^(forbidden|not found|conflict|permission denied|connection refused|timed out|timeout|unauthorized|unprocessable entity|service unavailable)\b/i.exec( + reason, + ); + const providerReason = + typeof body.reason === "string" + ? /^(Forbidden|NotFound|AlreadyExists|Conflict|Unauthorized|Invalid|ServiceUnavailable)$/.exec( + body.reason, + )?.[1] + : undefined; + const fallback = safePhase === "claim discovery" + ? "Check PostgreSQL claim declarations" + : safePhase === "operation execution" && !status + ? "Unexpected failure; check server logs using the operation ID" + : "Kubernetes request failed"; + const safeReason = `${knownReason?.[1] ?? providerReason ?? fallback}${status}`; + super(`Database reconciliation failed during ${safePhase}${context}: ${safeReason}`, { cause: error, }); this.name = "DatabaseReconciliationError"; + this.phase = safePhase; } } @@ -210,8 +251,9 @@ async function readObject( async function ensureRoleSecret( claim: PostgresClaim, ): Promise { + let existing: V1Secret | undefined; try { - const existing = await readObject({ + existing = await readObject({ apiVersion: "v1", kind: "Secret", metadata: { @@ -219,10 +261,12 @@ async function ensureRoleSecret( namespace: DATABASE_NAMESPACE, }, }); - + } catch (error) { + throw new DatabaseReconciliationError("role secret lookup", error, claim); + } + try { const username = claim.username; const password = decodeSecretValue(existing?.data?.password) ?? randomUUID(); - await applyResource({ apiVersion: "v1", kind: "Secret", @@ -239,7 +283,7 @@ async function ensureRoleSecret( return { username, password }; } catch (error) { - throw new DatabaseReconciliationError("role secret setup", error, claim); + throw new DatabaseReconciliationError("role secret apply", error, claim); } } @@ -290,17 +334,22 @@ async function reconcileManagedRoles( if (!cluster) { throw new DatabaseReconciliationError( - `CNPG cluster ${DATABASE_NAMESPACE}/${DATABASE_CLUSTER} lookup`, + "cluster lookup", new Error("Not found"), claims[0], ); } - const roles = new Map( - (cluster.spec?.managed?.roles ?? []).map((role) => [role.name, role]), - ); - for (const claim of claims) { - roles.set(claim.username, toManagedRole(claim)); + let roles: Map; + try { + roles = new Map( + (cluster.spec?.managed?.roles ?? []).map((role) => [role.name, role]), + ); + for (const claim of claims) { + roles.set(claim.username, toManagedRole(claim)); + } + } catch (error) { + throw new DatabaseReconciliationError("managed role preparation", error, claims[0]); } try { @@ -326,17 +375,26 @@ async function reconcileManagedRoles( async function reconcileDatabases( project: string, claims: PostgresClaim[], + existingDatabases: Set, signal?: AbortSignal, ): Promise { const uniqueDatabases = new Map(); - for (const claim of claims) { - uniqueDatabases.set(`${claim.database}:${claim.username}`, claim); + try { + for (const claim of claims) { + uniqueDatabases.set(`${claim.database}:${claim.username}`, claim); + } + } catch (error) { + throw new DatabaseReconciliationError("database preparation", error); } for (const claim of uniqueDatabases.values()) { + // An existing Database may be shared with other workspaces. Applying even + // an identical object can prune labels owned by the SSA field manager. + if (existingDatabases.has(claim.database)) continue; try { throwIfAborted(signal); - await applyResource({ + const { objectApi } = await import("./k8s"); + await objectApi.create({ apiVersion: "postgresql.cnpg.io/v1", kind: "Database", metadata: { @@ -357,22 +415,63 @@ async function reconcileDatabases( name: claim.database, owner: claim.username, }, - }); + } as KubernetesObject); } catch (error) { throw new DatabaseReconciliationError("database apply", error, claim); } } } +async function inspectDatabases( + claims: PostgresClaim[], + signal?: AbortSignal, +): Promise> { + const existing = new Set(); + for (const claim of claims) { + if (existing.has(claim.database)) continue; + let database: (KubernetesObject & { + spec?: { owner?: string; cluster?: { name?: string } }; + }) | undefined; + try { + throwIfAborted(signal); + database = await readObject({ + apiVersion: "postgresql.cnpg.io/v1", + kind: "Database", + metadata: { name: claim.database, namespace: DATABASE_NAMESPACE }, + }); + } catch (error) { + throw new DatabaseReconciliationError("database lookup", error, claim); + } + if (!database) continue; + if ( + database.spec?.owner !== claim.username || + database.spec?.cluster?.name !== DATABASE_CLUSTER + ) { + throw new DatabaseReconciliationError( + "database ownership", + new Error("Existing database does not match the requested role and cluster"), + claim, + ); + } + existing.add(claim.database); + } + return existing; +} + export async function reconcilePostgresClaim( project: string, claim: PostgresClaim, signal?: AbortSignal, ): Promise { - throwIfAborted(signal); + try { + throwIfAborted(signal); + } catch (error) { + throw new DatabaseReconciliationError("credential preparation", error, claim); + } + const existing = await inspectDatabases([claim], signal); const credentials = await ensureRoleSecret(claim); await reconcileManagedRoles([claim], signal); - await reconcileDatabases(project, [claim], signal); + await reconcileDatabases(project, [claim], existing, signal); return credentials; } @@ -381,27 +480,48 @@ export async function reconcilePostgresClaims( compose: ComposeSpecification, signal?: AbortSignal, ): Promise>> { - const claims = getComposePostgresClaims(compose); + let claims: PostgresClaim[]; + try { + claims = getComposePostgresClaims(compose); + } catch (error) { + throw new DatabaseReconciliationError("claim discovery", error); + } if (claims.length === 0) return {}; + // Validate every requested database before touching any role Secret or the + // shared Cluster; one conflicting claim must leave all roles untouched. + const existing = await inspectDatabases(claims, signal); + const credentialsBySecret = new Map(); for (const claim of claims) { if (credentialsBySecret.has(claim.secretName)) continue; - throwIfAborted(signal); + try { + throwIfAborted(signal); + } catch (error) { + throw new DatabaseReconciliationError("credential preparation", error, claim); + } credentialsBySecret.set(claim.secretName, await ensureRoleSecret(claim)); } await reconcileManagedRoles(claims, signal); - await reconcileDatabases(project, claims, signal); + await reconcileDatabases(project, claims, existing, signal); return Object.fromEntries( claims.map((claim) => { const credentials = credentialsBySecret.get(claim.secretName); if (!credentials) { - throw new Error(`Missing credentials for ${claim.secretName}`); + throw new DatabaseReconciliationError( + "credential lookup", + new Error("Missing credentials"), + claim, + ); } - return [claim.service, buildPostgresEnvironment(claim, credentials)]; + try { + return [claim.service, buildPostgresEnvironment(claim, credentials)]; + } catch (error) { + throw new DatabaseReconciliationError("environment assembly", error, claim); + } }), ); } diff --git a/lib/version-update.ts b/lib/version-update.ts new file mode 100644 index 0000000..eb73504 --- /dev/null +++ b/lib/version-update.ts @@ -0,0 +1,229 @@ +import { KUBER_VERSION } from "../shared/version"; +import { readlinkSync, statSync } from "node:fs"; + +type SemVer = { + major: bigint; + minor: bigint; + patch: bigint; + prerelease: string[]; +}; + +const SEMVER = + /^(0|[1-9]\d*)\.(0|[1-9]\d*)\.(0|[1-9]\d*)(?:-([0-9A-Za-z-]+(?:\.[0-9A-Za-z-]+)*))?(?:\+([0-9A-Za-z-]+(?:\.[0-9A-Za-z-]+)*))?$/; +const NUMERIC = /^(0|[1-9]\d*)$/; +const INSTALL_TIMEOUT_SECONDS = 30; +const gray = (line: string) => `\x1b[90m${line}\x1b[0m\n`; +const reportToStderr = (line: string) => process.stderr.write(line); + +// This process owns the install result when the invoking CLI has already exited. +// It opens only the original terminal (never the parent's pipe or stdout), and +// checks its identity before writing in case the pty path has been recycled. +const TTY_HELPER = ` +const { openSync, closeSync, fstatSync, writeSync, constants } = require('node:fs'); +const [version, seconds, path, dev, ino, rdev, uid] = process.argv.slice(1); +let success = false; +try { + const child = Bun.spawn(['timeout', '--signal=TERM', '--kill-after=2s', seconds + 's', + 'bun', 'i', '-g', '--no-cache', '@dmgnr/kuber@' + version], + { stdin: 'ignore', stdout: 'ignore', stderr: 'ignore' }); + success = (await child.exited) === 0; +} catch {} +try { + const fd = openSync(path, constants.O_WRONLY | constants.O_NOFOLLOW | constants.O_NONBLOCK); + try { + const stat = fstatSync(fd); + if (stat.isCharacterDevice() && String(stat.dev) === dev && + String(stat.ino) === ino && String(stat.rdev) === rdev && String(stat.uid) === uid) { + writeSync(fd, '\\x1b[90m+ ' + (success ? 'Updated to ' : 'New version available: ') + + version + '\\x1b[0m\\n'); + } + } finally { closeSync(fd); } +} catch {} +`; + +function originalTerminal(): string[] | undefined { + if (!process.stderr.isTTY) return; + try { + const path = readlinkSync("/proc/self/fd/2"); + if (!/^\/dev\/pts\/[0-9]+$/.test(path)) return; + const stat = statSync(path); + if (!stat.isCharacterDevice()) return; + return [ + path, + String(stat.dev), + String(stat.ino), + String(stat.rdev), + String(stat.uid), + ]; + } catch { + return; + } +} + +function installWithTerminalReport(version: string, tty: string[]): void { + const child = Bun.spawn( + [ + process.execPath, + "-e", + TTY_HELPER, + version, + String(INSTALL_TIMEOUT_SECONDS), + ...tty, + ], + { stdin: "ignore", stdout: "ignore", stderr: "ignore", detached: true }, + ); + child.unref(); +} + +function parseVersion(value: string | null): SemVer | undefined { + if (!value || value.length > 128) return; + const match = SEMVER.exec(value); + if (!match || match[0] !== value) return; + const prerelease = match[4]?.split(".") ?? []; + if (prerelease.some((part) => /^\d+$/.test(part) && !NUMERIC.test(part))) + return; + return { + major: BigInt(match[1]!), + minor: BigInt(match[2]!), + patch: BigInt(match[3]!), + prerelease, + }; +} + +/** A positive result means candidate is newer; build metadata has no precedence. */ +export function compareVersions( + candidate: string, + current: string, +): number | undefined { + const a = parseVersion(candidate); + const b = parseVersion(current); + if (!a || !b) return; + for (const field of ["major", "minor", "patch"] as const) { + if (a[field] !== b[field]) return a[field] > b[field] ? 1 : -1; + } + if (!a.prerelease.length || !b.prerelease.length) { + return Number(!a.prerelease.length) - Number(!b.prerelease.length); + } + for (let i = 0; i < Math.max(a.prerelease.length, b.prerelease.length); i++) { + const left = a.prerelease[i]; + const right = b.prerelease[i]; + if (left === undefined || right === undefined) + return left === undefined ? -1 : 1; + if (left === right) continue; + const leftNumeric = NUMERIC.test(left); + const rightNumeric = NUMERIC.test(right); + if (leftNumeric && rightNumeric) + return BigInt(left) > BigInt(right) ? 1 : -1; + if (leftNumeric !== rightNumeric) return leftNumeric ? -1 : 1; + return left < right ? -1 : 1; + } + return 0; +} + +export type VersionInstallRunner = (version: string) => Promise; + +type InstallProcess = { + exited: Promise; + unref?(): void; +}; +type InstallSpawn = ( + argv: string[], + options: { + stdin: "ignore"; + stdout: "ignore"; + stderr: "ignore"; + detached: true; + }, +) => InstallProcess; + +export async function installVersion( + version: string, + spawn: InstallSpawn = Bun.spawn, + timeoutSeconds = INSTALL_TIMEOUT_SECONDS, +): Promise { + // Defense in depth: never pass an unvalidated header to a subprocess. + if ( + !parseVersion(version) || + !Number.isSafeInteger(timeoutSeconds) || + timeoutSeconds < 1 || + timeoutSeconds > 300 + ) + return false; + const child = spawn( + [ + "timeout", + "--signal=TERM", + "--kill-after=2s", + `${timeoutSeconds}s`, + "bun", + "i", + "-g", + "--no-cache", + `@dmgnr/kuber@${version}`, + ], + { stdin: "ignore", stdout: "ignore", stderr: "ignore", detached: true }, + ); + // The detached timeout owns its bounded lifetime; it must not keep a short + // CLI invocation alive. Its result can still be observed while the parent lives. + child.unref?.(); + return child.exited.then((code) => code === 0); +} + +export function createVersionObserver({ + currentVersion = KUBER_VERSION, + runner = installVersion, + report = reportToStderr, +}: { + currentVersion?: string; + runner?: VersionInstallRunner; + report?: (line: string) => void; +} = {}): (version: string | null) => void { + let attempted = false; + return (version) => { + if (attempted || !version || compareVersions(version, currentVersion) !== 1) + return; + attempted = true; + // Defer install work beyond the response headers; never block body consumption. + setTimeout(() => { + // The global observer must not install packages in tests or source-tree + // development commands. Explicitly injected runners remain testable. + if ( + runner === installVersion && + (process.env.NODE_ENV === "test" || + process.env.NODE_ENV === "development" || + process.argv[1]?.endsWith(".ts")) + ) + return; + if (runner === installVersion && report === reportToStderr) { + const tty = originalTerminal(); + if (tty) { + try { + installWithTerminalReport(version, tty); + } catch { + try { + report(gray(`+ New version available: ${version}`)); + } catch {} + } + return; + } + } + void Promise.resolve() + .then(() => runner(version)) + .then( + (success) => { + report( + gray( + `+ ${success ? "Updated to " : "New version available: "}${version}`, + ), + ); + }, + () => report(gray(`+ New version available: ${version}`)), + ) + .catch(() => { + // A broken stderr must never affect an API request or command exit status. + }); + }, 0); + }; +} + +export const observeServerVersion = createVersionObserver(); diff --git a/package.json b/package.json index 31c6860..810ca85 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.6.0", + "version": "2.6.1-rc5", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/server/app.ts b/server/app.ts index e4878e0..7df0a36 100644 --- a/server/app.ts +++ b/server/app.ts @@ -1,9 +1,10 @@ import type { KubernetesObject } from "@kubernetes/client-node"; import { randomUUID } from "node:crypto"; import type { ComposeSpecification } from "../schema/docker.d"; -import { DatabaseReconciliationError } from "../lib/database"; +import { DATABASE_RECONCILE_PHASES, DatabaseReconciliationError } from "../lib/database"; import type { Operation as PublicOperation } from "../shared/api"; import type { BuildRequest, Sha256Digest } from "../shared/build-protocol"; +import { KUBER_VERSION, KUBER_VERSION_HEADER } from "../shared/version"; import { createToken, hashToken, @@ -28,6 +29,7 @@ import { BuildValidationError, type BuildController, } from "./build-controller"; +import { BuildEventStreamHub } from "./build-event-stream"; import { KubernetesLogError, type LogService } from "./log-service"; import { ExecService, @@ -142,6 +144,11 @@ export type UnknownFailureLog = { export type RequestErrorLog = Omit & { event: string; operationId?: string; + phase?: string; + errorClass?: string; + providerStatus?: number; + providerCode?: string; + topFrame?: string; }; export interface AppLogger { @@ -200,6 +207,62 @@ class HttpError extends Error { } } +const DATABASE_PHASES = new Set(DATABASE_RECONCILE_PHASES); +const DATABASE_ERROR_CLASSES = new Set([ + "Error", "TypeError", "RangeError", "SyntaxError", "AbortError", + "ApiException", "ResponseError", "FetchError", "TimeoutError", +]); +const DATABASE_ERROR_CODES = new Set([ + "ECONNREFUSED", "ECONNRESET", "ETIMEDOUT", "EHOSTUNREACH", + "ENOTFOUND", "EAI_AGAIN", "ABORT_ERR", "UND_ERR_CONNECT_TIMEOUT", +]); + +/** Never serialize an unknown error: its message, stack and provider fields may contain secrets. */ +function databaseFailureDiagnostics(error: unknown): Record { + let current = error; + let phase = "operation execution"; + const seen = new Set(); + for ( + let depth = 0; + depth < 5 && current && typeof current === "object" && !seen.has(current); + depth++ + ) { + seen.add(current); + if ( + current instanceof DatabaseReconciliationError && + DATABASE_PHASES.has(current.phase) + ) { + phase = current.phase; + } + const cause = (current as { cause?: unknown }).cause; + if (!cause || typeof cause !== "object" || seen.has(cause)) break; + current = cause; + } + const source = current && typeof current === "object" + ? current as Record : {}; + const errorClass = current instanceof Error && DATABASE_ERROR_CLASSES.has(current.name) + ? current.name + : "UnknownError"; + const status = [source.statusCode, source.status, source.code].find( + (value) => typeof value === "number" && Number.isInteger(value) && + value >= 400 && value < 600, + ); + const errorCode = typeof source.code === "string" && DATABASE_ERROR_CODES.has(source.code) + ? source.code + : undefined; + // Only the first frame and a fixed set of our source files; never log raw stacks or paths. + const firstFrame = current instanceof Error ? current.stack?.split("\n")[1] : undefined; + const frame = firstFrame && + /(?:^|\/)\b((?:lib\/database|server\/management|server\/app)\.ts):(\d+):(\d+)\b/.exec(firstFrame); + return { + phase, + errorClass, + ...(status !== undefined && { providerStatus: status }), + ...(errorCode && { providerCode: errorCode }), + ...(frame && { topFrame: `${frame[1]}:${frame[2]}:${frame[3]}` }), + }; +} + function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } @@ -282,6 +345,9 @@ export function createApp( const now = options.now ?? Date.now; const makeRequestId = options.requestId ?? randomUUID; const logger = options.logger ?? defaultAppLogger; + const buildEventStreams = options.builds + ? new BuildEventStreamHub(options.builds) + : undefined; const bodyLimit = options.jsonBodyLimit ?? DEFAULT_JSON_LIMIT; const allowedOrigins = new Set(options.allowedOrigins ?? []); const hasOriginConfiguration = options.allowedOrigins !== undefined; @@ -1150,25 +1216,43 @@ export function createApp( }); } catch (error) { if (!operationStarted) throw error; - const message = error instanceof Error ? error.message : String(error); const leaseLost = leaseOwnershipLost || (error instanceof HttpError && error.code === "WORKSPACE_LEASE_LOST"); + const databaseFailure = + action === "databases.reconcile" && !leaseLost + ? error instanceof DatabaseReconciliationError + ? error + : new DatabaseReconciliationError("operation execution", error) + : undefined; + if (databaseFailure) { + logger.error({ + event: "operation.database_reconcile.failed", + requestId: requestIds.get(request) ?? makeRequestId(), + method: "POST", + pathname: "/api/v2/workspaces/:workspaceId/databases", + operationId: operation.metadata.name, + status: 500, + code: "DATABASE_RECONCILE_FAILED", + errorName: "DatabaseReconciliationError", + message: "Database reconciliation failed", + ...databaseFailureDiagnostics(databaseFailure), + }); + } + const message = + databaseFailure?.message ?? + (error instanceof Error ? error.message : String(error)); const failed = await transitionOperationToFailure( operation.metadata.name, { code: leaseLost ? "WORKSPACE_LEASE_LOST" - : action === "databases.reconcile" && - error instanceof DatabaseReconciliationError + : databaseFailure ? "DATABASE_RECONCILE_FAILED" : "OPERATION_FAILED", message: leaseLost ? "Workspace operation lease ownership was lost" - : action === "databases.reconcile" && - !(error instanceof DatabaseReconciliationError) - ? "Database reconciliation failed" - : message, + : message, }, ); const failure = failed.status.error; @@ -1470,13 +1554,31 @@ export function createApp( ); if (!action && request.method === "GET") return response(await builds.getBuildStatus(id)); - if (action === "events" && request.method === "GET") + if (action === "events" && request.method === "GET") { + if (request.headers.get("accept")?.split(",").some((value) => + value.trim().split(";")[0]?.trim().toLowerCase() === "text/event-stream", + )) { + const cursor = request.headers.get("last-event-id") ?? url.searchParams.get("after"); + if ( + cursor !== null && + (!/^(0|[1-9]\d*)$/.test(cursor) || + !Number.isSafeInteger(Number(cursor))) + ) + throw new HttpError( + 400, + "Invalid cursor", + "INVALID_QUERY", + "Event cursor must be a non-negative safe integer", + ); + return buildEventStreams!.open(id, request); + } return response( await builds.getBuildEvents( id, nonNegativeInteger(url.searchParams.get("after"), "after") ?? 0, ), ); + } if (action === "reconcile" && request.method === "POST") return response(await builds.reconcileBuild(id)); if (action === "cancel" && request.method === "POST") @@ -2553,6 +2655,7 @@ export function createApp( result = problem(normalized, requestId); } result.headers.set("x-request-id", requestId); + result.headers.set(KUBER_VERSION_HEADER, KUBER_VERSION); if (origin && originAllowed) { result.headers.set("access-control-allow-origin", origin); result.headers.set("vary", "Origin"); @@ -2689,6 +2792,7 @@ export function execProblem(error: unknown, requestId = ""): Response { }, ); result.headers.set("content-type", "application/problem+json"); + result.headers.set(KUBER_VERSION_HEADER, KUBER_VERSION); return result; } diff --git a/server/build-event-stream.ts b/server/build-event-stream.ts new file mode 100644 index 0000000..2e12761 --- /dev/null +++ b/server/build-event-stream.ts @@ -0,0 +1,255 @@ +import type { BuildController } from "./build-controller"; +import type { BuildStatus } from "../shared/build-protocol"; + +const terminal = (status: BuildStatus) => + status.state === "succeeded" || status.state === "failed"; + +function cursor(value: string | null): number { + if (value === null) return 0; + if (!/^(0|[1-9]\d*)$/.test(value) || !Number.isSafeInteger(Number(value))) + throw new RangeError("Event cursor must be a non-negative safe integer"); + return Number(value); +} + +type Subscriber = { + sequence: number; + status?: string; + lastActivity: number; + controller: ReadableStreamDefaultController; + stop: () => void; +}; + +type Group = { + members: Set; + abort: AbortController; + nextReconcileAt: number; + reconciling: boolean; +}; + +/** One hub per app/controller instance; call only after the route's capability and workspace checks. + * The authenticated GET /api/v2/builds/:id/events route should branch on + * Accept: text/event-stream, then return hub.open(id, request); retain JSON otherwise. + * A hub coalesces reconciliation and persisted event reads per physical build on this + * replica. BuildController's lease coordinates reconciliation across replicas; its + * persisted event store supplies logs even when a different replica owns the job. + */ +export class BuildEventStreamHub { + private readonly builds = new Map(); + + constructor( + private readonly controller: Pick< + BuildController, + "getBuildStatus" | "getBuildEvents" | "reconcileBuild" + >, + private readonly options: { + pollMs?: number; + heartbeatMs?: number; + reconcileMs?: number; + } = {}, + ) {} + + async open(id: string, request: Request): Promise { + const header = request.headers.get("last-event-id"); + const query = new URL(request.url).searchParams.get("after"); + const after = cursor(header ?? query); + // Validate existence before committing response headers. + await this.controller.getBuildStatus(id); + let subscriber: Subscriber; + const stream = new ReadableStream( + { + start: (controller) => { + let group: Group; + const stop = () => { + request.signal.removeEventListener("abort", stop); + group.members.delete(subscriber); + if (group.members.size === 0 && this.builds.get(id) === group) { + this.builds.delete(id); + group.abort.abort(); + } + try { + controller.close(); + } catch { + /* Already closed or cancelled. */ + } + }; + subscriber = { + sequence: after, + lastActivity: Date.now(), + controller, + stop, + }; + const existing = this.builds.get(id); + if (existing) { + group = existing; + group.members.add(subscriber); + } else { + group = { + members: new Set([subscriber]), + abort: new AbortController(), + nextReconcileAt: 0, + reconciling: false, + }; + this.builds.set(id, group); + void this.run(id, group); + } + request.signal.addEventListener("abort", stop, { once: true }); + if (request.signal.aborted) stop(); + }, + cancel: () => subscriber.stop(), + }, + { highWaterMark: 16 }, + ); + return new Response(stream, { + headers: { + "content-type": "text/event-stream; charset=utf-8", + "cache-control": "no-cache, no-transform", + connection: "keep-alive", + "x-accel-buffering": "no", + }, + }); + } + + private async run(id: string, group: Group): Promise { + const encoder = new TextEncoder(); + const pollMs = this.options.pollMs ?? 1_000; + const reconcileMs = this.options.reconcileMs ?? 7_000; + const { members } = group; + // Bun's default HTTP idle timeout is 10 seconds; keep bytes flowing well + // inside it even while Kubernetes reconciliation makes no progress. + const heartbeatMs = this.options.heartbeatMs ?? 5_000; + const heartbeat = setInterval(() => { + if (group.abort.signal.aborted) return; + for (const member of members) { + if (Date.now() - member.lastActivity < heartbeatMs) continue; + if ((member.controller.desiredSize ?? 0) <= 0) member.stop(); + else { + member.controller.enqueue(encoder.encode(": heartbeat\n\n")); + member.lastActivity = Date.now(); + } + } + }, heartbeatMs); + heartbeat.unref?.(); + try { + while (!group.abort.signal.aborted) { + try { + const initial = await this.controller.getBuildStatus(id); + if (group.abort.signal.aborted) break; + // Reconciliation is shared per build and never blocks persisted progress. + // The controller's lease still coordinates attempts across replicas. + if ( + !terminal(initial) && + !group.reconciling && + Date.now() >= group.nextReconcileAt + ) { + group.reconciling = true; + group.nextReconcileAt = Date.now() + reconcileMs; + const signal = group.abort.signal; + let onAbort = () => {}; + const aborted = new Promise((resolve) => { + onAbort = resolve; + signal.addEventListener("abort", onAbort, { once: true }); + }); + const reconciliation = this.controller + .reconcileBuild(id, { signal }) + .catch(() => { + // Another replica may hold the lease; keep observing its progress. + }); + // Drop this group's references even if an external API ignores abort + // and leaves the reconciliation promise unresolved indefinitely. + void Promise.race([reconciliation, aborted]).finally(() => { + signal.removeEventListener("abort", onAbort); + group.reconciling = false; + }); + } + const status = await this.controller.getBuildStatus(id); + if (group.abort.signal.aborted) break; + const oldest = Math.min( + ...[...members].map((member) => member.sequence), + ); + const events = await this.controller.getBuildEvents(id, oldest); + if (group.abort.signal.aborted) break; + const firstLogSequence = events.find( + (event) => event.type === "log", + )?.sequence; + for (const member of members) { + const send = (frame: string) => { + // A slow consumer reconnects from its last delivered sequence instead + // of retaining an unbounded in-memory backlog on this replica. + if ((member.controller.desiredSize ?? 0) <= 0) { + member.stop(); + return false; + } + member.controller.enqueue(encoder.encode(frame)); + member.lastActivity = Date.now(); + return true; + }; + if ( + firstLogSequence !== undefined && + firstLogSequence > member.sequence + 1 + ) { + const gap = { + type: "gap", + after: member.sequence, + before: firstLogSequence, + missing: firstLogSequence - member.sequence - 1, + message: `Build log history was trimmed; ${firstLogSequence - member.sequence - 1} log event(s) before sequence ${firstLogSequence} are unavailable.`, + }; + if ( + !send( + `id: ${firstLogSequence - 1}\nevent: gap\ndata: ${JSON.stringify(gap)}\n\n`, + ) + ) + continue; + // Advance over the unavailable range so reconnects do not repeat + // the marker. The cursor now refers to the last unavailable + // sequence; delivered log records retain their normal IDs. + member.sequence = firstLogSequence - 1; + } + for (const event of events) { + if (event.type !== "log" || event.sequence <= member.sequence) + continue; + if ( + !send( + `id: ${event.sequence}\nevent: log\ndata: ${JSON.stringify(event)}\n\n`, + ) + ) + break; + member.sequence = event.sequence; + } + if (!members.has(member)) continue; + const fingerprint = JSON.stringify(status); + if (member.status !== fingerprint) { + if ( + !send( + `event: status\ndata: ${JSON.stringify({ type: "status", status })}\n\n`, + ) + ) + continue; + member.status = fingerprint; + } + if (terminal(status)) member.stop(); + } + } catch { + // A transient API/store failure closes streams so clients reconnect with + // their last sequence; no silent terminal state is fabricated. + for (const member of members) member.stop(); + } + if (group.abort.signal.aborted) break; + await new Promise((resolve) => { + const stop = () => { + clearTimeout(timer); + resolve(); + }; + const timer = setTimeout(() => { + group.abort.signal.removeEventListener("abort", stop); + resolve(); + }, pollMs); + group.abort.signal.addEventListener("abort", stop, { once: true }); + }); + } + } finally { + clearInterval(heartbeat); + group.abort.abort(); + } + } +} diff --git a/server/management.ts b/server/management.ts index 623f0cb..1be4f64 100644 --- a/server/management.ts +++ b/server/management.ts @@ -8,7 +8,9 @@ import { sortResources, } from "../lib/apply"; import { + DATABASE_CLUSTER, DATABASE_NAMESPACE, + DatabaseReconciliationError, getComposePostgresClaims, getRoleCredentials, listManagedDatabaseResources, @@ -494,8 +496,15 @@ export function createManagementService(dependencies: ManagementDependencies) { ); if (full) { + // Database CRs may have consumers in other projects. Labels are not a + // reference count, including when this workspace created the CR. + const databases = await dependencies.listDatabaseResources(workspace.project); + retainedResources.push(...databases.filter( + (resource) => resource.kind === "Database" && + resource.metadata?.labels?.[WORKSPACE_UID_LABEL] === workspace.uid, + )); deleteResources.push( - ...(await dependencies.listDatabaseResources(workspace.project)), + ...databases.filter((resource) => resource.kind !== "Database"), ...(await dependencies.listStorageResources(workspace.project)), ); if (!safety.namespaceUid) { @@ -640,19 +649,56 @@ export function createManagementService(dependencies: ManagementDependencies) { compose: ComposeSpecification, execution?: OperationExecution, ) { - throwIfExecutionAborted(execution); - await assertSafe(workspace, true); - const environment = await dependencies.reconcileDatabases( - workspace.project, - compose, - execution, - ); - await ownExternalResources( - workspace, - DATABASE_NAMESPACE, - await dependencies.listDatabaseResources(workspace.project), - execution, - ); + try { + throwIfExecutionAborted(execution); + await assertSafe(workspace, true); + } catch (error) { + throw new DatabaseReconciliationError("namespace precheck", error); + } + let environment: Record>; + try { + environment = await dependencies.reconcileDatabases( + workspace.project, + compose, + execution, + ); + } catch (error) { + if (error instanceof DatabaseReconciliationError) throw error; + throw new DatabaseReconciliationError("database dependency", error); + } + let resources: KubernetesObject[]; + try { + resources = await dependencies.listDatabaseResources(workspace.project); + } catch (error) { + throw new DatabaseReconciliationError("database resource listing", error); + } + try { + const claims = getComposePostgresClaims(compose); + const claimed = new Set(claims.map((claim) => claim.database)); + for (const resource of resources) { + // Only a declared database may be shared. Never adopt its foreign + // workspace UID or apply its metadata through this field manager. + const existingOwner = resource.metadata?.labels?.[WORKSPACE_UID_LABEL]; + if (resource.kind === "Database" && existingOwner && + existingOwner !== workspace.uid && + claimed.has(resource.metadata?.name ?? "")) { + if (resource.metadata?.namespace !== DATABASE_NAMESPACE) { + throw new Error("Database is outside expected namespace"); + } + const claim = claims.find((item) => item.database === resource.metadata?.name); + const spec = (resource as KubernetesObject & { + spec?: { owner?: string; cluster?: { name?: string } }; + }).spec; + if (spec?.owner !== claim?.username || spec?.cluster?.name !== DATABASE_CLUSTER) { + throw new Error("Database does not match the declared claim"); + } + continue; + } + await ownExternalResources(workspace, DATABASE_NAMESPACE, [resource], execution); + } + } catch (error) { + throw new DatabaseReconciliationError("database resource ownership", error); + } return environment; }, diff --git a/shared/version.ts b/shared/version.ts new file mode 100644 index 0000000..1a5a3e5 --- /dev/null +++ b/shared/version.ts @@ -0,0 +1,5 @@ +import packageMetadata from "../package.json"; + +// Bun bundles JSON imports into the CLI and server artifacts at build time. +export const KUBER_VERSION: string = packageMetadata.version; +export const KUBER_VERSION_HEADER = "x-kuber-version"; diff --git a/tests/command/up-api.test.ts b/tests/command/up-api.test.ts index 0ec21d4..9301f21 100644 --- a/tests/command/up-api.test.ts +++ b/tests/command/up-api.test.ts @@ -4,6 +4,7 @@ import { mkdtemp, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { Writable } from "node:stream"; +import { stripVTControlCharacters } from "node:util"; import type { ApiRequestInit, ApiRequestOptions } from "../../lib/api"; import { KuberApiError } from "../../lib/api"; import type { ApiRequester } from "../../lib/build"; @@ -25,8 +26,21 @@ const snapshot = { blobs: [], }; +// DefaultRenderer.create returns a complete TTY frame before the cursor redraw. +function frameRows(frame: string): string[] { + return stripVTControlCharacters(frame).split("\n"); +} + +function taskRow(rows: string[], title: string): string { + return rows.find((row) => row.includes(title)) ?? ""; +} + +function depth(row: string): number { + return row.match(/^\s*/)?.[0].length ?? 0; +} + describe("up API pipeline", () => { - test("keeps queued builds in the TTY render tree without duplicate rows", async () => { + test("renders build lifecycle in child titles without duplicate phase output", async () => { const root = await mkdtemp(join(tmpdir(), "kuber-up-api-")); const previousCwd = process.cwd(); const tty = Object.getOwnPropertyDescriptor(process.stdout, "isTTY"); @@ -37,9 +51,7 @@ describe("up API pipeline", () => { frames.push(frame); return frame; }); - const writes = spyOn(process.stdout, "write").mockImplementation(((chunk: string | Uint8Array) => { - return true; - }) as typeof process.stdout.write); + const writes = spyOn(process.stdout, "write").mockImplementation((() => true) as typeof process.stdout.write); const polls = new Map(); const builds = new Map(); try { @@ -60,13 +72,13 @@ describe("up API pipeline", () => { if (path.includes("/events")) return [{ type: "status", status: { state: "queued", - phase: ["queued", "preparing", "waiting", "waiting"][polls.get(service) ?? 0], + phase: ["queued", "creating", "starting", "running", "done"][polls.get(service) ?? 0], } }] as T; if (path.endsWith("/reconcile")) { const count = (polls.get(service) ?? 0) + 1; polls.set(service, count); await Bun.sleep(110); - return count < 3 + return count < 4 ? { state: "queued" } as T : service === "client" ? { state: "failed", error: "build stopped\nstack detail" } as T @@ -76,29 +88,27 @@ describe("up API pipeline", () => { throw new Error(path); }; await expect(provideContext(() => runUp(true, request, { trust }))).rejects.toThrow("build stopped"); - expect(polls.get("client")).toBe(3); - const active = frames.filter((frame) => - frame.includes("Build images") && frame.includes("Build app") && frame.includes("Build client") && - frame.split("\n").filter((row) => row.includes("Build queued")).length === 2); + expect(polls.get("client")).toBe(4); + const active = frames.map(frameRows).filter((rows) => + taskRow(rows, "Build images") && taskRow(rows, "Build app") && taskRow(rows, "Build client")); expect(active.length).toBeGreaterThan(0); - for (const frame of frames.filter((frame) => frame.includes("Build images") && frame.includes("Build app") && frame.includes("Build client"))) { - const rows = frame.split("\n"); - const parent = rows.find((row) => row.includes("Build images"))!; + for (const rows of active) { + const parent = taskRow(rows, "Build images"); for (const name of ["app", "client"]) { const children = rows.filter((row) => row.includes(`Build ${name}`)); expect(children).toHaveLength(1); - expect(children[0]!.search(/\S/)).toBeGreaterThan(parent.search(/\S/)); + expect(depth(children[0]!)).toBeGreaterThan(depth(parent)); } - expect(rows.filter((row) => row.includes("Build queued")).length).toBeLessThanOrEqual(2); + expect(rows.filter((row) => /^\s*› Build (?:queued|creating|starting|running|done)/.test(row))).toHaveLength(0); } - expect(active[0]).toMatch(/Build app\n {2,}.*Build queued[\s\S]*Build client\n {2,}.*Build queued/); - const final = frames.at(-1)!; - expect(final).toMatch(/Build images[\s\S]* .*Build app[\s\S]* .*Build client/); - expect(final).toMatch(/✔.*Build app/); - expect(final).toMatch(/✖.*Build client/); - expect(frames.some((frame) => frame.includes("Build preparing") && frame.includes("Build waiting"))).toBe(true); - expect(final).toContain("build stopped"); - expect(final).toContain("stack detail"); + for (const phase of ["queued", "creating", "starting", "running", "done"]) + expect(frames.some((frame) => frame.includes(`Build app: ${phase}`))).toBe(true); + const final = frameRows(frames.at(-1)!); + expect(taskRow(final, "Build app: done")).toMatch(/✔/); + expect(taskRow(final, "Build client: failed")).toMatch(/✖/); + expect(final.join("\n")).toContain("build stopped"); + expect(final.join("\n")).toContain("stack detail"); + expect(final.join("\n").split("build stopped")).toHaveLength(2); } finally { renderSpy.mockRestore(); writes.mockRestore(); @@ -292,11 +302,18 @@ describe("up API pipeline", () => { const order: string[] = []; const started: string[] = []; let rootTasks: Listr["tasks"] | undefined; - let backingTasks: Listr["tasks"] | undefined; const originalRun = Listr.prototype.run; + const frames: string[] = []; + const create = DefaultRenderer.prototype.create; + const renderSpy = spyOn(DefaultRenderer.prototype, "create").mockImplementation(function (this: DefaultRenderer, options) { + const frame = create.call(this, options); + frames.push(frame); + return frame; + }); + const tty = Object.getOwnPropertyDescriptor(process.stdout, "isTTY"); + const writes = spyOn(process.stdout, "write").mockImplementation((() => true) as typeof process.stdout.write); const runSpy = spyOn(Listr.prototype, "run").mockImplementation(function (this: Listr) { if (this.tasks[0]?.title === "Read compose") rootTasks = this.tasks; - if (this.tasks[0]?.title === "Reconcile databases") backingTasks = this.tasks; return originalRun.call(this); }); let bothReady!: () => void; @@ -312,6 +329,7 @@ describe("up API pipeline", () => { releaseStorage = resolve; }); try { + Object.defineProperty(process.stdout, "isTTY", { configurable: true, value: true }); await writeFile( join(root, "compose.yml"), "services:\n app:\n image: nginx\n volumes:\n - postgresql:app\n - s3:assets\n", @@ -356,10 +374,21 @@ describe("up API pipeline", () => { expect(titles.indexOf("Render manifests")).toBeLessThan( titles.indexOf("Reconcile resources"), ); - expect(backingTasks!.map(({ title }) => title)).toEqual([ + const backing = rootTasks!.find(({ title }) => title === "Reconcile backing services")!; + expect(backing.subtasks.map(({ title }) => title)).toEqual([ "Reconcile databases", "Reconcile S3 storage", ]); + const active = frames.map(frameRows).filter((rows) => + taskRow(rows, "Reconcile backing services") && + taskRow(rows, "Reconcile databases") && + taskRow(rows, "Reconcile S3 storage")); + expect(active.length).toBeGreaterThan(0); + for (const rows of active) { + const parent = taskRow(rows, "Reconcile backing services"); + expect(depth(taskRow(rows, "Reconcile databases"))).toBeGreaterThan(depth(parent)); + expect(depth(taskRow(rows, "Reconcile S3 storage"))).toBeGreaterThan(depth(parent)); + } expect(rootTasks!.filter(({ title }) => title === "Reconcile databases")).toEqual([]); expect(started).toEqual(["database", "storage"]); expect(order.indexOf("/workspaces/shop/adopt")).toBeLessThan( @@ -388,6 +417,10 @@ describe("up API pipeline", () => { } } finally { runSpy.mockRestore(); + renderSpy.mockRestore(); + writes.mockRestore(); + if (tty) Object.defineProperty(process.stdout, "isTTY", tty); + else Reflect.deleteProperty(process.stdout, "isTTY"); process.chdir(previousCwd); await rm(root, { recursive: true, force: true }); } @@ -396,17 +429,18 @@ describe("up API pipeline", () => { test("shows database API problem details on the database task without starting resources", async () => { const root = await mkdtemp(join(tmpdir(), "kuber-up-api-")); const previousCwd = process.cwd(); - const rendered: string[] = []; - const writes = spyOn(process.stdout, "write").mockImplementation(((chunk: string | Uint8Array) => { - rendered.push(String(chunk)); - return true; - }) as typeof process.stdout.write); - const errors = spyOn(process.stderr, "write").mockImplementation(((chunk: string | Uint8Array) => { - rendered.push(String(chunk)); - return true; - }) as typeof process.stderr.write); + const frames: string[] = []; + const create = DefaultRenderer.prototype.create; + const renderSpy = spyOn(DefaultRenderer.prototype, "create").mockImplementation(function (this: DefaultRenderer, options) { + const frame = create.call(this, options); + frames.push(frame); + return frame; + }); + const tty = Object.getOwnPropertyDescriptor(process.stdout, "isTTY"); + const writes = spyOn(process.stdout, "write").mockImplementation((() => true) as typeof process.stdout.write); const calls: string[] = []; try { + Object.defineProperty(process.stdout, "isTTY", { configurable: true, value: true }); await writeFile(join(root, "compose.yml"), "services:\n web:\n image: nginx\n volumes:\n - postgresql:web_db\n"); await writeFile(join(root, ".kuberrc.ts"), 'export default { project: "shop" };\n'); process.chdir(root); @@ -423,11 +457,17 @@ describe("up API pipeline", () => { throw new Error(`Unexpected request: ${path}`); }; await expect(provideContext(() => runUp(false, request, { trust }))).rejects.toThrow(detail); - expect(rendered.join("")).toContain("database apply for database web_db"); + const rows = frameRows(frames.at(-1)!); + expect(depth(taskRow(rows, "Reconcile databases"))).toBeGreaterThan(depth(taskRow(rows, "Reconcile backing services"))); + expect(taskRow(rows, "Reconcile databases")).toMatch(/✖/); + const visible = rows.join(" ").replace(/\s+/g, " "); + expect(visible.split(detail)).toHaveLength(2); expect(calls).not.toContain("/workspaces/shop/resources/plan"); } finally { - errors.mockRestore(); + renderSpy.mockRestore(); writes.mockRestore(); + if (tty) Object.defineProperty(process.stdout, "isTTY", tty); + else Reflect.deleteProperty(process.stdout, "isTTY"); process.chdir(previousCwd); await rm(root, { recursive: true, force: true }); } diff --git a/tests/lib/api-stream.test.ts b/tests/lib/api-stream.test.ts new file mode 100644 index 0000000..49433d6 --- /dev/null +++ b/tests/lib/api-stream.test.ts @@ -0,0 +1,67 @@ +import { afterEach, expect, test } from "bun:test"; +import { apiStreamEvents, ApiStreamUnsupportedError } from "../../lib/api"; + +const originalFetch = globalThis.fetch; +afterEach(() => { + globalThis.fetch = originalFetch; +}); + +test("parses split UTF-8 SSE frames, comments and ids with authenticated request", async () => { + let headers: Headers | undefined; + const versions: Array = []; + globalThis.fetch = (async ( + _url: string | URL | Request, + init?: RequestInit, + ) => { + headers = new Headers(init?.headers); + const bytes = new TextEncoder().encode( + ': ping\r\nid: 3\r\nevent: log\r\ndata: {"type":"log","message":"hé"}\r\n\r\nevent: status\ndata: {"type":"status"}\n\n', + ); + return new Response( + new ReadableStream({ + start(controller) { + for (let i = 0; i < bytes.length; i += 2) + controller.enqueue(bytes.slice(i, i + 2)); + controller.close(); + }, + }), + { + headers: { + "content-type": "text/event-stream", + "x-kuber-version": "2.6.1-rc5", + }, + }, + ); + }) as unknown as typeof fetch; + const events = []; + for await (const event of apiStreamEvents("/builds/a/events", 2, undefined, { + baseUrl: "https://test", + session: { token: "secret" } as never, + onServerVersion: (value) => versions.push(value), + })) + events.push(event); + expect(headers?.get("authorization")).toBe("Bearer secret"); + expect(headers?.get("last-event-id")).toBe("2"); + expect(versions).toEqual(["2.6.1-rc5"]); + expect(events).toEqual([ + { id: 3, event: { type: "log", message: "hé" } }, + { id: undefined, event: { type: "status" } }, + ]); +}); + +test("identifies existing JSON-only servers without parsing their response as SSE", async () => { + globalThis.fetch = (async () => + new Response("[]", { + headers: { "content-type": "application/json" }, + })) as unknown as typeof fetch; + await expect( + (async () => { + for await (const _ of apiStreamEvents("/builds/a/events", 0, undefined, { + baseUrl: "https://test", + session: { token: "secret" } as never, + })) { + /* consume */ + } + })(), + ).rejects.toBeInstanceOf(ApiStreamUnsupportedError); +}); diff --git a/tests/lib/api.test.ts b/tests/lib/api.test.ts index b8071a6..609e67c 100644 --- a/tests/lib/api.test.ts +++ b/tests/lib/api.test.ts @@ -7,6 +7,7 @@ import { apiUpload, } from "../../lib/api"; import type { KuberSession } from "../../lib/session"; +import { createVersionObserver } from "../../lib/version-update"; const originalFetch = globalThis.fetch; const session: KuberSession = { @@ -19,6 +20,73 @@ afterEach(() => { globalThis.fetch = originalFetch; }); +test("observes the response header on JSON, errors, uploads and streaming without delaying payloads", async () => { + const seen: Array = []; + const onServerVersion = (version: string | null) => { + seen.push(version); + }; + const options = { + authenticated: false, + baseUrl: "https://api.test", + onServerVersion, + }; + const encoder = new TextEncoder(); + globalThis.fetch = mock(async (input) => { + const path = new URL(String(input)).pathname; + const headers = { "x-kuber-version": "2.7.0" }; + if (path === "/error") + return Response.json({ detail: "failure" }, { status: 503, headers }); + if (path === "/upload") return new Response(null, { status: 204, headers }); + if (path === "/stream") + return new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode('{"n":1}\n')); + }, + }), + { headers }, + ); + return Response.json({ ok: true }, { headers }); + }) as unknown as typeof fetch; + + expect(await apiRequest<{ ok: boolean }>("/json", {}, options)).toEqual({ + ok: true, + }); + await expect(apiRequest("/error", {}, options)).rejects.toBeInstanceOf( + KuberApiError, + ); + await apiUpload("/upload", new Uint8Array([1]), { ...options, offset: 0 }); + + let complete!: (success: boolean) => void; + const pending = new Promise((resolve) => { + complete = resolve; + }); + const lines: string[] = []; + const observer = createVersionObserver({ + currentVersion: "2.6.1", + runner: () => pending, + report: (line) => lines.push(line), + }); + const records = apiStreamNdjson<{ n: number }>( + "/stream", + {}, + { + ...options, + onServerVersion: (version) => { + onServerVersion(version); + observer(version); + }, + }, + ); + expect((await records.next()).value).toEqual({ n: 1 }); + expect(seen).toEqual(["2.7.0", "2.7.0", "2.7.0", "2.7.0"]); + expect(lines).toEqual([]); + await records.return(); + complete(false); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(lines).toEqual(["\x1b[90m+ New version available: 2.7.0\x1b[0m\n"]); +}); + describe("authenticated API transport", () => { test("sends authentication and serializes JSON", async () => { const fetchMock = mock( diff --git a/tests/lib/build-api.test.ts b/tests/lib/build-api.test.ts index 17f5939..6aa4d90 100644 --- a/tests/lib/build-api.test.ts +++ b/tests/lib/build-api.test.ts @@ -18,6 +18,7 @@ import { type WorkspaceSnapshot, } from "../../lib/workspace"; import type { BuildRequest } from "../../shared/build-protocol"; +import { ApiStreamUnsupportedError } from "../../lib/api"; const directories: string[] = []; @@ -582,6 +583,88 @@ describe("authenticated build API pipeline", () => { expect(calls.some(({ path }) => path.endsWith("/result"))).toBe(true); }); + test("streams one physical build, resumes cursor, and never polls JSON while SSE is active", async () => { + const snapshot = emptySnapshot(); + const calls: string[] = []; + const cursors: number[] = []; + const output: string[] = []; + const request: ApiRequester = async (path: string) => { + calls.push(path); + if (path === "/snapshots/negotiate") return { ready: true } as T; + if (path === "/builds") return { state: "queued" } as T; + if (path.endsWith("/result")) return { reference: "image:web", references: { worker: "image:worker" } } as T; + throw new Error(`Unexpected JSON request ${path}`); + }; + const result = await buildServices("shop", { + services: { web: { build: "." }, worker: { build: "." } }, + }, process.cwd(), { stream: new Writable({ write(chunk, _, done) { + output.push(String(chunk)); done(); + } }) }, { request, snapshot, pollIntervalMs: 0, sleep: async () => {}, + streamEvents: async function* (_path, after) { + cursors.push(after); + if (after === 0) { + yield { id: 1, event: { type: "log", sequence: 1, id: "build", message: "once\n" } }; + yield { event: { type: "status", status: { version: 1, id: "build", state: "running", createdAt: "now" } } }; + } else { + yield { id: 1, event: { type: "log", sequence: 1, id: "build", message: "once\n" } }; + yield { event: { type: "status", status: { version: 1, id: "build", state: "succeeded", createdAt: "now" } } }; + } + }, + }); + expect(result.images).toEqual({ web: "image:web", worker: "image:worker" }); + expect(cursors).toEqual([0, 1]); + expect(output).toEqual(["[web] once\n"]); + expect(calls.filter((path) => path.includes("/events") || path.includes("/reconcile"))).toEqual([]); + }); + + test("falls back to JSON events/reconcile only after SSE is unsupported", async () => { + const calls: string[] = []; + const request: ApiRequester = async (path: string) => { + calls.push(path); + if (path === "/snapshots/negotiate") return { ready: true } as T; + if (path === "/builds") return { state: "queued" } as T; + if (path.includes("/events")) return [{ type: "log", sequence: 1, message: "fallback\n" }] as T; + if (path.endsWith("/reconcile")) return { state: "succeeded" } as T; + if (path.endsWith("/result")) return { reference: "image:web" } as T; + throw new Error(path); + }; + expect((await buildServices("shop", { services: { web: { build: "." } } }, + process.cwd(), undefined, { request, snapshot: emptySnapshot(), + streamEvents: async function* () { yield* []; throw new ApiStreamUnsupportedError(); }, + })).images).toEqual({ web: "image:web" }); + expect(calls.filter((path) => path.includes("/events") || path.includes("/reconcile"))).toEqual([ + expect.stringContaining("/events?after=0"), + expect.stringContaining("/reconcile"), + expect.stringContaining("/events?after=1"), + ]); + }); + + test("aborts an active stream and does not issue JSON polling requests", async () => { + const controller = new AbortController(); + const reason = new DOMException("Stopped", "AbortError"); + let listening!: () => void; + const ready = new Promise((resolve) => { listening = resolve; }); + const calls: string[] = []; + const request: ApiRequester = async (path: string) => { + calls.push(path); + if (path === "/snapshots/negotiate") return { ready: true } as T; + if (path === "/builds") return { state: "queued" } as T; + throw new Error(path); + }; + const result = buildServices("shop", { services: { web: { build: "." } } }, + process.cwd(), undefined, { request, snapshot: emptySnapshot(), signal: controller.signal, + streamEvents: async function* (_path, _after, signal) { + yield* []; + listening(); + await new Promise((_resolve, reject) => signal?.addEventListener("abort", () => reject(signal.reason), { once: true })); + }, + }); + await ready; + controller.abort(reason); + await expect(result).rejects.toBe(reason); + expect(calls).toEqual(["/snapshots/negotiate", "/builds"]); + }); + test("builds equivalent services once and attributes events and image results to each", async () => { const snapshot = emptySnapshot(); const submitted: BuildRequest[] = []; diff --git a/tests/lib/database.test.ts b/tests/lib/database.test.ts index 090245f..3bee14b 100644 --- a/tests/lib/database.test.ts +++ b/tests/lib/database.test.ts @@ -8,6 +8,7 @@ import { getServicePostgresClaim, isPostgresVolumeEntry, reconcilePostgresClaim, + reconcilePostgresClaims, } from "../../lib/database"; import { objectApi } from "../../lib/k8s"; @@ -127,6 +128,7 @@ describe("managed PostgreSQL claims", () => { secretName: "postgres-app", }; spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Database") throw { code: 404 }; if (resource.kind === "Secret") { return { ...resource, @@ -141,6 +143,7 @@ describe("managed PostgreSQL claims", () => { const patch = spyOn(objectApi, "patch").mockImplementation( async (resource) => resource as never, ); + const create = spyOn(objectApi, "create").mockImplementation(async (resource) => resource as never); expect(await reconcilePostgresClaim("project", claim)).toEqual({ username: "app", @@ -149,8 +152,8 @@ describe("managed PostgreSQL claims", () => { expect(patch.mock.calls.map(([resource]) => resource.kind)).toEqual([ "Secret", "Cluster", - "Database", ]); + expect(create.mock.calls.map(([resource]) => resource.kind)).toEqual(["Database"]); }); test("reports database apply phase and claim without exposing provider credentials", async () => { @@ -161,12 +164,14 @@ describe("managed PostgreSQL claims", () => { secretName: "postgres-app_role", }; spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Database") throw { code: 404 }; if (resource.kind === "Secret") { return { ...resource, data: { password: Buffer.from("private-value").toString("base64") } } as never; } return { ...resource, spec: { managed: { roles: [] } } } as never; }); - spyOn(objectApi, "patch").mockImplementation(async (resource) => { + spyOn(objectApi, "patch").mockImplementation(async (resource) => resource as never); + spyOn(objectApi, "create").mockImplementation(async (resource) => { if (resource.kind === "Database") { throw Object.assign(new Error("Forbidden: password=private-value"), { code: 403 }); } @@ -180,11 +185,78 @@ describe("managed PostgreSQL claims", () => { } expect(failure).toBeInstanceOf(DatabaseReconciliationError); expect((failure as Error).message).toBe( - "Database reconciliation failed during database apply for database app_db (service app, role app_role): Forbidden (HTTP 403)", + "Database reconciliation failed during database apply for requested database claim: Forbidden (HTTP 403)", ); expect((failure as Error).message).not.toContain("private-value"); }); + test.each(["read", "patch"])("distinguishes role secret %s failure", async (method) => { + const secretLike = "DB_PASSWORD_private123"; + spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Database") throw { code: 404 }; + if (method === "read") throw new Error(`token=${secretLike}`); + return undefined as never; + }); + spyOn(objectApi, "patch").mockImplementation(async () => { + throw new Error(`token=${secretLike}`); + }); + let failure: unknown; + try { + await reconcilePostgresClaim("project", { + service: secretLike, username: secretLike, database: secretLike, secretName: `postgres-${secretLike}`, + }); + } catch (error) { failure = error; } + expect(failure).toBeInstanceOf(DatabaseReconciliationError); + expect((failure as DatabaseReconciliationError).phase).toBe(`role secret ${method === "read" ? "lookup" : "apply"}`); + expect((failure as Error).message).not.toContain(secretLike); + }); + + test("uses safe provider status and reason without exposing request bodies or unsafe claim identifiers", () => { + const cause = Object.assign(new Error("request body DATABASE_URL=postgresql://admin:private@db.example/app"), { + statusCode: 422, + body: { reason: "Invalid", code: 422, message: "token=private" }, + }); + const failure = new DatabaseReconciliationError("managed role update", cause, { + service: "web", database: "postgresql://admin:private@db.example/app", username: "web_role", secretName: "secret", + }); + expect(failure.message).toBe( + "Database reconciliation failed during managed role update for requested database claim: Invalid (HTTP 422)", + ); + expect(failure.cause).toBe(cause); + expect(failure.message).not.toContain("private"); + }); + + test("rejects an unrecognized phase containing a syntactically valid secret-like name", () => { + const failure = new DatabaseReconciliationError("DB_PASSWORD_private123", new Error("private")); + expect(failure.phase).toBe("operation execution"); + expect(failure.message).not.toContain("DB_PASSWORD_private123"); + }); + + test("identifies failure while preparing malformed managed roles before a Cluster write", async () => { + spyOn(objectApi, "read").mockImplementation(async (resource) => resource.kind === "Secret" + ? { ...resource, data: { password: Buffer.from("private").toString("base64") } } as never + : resource.kind === "Database" ? Promise.reject({ code: 404 }) : { ...resource, spec: { managed: { roles: {} } } } as never); + const patch = spyOn(objectApi, "patch").mockImplementation(async (resource) => resource as never); + let failure: unknown; + try { + await reconcilePostgresClaim("project", { + service: "DB_PASSWORD_private123", username: "role", database: "db", secretName: "postgres-role", + }); + } catch (error) { failure = error; } + expect(failure).toBeInstanceOf(DatabaseReconciliationError); + expect((failure as DatabaseReconciliationError).phase).toBe("managed role preparation"); + expect((failure as Error).message).not.toContain("DB_PASSWORD_private123"); + expect(patch.mock.calls.map(([resource]) => resource.kind)).toEqual(["Secret"]); + }); + + test("reports malformed claims as claim discovery failures without echoing compose input", async () => { + await expect(reconcilePostgresClaims("project", { + services: { web: { volumes: ["postgresql:user:password=private"] } }, + } as ComposeSpecification)).rejects.toThrow( + "Database reconciliation failed during claim discovery: Check PostgreSQL claim declarations", + ); + }); + test("reconciles a CNPG cluster returned with managed fields without sending them back", async () => { const claim = { service: "app", @@ -193,6 +265,7 @@ describe("managed PostgreSQL claims", () => { secretName: "postgres-app", }; spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Database") throw { code: 404 }; if (resource.kind === "Secret") { return { ...resource, @@ -223,6 +296,7 @@ describe("managed PostgreSQL claims", () => { return resource as never; }, ); + spyOn(objectApi, "create").mockImplementation(async (resource) => resource as never); await reconcilePostgresClaim("project", claim); @@ -239,4 +313,85 @@ describe("managed PostgreSQL claims", () => { expect(cluster?.metadata).not.toHaveProperty("resourceVersion"); expect(cluster).not.toHaveProperty("status"); }); + + test("two projects reuse a foreign database without applying its metadata", async () => { + const claim = { service: "web", username: "sastify", database: "sastify-store", secretName: "postgres-sastify" }; + const live = { + apiVersion: "postgresql.cnpg.io/v1", kind: "Database", + metadata: { name: claim.database, namespace: "database", labels: { + "kuber.dev/project": "sastify-api", "kuber.dev/workspace-uid": "foreign-uid", + } }, + spec: { owner: claim.username, cluster: { name: "postgres" } }, + }; + const original = structuredClone(live); + spyOn(objectApi, "read").mockImplementation(async (resource) => resource.kind === "Database" + ? live as never + : resource.kind === "Secret" + ? { ...resource, data: { password: Buffer.from("shared-password").toString("base64") } } as never + : { ...resource, spec: { managed: { roles: [] } } } as never); + const patch = spyOn(objectApi, "patch").mockImplementation(async (resource) => resource as never); + const create = spyOn(objectApi, "create").mockImplementation(async (resource) => resource as never); + const compose = { services: { web: { volumes: ["postgresql:sastify/sastify-store"] } } } as ComposeSpecification; + for (const project of ["sastify-api", "another-project"]) { + const env = await reconcilePostgresClaims(project, compose); + expect(env.web?.DATABASE_URL).toBe("postgresql://sastify:shared-password@c.database.svc.cluster.local:5432/sastify-store"); + } + expect(live).toEqual(original); + expect(patch.mock.calls.map(([resource]) => resource.kind)).toEqual(["Secret", "Cluster", "Secret", "Cluster"]); + expect(create).not.toHaveBeenCalled(); + }); + + test("validates all databases in a compose file before writing the first role", async () => { + const read = spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Database" && resource.metadata?.name === "first") throw { code: 404 }; + if (resource.kind === "Database") return { ...resource, spec: { + owner: "another-role", cluster: { name: "postgres" }, + } } as never; + throw new Error("Role lookup must not run"); + }); + const patch = spyOn(objectApi, "patch").mockImplementation(async (resource) => resource as never); + const create = spyOn(objectApi, "create").mockImplementation(async (resource) => resource as never); + await expect(reconcilePostgresClaims("project", { services: { + first: { volumes: ["postgresql:first"] }, + second: { volumes: ["postgresql:second"] }, + } } as ComposeSpecification)).rejects.toMatchObject({ phase: "database ownership" }); + expect(read.mock.calls.map(([resource]) => resource.kind)).toEqual(["Database", "Database"]); + expect(patch).not.toHaveBeenCalled(); + expect(create).not.toHaveBeenCalled(); + }); + + test.each([ + ["wrong role", "other", "postgres"], + ["wrong cluster", "sastify", "other"], + ])("rejects existing database with %s before any writes", async (_case, owner, cluster) => { + const claim = { service: "web", username: "sastify", database: "sastify-store", secretName: "postgres-sastify" }; + spyOn(objectApi, "read").mockImplementation(async (resource) => resource.kind === "Database" + ? { ...resource, spec: { owner, cluster: { name: cluster } } } as never + : undefined as never); + const patch = spyOn(objectApi, "patch").mockImplementation(async (resource) => resource as never); + const create = spyOn(objectApi, "create").mockImplementation(async (resource) => resource as never); + await expect(reconcilePostgresClaim("project", claim)).rejects.toMatchObject({ phase: "database ownership" }); + expect(patch).not.toHaveBeenCalled(); + expect(create).not.toHaveBeenCalled(); + }); + + test("creates missing databases only once and fails closed on create races", async () => { + spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Database") throw { code: 404 }; + return { ...resource, spec: { managed: { roles: [] } } } as never; + }); + spyOn(objectApi, "patch").mockImplementation(async (resource) => resource as never); + const create = spyOn(objectApi, "create").mockImplementation(async () => { + throw Object.assign(new Error("Conflict"), { code: 409 }); + }); + await expect(reconcilePostgresClaims("project", { services: { + web: { volumes: ["postgresql:sastify/sastify-store"] }, + worker: { volumes: ["postgresql:sastify/sastify-store"] }, + } } as ComposeSpecification)).rejects.toMatchObject({ phase: "database apply" }); + expect(create).toHaveBeenCalledTimes(1); + expect(create.mock.calls[0]?.[0]).toMatchObject({ + metadata: { name: "sastify-store", labels: { "kuber.dev/project": "project" } }, + spec: { owner: "sastify", cluster: { name: "postgres" } }, + }); + }); }); diff --git a/tests/lib/version-update-process.test.ts b/tests/lib/version-update-process.test.ts new file mode 100644 index 0000000..f34d75c --- /dev/null +++ b/tests/lib/version-update-process.test.ts @@ -0,0 +1,51 @@ +import { expect, test } from "bun:test"; + +test("a short CLI exits while its version updater remains pending", async () => { + const moduleUrl = new URL("../../lib/version-update.ts", import.meta.url) + .href; + const script = ` + const { createVersionObserver } = await import(${JSON.stringify(moduleUrl)}); + createVersionObserver({ + currentVersion: "2.6.1-rc3", + runner: () => new Promise(() => {}), + })("2.7.0"); + `; + const child = Bun.spawn([process.execPath, "-e", script], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + const exit = await Promise.race([ + child.exited, + new Promise((_, reject) => + setTimeout(() => reject(new Error("short CLI stayed alive")), 1_000), + ), + ]); + + expect(exit).toBe(0); + expect(await new Response(child.stderr).text()).toBe(""); +}); + +test.each([ + ["success", "return true", "+ Updated to 2.7.0"], + ["failure", "throw new Error('install failed')", "+ New version available: 2.7.0"], +])("reports exactly one %s outcome when it settles", async (_name, outcome, message) => { + const moduleUrl = new URL("../../lib/version-update.ts", import.meta.url).href; + const script = ` + const { createVersionObserver } = await import(${JSON.stringify(moduleUrl)}); + createVersionObserver({ + currentVersion: "2.6.1-rc3", + runner: async () => { ${outcome}; }, + })("2.7.0"); + await new Promise((resolve) => setTimeout(resolve, 25)); + `; + const child = Bun.spawn([process.execPath, "-e", script], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + + expect(await child.exited).toBe(0); + const stderr = await new Response(child.stderr).text(); + expect(stderr.split("\n").filter((line) => line.includes(message))).toHaveLength(1); +}); diff --git a/tests/lib/version-update.test.ts b/tests/lib/version-update.test.ts new file mode 100644 index 0000000..fe52adc --- /dev/null +++ b/tests/lib/version-update.test.ts @@ -0,0 +1,201 @@ +import { describe, expect, test } from "bun:test"; +import { + compareVersions, + createVersionObserver, + installVersion, +} from "../../lib/version-update"; + +const tick = () => new Promise((resolve) => setTimeout(resolve, 0)); + +describe("server version updates", () => { + test("compares strict SemVer precedence including prereleases and ignores build metadata", () => { + expect(compareVersions("2.6.1", "2.6.1-rc3")).toBe(1); + expect(compareVersions("2.6.1-rc.10", "2.6.1-rc.3")).toBe(1); + expect(compareVersions("2.6.1-alpha.1", "2.6.1-alpha.beta")).toBe(-1); + expect(compareVersions("2.6.1-rc3.1", "2.6.1-rc3")).toBe(1); + expect(compareVersions("2.6.1+new", "2.6.1+old")).toBe(0); + expect(compareVersions("3.0.0", "2.6.1")).toBe(1); + for (const invalid of [ + "", + "v3.0.0", + "3.0", + "03.0.0", + "3.0.0-rc.01", + "3.0.0-", + "3.0.0+", + "3.0.0;touch /tmp/pwned", + "3.0.0\n", + "3.0.0-💥", + "9".repeat(129) + ".0.0", + ]) { + expect(compareVersions(invalid, "2.6.1-rc3")).toBeUndefined(); + } + }); + + test("rejects invalid install arguments before spawning", async () => { + expect(await installVersion("2.7.0;echo hacked")).toBe(false); + }); + + test("uses a quiet detached timeout argument vector", async () => { + const calls: string[][] = []; + let unrefs = 0; + const spawn = ( + argv: string[], + options: { + stdin: "ignore"; + stdout: "ignore"; + stderr: "ignore"; + detached: true; + }, + ) => { + calls.push(argv); + expect(options).toEqual({ + stdin: "ignore", + stdout: "ignore", + stderr: "ignore", + detached: true, + }); + return { + exited: Promise.resolve(0), + unref: () => { + unrefs++; + }, + }; + }; + expect(await installVersion("2.7.0-rc.1+build.2", spawn, 1)).toBe(true); + expect(calls).toEqual([ + [ + "timeout", + "--signal=TERM", + "--kill-after=2s", + "1s", + "bun", + "i", + "-g", + "--no-cache", + "@dmgnr/kuber@2.7.0-rc.1+build.2", + ], + ]); + expect(unrefs).toBe(1); + }); + + test("ignores malformed, equal, and older headers without attempting an install", async () => { + const runs: string[] = []; + const lines: string[] = []; + const observe = createVersionObserver({ + currentVersion: "2.6.1-rc3", + runner: async (version) => { + runs.push(version); + return true; + }, + report: (line) => lines.push(line), + }); + for (const version of [ + null, + "garbage", + "2.6.1-rc3", + "2.6.0", + "2.6.1-rc2", + "2.6.1-rc3+build", + ]) + observe(version); + await tick(); + expect(runs).toEqual([]); + expect(lines).toEqual([]); + }); + + test("starts once for concurrent newer responses, reports success in gray on stderr", async () => { + let finish!: (value: boolean) => void; + const pending = new Promise((resolve) => { + finish = resolve; + }); + const runs: string[] = []; + const lines: string[] = []; + const observe = createVersionObserver({ + currentVersion: "2.6.1-rc3", + runner: (version) => { + runs.push(version); + return pending; + }, + report: (line) => lines.push(line), + }); + observe("2.6.1"); + observe("9.0.0"); + expect(runs).toEqual([]); + await tick(); + expect(runs).toEqual(["2.6.1"]); + expect(lines).toEqual([]); + finish(true); + await tick(); + expect(lines).toEqual(["\x1b[90m+ Updated to 2.6.1\x1b[0m\n"]); + observe("10.0.0"); + await tick(); + expect(runs).toHaveLength(1); + }); + + test("reports failure or rejection once without throwing into requests", async () => { + for (const runner of [ + async () => false, + async () => { + throw new Error("offline"); + }, + ]) { + const lines: string[] = []; + const observe = createVersionObserver({ + currentVersion: "2.6.1-rc3", + runner, + report: (line) => lines.push(line), + }); + observe("2.7.0"); + observe("2.8.0"); + await tick(); + expect(lines).toEqual([ + "\x1b[90m+ New version available: 2.7.0\x1b[0m\n", + ]); + } + }); + + test("a detached pending install does not delay a short CLI or pollute its streams", async () => { + const url = new URL("../../lib/version-update.ts", import.meta.url).href; + const script = ` + const { createVersionObserver } = await import(${JSON.stringify(url)}); + createVersionObserver({ currentVersion: '2.6.1', + runner: () => new Promise(() => {}) })('2.7.0'); + `; + const child = Bun.spawn([process.execPath, "-e", script], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + const timeout = setTimeout(() => child.kill(), 1_000); + try { + expect(await child.exited).toBe(0); + expect(await new Response(child.stdout).text()).toBe(""); + expect(await new Response(child.stderr).text()).toBe(""); + } finally { + clearTimeout(timeout); + } + }); + + test("a short subprocess emits one outcome to stderr when its runner settles", async () => { + const url = new URL("../../lib/version-update.ts", import.meta.url).href; + for (const result of ["true", "false", "reject"]) { + const script = ` + const { createVersionObserver } = await import(${JSON.stringify(url)}); + createVersionObserver({ currentVersion: '2.6.1', + runner: async () => ${result === "reject" ? "Promise.reject(Error('offline'))" : result} })('2.7.0'); + await new Promise(resolve => setTimeout(resolve, 20)); + `; + const child = Bun.spawn([process.execPath, "-e", script], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + expect(await child.exited).toBe(0); + expect(await new Response(child.stdout).text()).toBe(""); + expect(await new Response(child.stderr).text()).toBe( + `\x1b[90m+ ${result === "true" ? "Updated to " : "New version available: "}2.7.0\x1b[0m\n`, + ); + } + }); +}); diff --git a/tests/server/api-keys.test.ts b/tests/server/api-keys.test.ts index 7895ee9..e3a1a2c 100644 --- a/tests/server/api-keys.test.ts +++ b/tests/server/api-keys.test.ts @@ -7,7 +7,11 @@ import type { ManagementService } from "../../server/management"; import { MemoryTrustStore } from "../../server/trust-store"; import { MemoryWorkspaceStore } from "../../server/workspace-store"; -const now = Date.parse("2026-09-05T00:00:00.000Z"); +const now = Date.parse("2030-09-05T00:00:00.000Z"); +const finiteExpiry = new Date(now + 30 * 24 * 60 * 60 * 1000).toISOString(); +const delegatedLaterExpiry = new Date( + now + 31 * 24 * 60 * 60 * 1000, +).toISOString(); function request(path: string, init: RequestInit = {}, token = "admin-token") { const headers = new Headers(init.headers); @@ -33,7 +37,7 @@ async function setup() { tokenHash: hashToken("admin-token"), username: "admin", authVersion: 1, - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); return { store, @@ -86,7 +90,7 @@ describe("API keys", () => { tokenHash: hashToken("ci-key"), username: "ci", capabilities: ["kubernetes:read"], - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); expect((await app(request("/api/v2/users", {}, "ci-key"))).status).toBe( 403, @@ -104,7 +108,7 @@ describe("API keys", () => { username: "ci", capabilities: ["users:write", "kubernetes:read"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); const create = (body: unknown) => app( @@ -149,14 +153,14 @@ describe("API keys", () => { const laterExpiry = await create({ capabilities: ["kubernetes:read"], workspace: "shop", - expiresAt: "2026-10-06T00:00:00.000Z", + expiresAt: delegatedLaterExpiry, }); expect(laterExpiry.status).toBe(403); const subset = await create({ capabilities: ["kubernetes:read"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); expect(subset.status).toBe(201); expect( @@ -194,7 +198,7 @@ describe("API keys", () => { tokenHash: hashToken("unscoped-delegation"), username: "ci", capabilities: ["users:write", "kubernetes:read"], - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); const unscopedSubset = await app( request( @@ -204,7 +208,7 @@ describe("API keys", () => { headers: { "content-type": "application/json" }, body: JSON.stringify({ capabilities: ["kubernetes:read"], - expiresAt: "2026-09-20T00:00:00.000Z", + expiresAt: new Date(now + 15 * 24 * 60 * 60 * 1000).toISOString(), }), }, "unscoped-delegation", @@ -248,7 +252,7 @@ describe("API keys", () => { tokenHash: hashToken("expired-key"), username: "ci", capabilities: ["kubernetes:read"], - expiresAt: "2026-09-04T00:00:00.000Z", + expiresAt: new Date(now - 24 * 60 * 60 * 1000).toISOString(), }); expect((await app(request("/api/v2/me", {}, "expired-key"))).status).toBe( 401, @@ -261,7 +265,7 @@ describe("API keys", () => { headers: { "content-type": "application/json" }, body: JSON.stringify({ capabilities: ["kubernetes:read"], - expiresAt: "2027-09-06T00:00:00.000Z", + expiresAt: new Date(now + 366 * 24 * 60 * 60 * 1000).toISOString(), }), }), ) @@ -291,7 +295,7 @@ describe("API keys", () => { headers: { "content-type": "application/json" }, body: JSON.stringify({ capabilities: ["kubernetes:read"], - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }), }), ); @@ -300,7 +304,7 @@ describe("API keys", () => { token: string; expiresAt: string; }; - expect(finiteKey.expiresAt).toBe("2026-10-05T00:00:00.000Z"); + expect(finiteKey.expiresAt).toBe(finiteExpiry); expect( (await app(request("/api/v2/logout", { method: "POST" }, key.token))) .status, @@ -326,7 +330,7 @@ describe("API keys", () => { username: "ci", capabilities: ["kubernetes:read"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); expect( (await app(request("/api/v2/workspaces/other", {}, "scoped-key"))).status, @@ -366,7 +370,7 @@ describe("API keys", () => { username: "ci", capabilities: ["kubernetes:write"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); const apply = (workspace: string) => app( @@ -402,7 +406,7 @@ describe("API keys", () => { username: "ci", capabilities: ["users:read"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); await auditStore.append({ actor: { username: "admin" }, @@ -459,7 +463,7 @@ describe("API keys", () => { username: "ci", capabilities: ["kubernetes:read"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); const unfiltered = await app( @@ -525,7 +529,7 @@ describe("API keys", () => { username: "admin", capabilities: ["kubernetes:write"], workspace: "kuber-system", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); await store.createApiKey({ id: "key_platform_scope", @@ -533,7 +537,7 @@ describe("API keys", () => { username: "admin", capabilities: ["platform:adopt"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); await store.createApiKey({ id: "key_platform_allowed", @@ -541,7 +545,7 @@ describe("API keys", () => { username: "admin", capabilities: ["platform:adopt"], workspace: "kuber-system", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); const adopt = (token: string) => app( @@ -569,7 +573,7 @@ describe("API keys", () => { username: "ci", capabilities: ["kubernetes:write"], workspace: "shop", - expiresAt: "2026-10-05T00:00:00.000Z", + expiresAt: finiteExpiry, }); const matching = await app( request( diff --git a/tests/server/app.test.ts b/tests/server/app.test.ts index b25c99e..5930334 100644 --- a/tests/server/app.test.ts +++ b/tests/server/app.test.ts @@ -9,6 +9,7 @@ import { } from "../../server/operation-store"; import { MemoryWorkspaceStore } from "../../server/workspace-store"; import { MemoryTrustStore } from "../../server/trust-store"; +import type { BuildController } from "../../server/build-controller"; function request( path: string, @@ -340,7 +341,8 @@ describe("operation response safety", () => { expect(immediate.status).toBe(500); expect(immediateBody).not.toContain("database-password"); expect(immediateBody).not.toContain("kube-secret"); - expect(immediateBody).toContain("OPERATION_FAILED"); + expect(immediateBody).toContain("DATABASE_RECONCILE_FAILED"); + expect(immediateBody).toContain("during operation execution: Unexpected failure; check server logs using the operation ID"); const operationId = (await operationStore.list())[0]!.metadata.name; const retrieved = await app( @@ -529,6 +531,98 @@ async function authenticatedStore(role: "viewer" | "operator" | "admin") { } describe("kuber v2 HTTP routes", () => { + test("streams build status and logs, supports resume cursors, and keeps version headers", async () => { + const controller = { + getBuildProject: async () => "demo", + getBuildStatus: async () => ({ state: "running", phase: "building" }), + getBuildEvents: async (_id: string, after = 0) => [ + { type: "log", sequence: 4, message: "build output" }, + ].filter((event) => event.sequence > after), + reconcileBuild: async () => ({ state: "running", phase: "building" }), + } as unknown as BuildController; + const app = createApp({ store: await authenticatedStore("operator"), builds: controller }); + const abort = new AbortController(); + const response = await app(request( + "/api/v2/builds/build-1/events?after=2", + { headers: { accept: "text/event-stream" }, signal: abort.signal }, + "token", + )); + expect(response.status).toBe(200); + expect(response.headers.get("content-type")).toContain("text/event-stream"); + expect(response.headers.get("x-kuber-version")).toBeTruthy(); + const reader = response.body!.getReader(); + let body = ""; + while (!body.includes("event: status")) { + const { done, value } = await reader.read(); + if (done) break; + body += new TextDecoder().decode(value); + } + expect(body).toContain("id: 4\nevent: log"); + expect(body).toContain("build output"); + expect(body).toContain("event: status"); + abort.abort(); + await reader.cancel(); + + const resumedAbort = new AbortController(); + const resumed = await app(request( + "/api/v2/builds/build-1/events", + { headers: { accept: "text/event-stream", "last-event-id": "4" }, signal: resumedAbort.signal }, + "token", + )); + const resumedReader = resumed.body!.getReader(); + const resumedChunk = await resumedReader.read(); + expect(new TextDecoder().decode(resumedChunk.value)).not.toContain("id: 4"); + resumedAbort.abort(); + await resumedReader.cancel(); + }); + + test("checks authentication and workspace scope before opening a build stream", async () => { + const controller = { + getBuildProject: async () => "other", + getBuildStatus: async () => ({ state: "running", phase: "building" }), + getBuildEvents: async () => [], + reconcileBuild: async () => ({ state: "running", phase: "building" }), + } as unknown as BuildController; + const app = createApp({ store: await authenticatedStore("operator"), builds: controller }); + const unauthenticated = await app(request( + "/api/v2/builds/build-1/events", + { headers: { accept: "text/event-stream" } }, + )); + expect(unauthenticated.status).toBe(401); + + const store = new MemoryAuthStore(); + await store.putUser({ username: "scoped", passwordHash: "hash", roles: ["operator"] }); + await store.createApiKey({ + id: "scoped-key-12345678", username: "scoped", tokenHash: hashToken("scoped-token"), + capabilities: ["kubernetes:write"], workspace: "demo", + }); + const scopedApp = createApp({ store, builds: controller }); + const forbidden = await scopedApp(request( + "/api/v2/builds/build-1/events", + { headers: { accept: "text/event-stream" } }, + "scoped-token", + )); + expect(forbidden.status).toBe(403); + expect(forbidden.headers.get("content-type")).not.toContain("text/event-stream"); + }); + + test("rejects malformed SSE cursors before opening the stream", async () => { + const controller = { + getBuildProject: async () => "demo", + getBuildStatus: async () => ({ state: "running", phase: "building" }), + getBuildEvents: async () => [], + reconcileBuild: async () => ({ state: "running", phase: "building" }), + } as unknown as BuildController; + const app = createApp({ store: await authenticatedStore("operator"), builds: controller }); + const response = await app(request( + "/api/v2/builds/build-1/events?after=1.5", + { headers: { accept: "text/event-stream" } }, + "token", + )); + expect(response.status).toBe(400); + expect(response.headers.get("content-type")).toContain("application/problem+json"); + }); + test("grants, lists, revokes, and enforces namespace trust for applies", async () => { const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid", @@ -1788,11 +1882,134 @@ describe("kuber v2 HTTP routes", () => { expect(result.status).toBe(500); const problem = await result.json() as { code: string; detail: string }; expect(problem.code).toBe("DATABASE_RECONCILE_FAILED"); - expect(problem.detail).toContain("database apply for database web_db (service web, role web_role): Forbidden"); + expect(problem.detail).toContain("database apply for requested database claim: Forbidden"); expect(JSON.stringify(problem)).not.toContain("private-value"); } expect((await operationStore.get("operation-db-failure"))?.status.error?.message) - .toContain("database apply for database web_db"); + .toContain("database apply for requested database claim"); + }); + + test("never exposes syntactically valid secret-like database claim identifiers", 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(undefined, () => "db-secret-identifier"); + const { DatabaseReconciliationError } = await import("../../lib/database"); + const secretLike = ["DB_PASSWORD_private123", "svc_api_token_abc", "role_secret_key_123"]; + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + operationStore, + management: { + reconcileDatabases: async () => { + throw new DatabaseReconciliationError( + "database apply", + Object.assign(new Error("Forbidden"), { statusCode: 403 }), + { service: secretLike[1]!, username: secretLike[2]!, database: secretLike[0]!, secretName: "safe-secret-name" }, + ); + }, + } as unknown as ManagementService, + }); + const response = await app(request( + "/api/v2/workspaces/demo/databases", + { method: "POST", headers: { "idempotency-key": "db-secret-identifier" }, body: JSON.stringify({ compose: {} }) }, + "token", + )); + const responseBody = JSON.stringify(await response.json()); + const operation = await operationStore.get("operation-db-secret-identifier"); + + expect(response.status).toBe(500); + expect(responseBody).toContain("database apply for requested database claim: Forbidden (HTTP 403)"); + expect(operation?.status.error?.message).toContain("database apply for requested database claim"); + for (const identifier of secretLike) { + expect(responseBody).not.toContain(identifier); + expect(JSON.stringify(operation)).not.toContain(identifier); + } + }); + + test("logs allowlisted database failure metadata with the operation ID, never provider data", 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(undefined, () => "diagnostic"); + const logs: Record[] = []; + const { DatabaseReconciliationError } = await import("../../lib/database"); + const secretLike = "DB_PASSWORD_private123"; + const cause = Object.assign(new Error(`token=${secretLike}`), { + name: secretLike, + code: "ECONNREFUSED", + statusCode: 503, + body: { message: secretLike, headers: { authorization: secretLike } }, + stack: `Error: ${secretLike}\n at ${secretLike} (/tmp/${secretLike}/lib/database.ts:22:7)`, + }); + const app = createApp({ + store: await authenticatedStore("operator"), workspaceStore, operationStore, + logger: { error: (entry) => logs.push(entry) }, + management: { reconcileDatabases: async () => { + throw new DatabaseReconciliationError("database apply", cause, { + service: secretLike, username: secretLike, database: secretLike, secretName: secretLike, + }); + } } as unknown as ManagementService, + }); + const response = await app(request("/api/v2/workspaces/demo/databases", { + method: "POST", body: JSON.stringify({ compose: {} }), + }, "token")); + expect(response.status).toBe(500); + expect(await response.json()).toMatchObject({ + code: "DATABASE_RECONCILE_FAILED", operationId: "operation-diagnostic", + }); + const diagnostic = logs.filter((entry) => entry.event === "operation.database_reconcile.failed"); + expect(diagnostic).toHaveLength(1); + expect(diagnostic[0]).toMatchObject({ + operationId: "operation-diagnostic", phase: "database apply", + errorClass: "UnknownError", providerCode: "ECONNREFUSED", providerStatus: 503, + topFrame: "lib/database.ts:22:7", + }); + expect(JSON.stringify(diagnostic)).not.toContain(secretLike); + expect(JSON.stringify(await operationStore.get("operation-diagnostic"))).not.toContain(secretLike); + }); + + test("retains a safe diagnosis for an ordinary database provider error through HTTP, persistence and polling", 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(undefined, () => "ordinary-db-error"); + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + operationStore, + management: { + reconcileDatabases: async () => { + throw Object.assign(new Error("provider failed: DATABASE_URL=postgresql://admin:private@db.example/app"), { + statusCode: 403, + body: { reason: "Forbidden", message: "password=private", code: 403 }, + }); + }, + } as unknown as ManagementService, + }); + const reconcile = () => app(request( + "/api/v2/workspaces/demo/databases", + { method: "POST", headers: { "idempotency-key": "ordinary-db-error" }, body: JSON.stringify({ compose: {} }) }, + "token", + )); + const expected = "Database reconciliation failed during operation execution: Forbidden (HTTP 403)"; + for (const result of [await reconcile(), await reconcile()]) { + expect(result.status).toBe(500); + expect(await result.json()).toMatchObject({ + code: "DATABASE_RECONCILE_FAILED", + detail: expected, + operationId: "operation-ordinary-db-error", + }); + } + const persisted = await operationStore.get("operation-ordinary-db-error"); + expect(persisted?.status.error).toEqual({ code: "DATABASE_RECONCILE_FAILED", message: expected }); + const polled = await app(request("/api/v2/operations/operation-ordinary-db-error", {}, "token")); + expect(polled.status).toBe(200); + expect(await polled.json()).toMatchObject({ status: { error: { message: expected } } }); + expect(JSON.stringify(persisted)).not.toContain("private"); }); test("routes workspace adoption and keeps platform adoption admin-only", async () => { diff --git a/tests/server/build-event-stream.test.ts b/tests/server/build-event-stream.test.ts new file mode 100644 index 0000000..a01a390 --- /dev/null +++ b/tests/server/build-event-stream.test.ts @@ -0,0 +1,330 @@ +import { describe, expect, test } from "bun:test"; +import { BuildEventStreamHub } from "../../server/build-event-stream"; +import type { BuildStatus } from "../../shared/build-protocol"; + +const base: BuildStatus = { + version: 1, + id: "a", + state: "running", + createdAt: "now", +}; +const tick = () => Bun.sleep(10); + +describe("build SSE hub", () => { + test("reports an explicit gap when the resume cursor predates retained logs", async () => { + const abort = new AbortController(); + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => base, + reconcileBuild: async () => base, + getBuildEvents: async (_id, after) => [ + { type: "status" as const, status: base }, + ...[4, 5] + .filter((sequence) => sequence > (after ?? 0)) + .map((sequence) => ({ + type: "log" as const, + id: "a", + sequence, + message: `log ${sequence}\n`, + })), + ], + }, + { pollMs: 100, heartbeatMs: 1_000 }, + ); + const response = await hub.open( + "a", + new Request("https://test/builds/a/events", { + headers: { "last-event-id": "1" }, + signal: abort.signal, + }), + ); + const reader = response.body!.getReader(); + const decoder = new TextDecoder(); + const first = decoder.decode((await reader.read()).value); + expect(first).toContain("event: gap\n"); + expect(first).toContain('"missing":2'); + expect(first).toContain("id: 3\n"); + expect(decoder.decode((await reader.read()).value)).toContain("id: 4\nevent: log"); + abort.abort(); + await reader.cancel(); + }); + + test("coalesces reconciliation, resumes logs, dedupes status and drains terminal", async () => { + let calls = 0; + let status = base; + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => status, + reconcileBuild: async () => { + calls++; + return status; + }, + getBuildEvents: async (_id, after) => [ + { type: "status" as const, status: base }, + ...[1, 2] + .filter((sequence) => sequence > (after ?? 0)) + .map((sequence) => ({ + type: "log" as const, + id: "a", + sequence, + message: `log ${sequence}\n`, + })), + ], + }, + { pollMs: 20, heartbeatMs: 40 }, + ); + const a = new AbortController(); + const b = new AbortController(); + const request = (signal: AbortSignal, after: string) => + new Request(`https://test/builds/a/events?after=${after}`, { signal }); + const first = await hub.open("a", request(a.signal, "1")); + const second = await hub.open("a", request(b.signal, "2")); + expect(first.headers.get("content-type")).toContain("text/event-stream"); + const reader1 = first.body!.getReader(); + const reader2 = second.body!.getReader(); + const text = new TextDecoder(); + expect(text.decode((await reader1.read()).value)).toContain( + "id: 2\nevent: log\n", + ); + expect(text.decode((await reader1.read()).value)).toContain( + "event: status\n", + ); + expect(text.decode((await reader2.read()).value)).toContain( + "event: status\n", + ); + await tick(); + expect(calls).toBe(1); + status = { ...base, state: "succeeded" }; + const remaining = async (reader: typeof reader1) => { + let output = ""; + for (;;) { + const { value, done } = await reader.read(); + if (done) return output; + output += text.decode(value); + } + }; + expect(await remaining(reader1)).toContain('"state":"succeeded"'); + expect(await remaining(reader2)).toContain('"state":"succeeded"'); + a.abort(); + b.abort(); + }); + + test("validates cursor, stops aborted subscribers and emits heartbeat comments", async () => { + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => base, + reconcileBuild: async () => base, + getBuildEvents: async () => [], + }, + { pollMs: 5, heartbeatMs: 5 }, + ); + await expect( + hub.open("a", new Request("https://test/events?after=-1")), + ).rejects.toThrow(RangeError); + const abort = new AbortController(); + const response = await hub.open( + "a", + new Request("https://test/events", { signal: abort.signal }), + ); + const reader = response.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toContain( + "event: status", + ); + expect(new TextDecoder().decode((await reader.read()).value)).toBe( + ": heartbeat\n\n", + ); + abort.abort(); + expect((await reader.read()).done).toBe(true); + }); + + test("reads persisted progress when another replica owns reconciliation", async () => { + let status = base; + let sequence = 0; + let reconciles = 0; + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => status, + reconcileBuild: async () => { + reconciles++; + throw new Error("build lease held by another replica"); + }, + getBuildEvents: async (_id, after) => + sequence > (after ?? 0) + ? [ + { + type: "log" as const, + id: "a", + sequence, + message: "remote log\n", + }, + ] + : [], + }, + { pollMs: 5, heartbeatMs: 50 }, + ); + const abort = new AbortController(); + const response = await hub.open( + "a", + new Request("https://test/events", { signal: abort.signal }), + ); + const reader = response.body!.getReader(); + const decoder = new TextDecoder(); + expect(decoder.decode((await reader.read()).value)).toContain( + '"state":"running"', + ); + sequence = 1; + expect(decoder.decode((await reader.read()).value)).toContain( + "id: 1\nevent: log", + ); + status = { ...base, state: "succeeded" }; + let terminalFrame = ""; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + terminalFrame += decoder.decode(value); + } + expect(terminalFrame).toContain('"state":"succeeded"'); + expect(reconciles).toBeGreaterThan(0); + abort.abort(); + }); + + test("heartbeats during slow reconciliation and tears down on abort", async () => { + const blocked = new Promise(() => {}); + const signals: AbortSignal[] = []; + let observations = 0; + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => base, + reconcileBuild: async (_id, options) => { + if (options?.signal) signals.push(options.signal); + return blocked; + }, + getBuildEvents: async () => { + observations++; + return []; + }, + }, + { pollMs: 5, heartbeatMs: 5 }, + ); + const abort = new AbortController(); + const response = await hub.open( + "a", + new Request("https://test/events", { signal: abort.signal }), + ); + const reader = response.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toContain( + "event: status", + ); + expect(new TextDecoder().decode((await reader.read()).value)).toBe( + ": heartbeat\n\n", + ); + await Bun.sleep(25); + expect(observations).toBeGreaterThan(1); + expect(signals).toHaveLength(1); + abort.abort(); + for (;;) { + if ((await reader.read()).done) break; + } + expect(signals[0]?.aborted).toBe(true); + const stoppedAt = observations; + await Bun.sleep(20); + expect(observations).toBe(stoppedAt); + const next = new AbortController(); + const reopened = await hub.open( + "a", + new Request("https://test/events", { signal: next.signal }), + ); + const nextReader = reopened.body!.getReader(); + expect(new TextDecoder().decode((await nextReader.read()).value)).toContain( + "event: status", + ); + expect(signals).toHaveLength(2); + expect(signals[1]?.aborted).toBe(false); + next.abort(); + await nextReader.cancel(); + }); + + test("bounds shared reconciliation writes while observing two subscribers independently", async () => { + let reconciles = 0; + let observations = 0; + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => base, + reconcileBuild: async () => { + reconciles++; + return base; + }, + getBuildEvents: async () => { + observations++; + return []; + }, + }, + { pollMs: 10, reconcileMs: 80, heartbeatMs: 1_000 }, + ); + const a = new AbortController(); + const b = new AbortController(); + const first = await hub.open( + "a", + new Request("https://test/events", { signal: a.signal }), + ); + const second = await hub.open( + "a", + new Request("https://test/events", { signal: b.signal }), + ); + expect( + new TextDecoder().decode((await first.body!.getReader().read()).value), + ).toContain("event: status"); + expect( + new TextDecoder().decode((await second.body!.getReader().read()).value), + ).toContain("event: status"); + await Bun.sleep(45); + expect(reconciles).toBe(1); // One lease/write attempt, not one per poll or subscriber. + expect(observations).toBeGreaterThanOrEqual(3); + for (let i = 0; i < 15 && reconciles < 2; i++) await tick(); + expect(reconciles).toBe(2); + a.abort(); + b.abort(); + }); + + test("closes a slow subscriber at the bounded queue limit for resumable logs", async () => { + const hub = new BuildEventStreamHub( + { + getBuildStatus: async () => base, + reconcileBuild: async () => base, + getBuildEvents: async (_id, after) => + Array.from({ length: 40 }, (_, index) => index + 1) + .filter((sequence) => sequence > (after ?? 0)) + .map((sequence) => ({ + type: "log" as const, + id: "a", + sequence, + message: `log ${sequence}\n`, + })), + }, + { pollMs: 5, heartbeatMs: 50 }, + ); + const response = await hub.open("a", new Request("https://test/events")); + await tick(); // Allow the producer to fill its queue before consumption. + const reader = response.body!.getReader(); + const frames: string[] = []; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + frames.push(new TextDecoder().decode(value)); + } + expect(frames).toHaveLength(16); + expect(frames[0]).toContain("id: 1\nevent: log"); + expect(frames[15]).toContain("id: 16\nevent: log"); + const resumed = await hub.open( + "a", + new Request("https://test/events", { + headers: { "last-event-id": "16" }, + }), + ); + const resumedReader = resumed.body!.getReader(); + expect( + new TextDecoder().decode((await resumedReader.read()).value), + ).toContain("id: 17\nevent: log"); + await resumedReader.cancel(); + }); +}); diff --git a/tests/server/management.test.ts b/tests/server/management.test.ts index 2362a53..cef28a0 100644 --- a/tests/server/management.test.ts +++ b/tests/server/management.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test } from "bun:test"; import type { KubernetesObject, V1Deployment } from "@kubernetes/client-node"; +import { DatabaseReconciliationError } from "../../lib/database"; import { createManagementService, type ManagementDependencies, @@ -67,6 +68,27 @@ function dependencies( } describe("server management service", () => { + test.each([ + ["namespace precheck", { readNamespace: async () => { throw new Error("token=private"); } }], + ["database dependency", { reconcileDatabases: async () => { throw new Error("token=private"); } }], + ["database resource listing", { listDatabaseResources: async () => { throw new Error("token=private"); } }], + ["database resource ownership", { applyResource: async () => { throw new Error("token=private"); }, listDatabaseResources: async () => [object("Database", "DB_PASSWORD_private123", "db-uid")] }], + ] as const)("labels database %s failures without leaking details", async (phase, overrides) => { + const service = createManagementService(dependencies(overrides)); + let failure: unknown; + try { await service.reconcileDatabases(workspace, { services: {} }); } + catch (error) { failure = error; } + expect(failure).toBeInstanceOf(DatabaseReconciliationError); + expect((failure as DatabaseReconciliationError).phase).toBe(phase); + expect((failure as Error).message).not.toContain("private"); + }); + + test("preserves an existing database substep rather than masking it", async () => { + const service = createManagementService(dependencies({ + reconcileDatabases: async () => { throw new DatabaseReconciliationError("database apply", new Error("private")); }, + })); + await expect(service.reconcileDatabases(workspace, { services: {} })).rejects.toMatchObject({ phase: "database apply" }); + }); test("waits on deployments together and cancels siblings after a crash", async () => { const started: string[] = []; let siblingCancelled = false; @@ -344,14 +366,18 @@ describe("server management service", () => { namespace: "database", }, }; - const service = createManagementService( - dependencies({ listDatabaseResources: async () => [database] }), - ); + const deleted: ResourceIdentity[] = []; + const service = createManagementService(dependencies({ + listDatabaseResources: async () => [database], + deleteResource: async (identity) => { deleted.push(identity); }, + })); const plan = await service.planDown(workspace, true); expect(plan.delete.map(({ kind, uid }) => [kind, uid])).toEqual([ - ["Database", "database-uid"], ["Namespace", "namespace-uid"], ]); + expect(plan.retained.map(({ kind, uid }) => [kind, uid])).toEqual([["Database", "database-uid"]]); + await service.down(workspace, true); + expect(deleted.map(({ kind }) => kind)).toEqual(["Namespace"]); const missingUid = createManagementService( dependencies({ @@ -368,6 +394,87 @@ describe("server management service", () => { ); }); + test("preserves foreign and owned database CRs on full down while deleting other resources", async () => { + const shared = { ...object("Database", "sastify-store", "shared-uid", { + "app.kubernetes.io/managed-by": "kuber", [WORKSPACE_PROJECT_LABEL]: workspace.project, + [WORKSPACE_UID_LABEL]: "other-workspace", + }), metadata: { ...object("Database", "sastify-store", "shared-uid").metadata, + namespace: "database", labels: { + "app.kubernetes.io/managed-by": "kuber", [WORKSPACE_PROJECT_LABEL]: workspace.project, + [WORKSPACE_UID_LABEL]: "other-workspace", + } } }; + const owned = { ...object("Database", "local", "local-uid"), metadata: { + ...object("Database", "local", "local-uid").metadata, namespace: "database", + } }; + const deleted: string[] = []; + const service = createManagementService(dependencies({ + listProjectResources: async () => [object("Deployment", "api", "deployment-uid")], + listDatabaseResources: async () => [shared, owned], + deleteResource: async (identity) => { deleted.push(identity.kind); }, + })); + const plan = await service.down(workspace, true); + expect(plan.delete.map(({ kind }) => kind)).toEqual(["Deployment", "Namespace"]); + expect(plan.retained.map(({ name }) => name)).toEqual(["local"]); + expect(deleted).toEqual(["Deployment", "Namespace"]); + }); + + test("reuses declared foreign database without taking its workspace UID", async () => { + const db = { ...object("Database", "sastify-store", "shared-uid", { + "app.kubernetes.io/managed-by": "kuber", [WORKSPACE_PROJECT_LABEL]: workspace.project, + [WORKSPACE_UID_LABEL]: "other-workspace", + }), metadata: { + ...object("Database", "sastify-store", "shared-uid").metadata, + namespace: "database", + labels: { "app.kubernetes.io/managed-by": "kuber", + [WORKSPACE_PROJECT_LABEL]: workspace.project, [WORKSPACE_UID_LABEL]: "other-workspace" }, + }, spec: { owner: "sastify", cluster: { name: "postgres" } } }; + const before = structuredClone(db); + const applied: KubernetesObject[] = []; + const service = createManagementService(dependencies({ + readNamespace: async (project) => ({ uid: "namespace-uid", labels: { + "app.kubernetes.io/managed-by": "kuber", + [WORKSPACE_UID_LABEL]: project === workspace.project ? workspace.uid : "second-uid", + } }), + listDatabaseResources: async () => [db], + applyResource: async (resource) => { applied.push(resource); return resource; }, + })); + const compose = { services: { web: { volumes: ["postgresql:sastify/sastify-store"] } } } as never; + for (const project of [workspace, { project: "second-project", uid: "second-uid" }]) { + await service.reconcileDatabases(project, compose); + } + expect(applied).toEqual([]); + expect(db).toEqual(before); + await expect(service.reconcileDatabases(workspace, { services: {} })).rejects.toMatchObject({ phase: "database resource ownership" }); + expect(applied).toEqual([]); + }); + + test("adopts unowned and same-workspace databases and rejects foreign conflicting specs", async () => { + const resources: KubernetesObject[] = [ + { ...object("Database", "new-db", "new-uid"), metadata: { + ...object("Database", "new-db", "new-uid").metadata, + namespace: "database", labels: { "app.kubernetes.io/managed-by": "kuber" }, + } }, + { ...object("Database", "owned-db", "owned-uid"), metadata: { + ...object("Database", "owned-db", "owned-uid").metadata, namespace: "database", + } }, + ]; + const applied: KubernetesObject[] = []; + const service = createManagementService(dependencies({ + listDatabaseResources: async () => resources, + applyResource: async (resource) => { applied.push(resource); return resource; }, + })); + await service.reconcileDatabases(workspace, { services: {} }); + expect(applied.map((resource) => resource.metadata?.labels?.[WORKSPACE_UID_LABEL])).toEqual([workspace.uid, workspace.uid]); + resources[0] = { ...resources[0]!, metadata: { ...resources[0]!.metadata, + labels: { [WORKSPACE_UID_LABEL]: "other-workspace" } }, spec: { + owner: "wrong", cluster: { name: "postgres" }, + } } as KubernetesObject; + await expect(service.reconcileDatabases(workspace, { services: { web: { + volumes: ["postgresql:sastify/new-db"], + } } } as never)).rejects.toMatchObject({ phase: "database resource ownership" }); + expect(applied).toHaveLength(2); + }); + test("stop deletes matching HPAs before scaling selected deployments", async () => { const calls: string[] = []; const service = createManagementService( diff --git a/tests/server/operation-store.test.ts b/tests/server/operation-store.test.ts index c6f3dd8..98af1d5 100644 --- a/tests/server/operation-store.test.ts +++ b/tests/server/operation-store.test.ts @@ -121,6 +121,13 @@ describe("operation store", () => { code: "RECONCILE_FAILED", message: "Database reconciliation failed", }); + expect(sanitizeOperationError({ + code: "DATABASE_RECONCILE_FAILED", + message: "Database reconciliation failed during database apply: Forbidden (HTTP 403) DB_PASSWORD=private", + }, "databases.reconcile")).toEqual({ + code: "DATABASE_RECONCILE_FAILED", + message: "Database reconciliation failed during database apply: Forbidden (HTTP 403) DB_PASSWORD=[REDACTED]", + }); }); test("enforces the operation state machine", async () => { diff --git a/tests/server/version-header.test.ts b/tests/server/version-header.test.ts new file mode 100644 index 0000000..7579cd1 --- /dev/null +++ b/tests/server/version-header.test.ts @@ -0,0 +1,74 @@ +import { expect, test } from "bun:test"; +import { createApp, execProblem } from "../../server/app"; +import { hashToken, MemoryAuthStore } from "../../server/auth"; +import { MemoryWorkspaceStore } from "../../server/workspace-store"; +import type { LogService } from "../../server/log-service"; +import { KUBER_VERSION, KUBER_VERSION_HEADER } from "../../shared/version"; + +test("returns the bundled version on success, errors, preflight, empty and streaming responses", async () => { + const store = new MemoryAuthStore(); + await store.putUser({ + username: "viewer", + passwordHash: "hash", + roles: ["viewer"], + }); + await store.putSession({ + tokenHash: hashToken("token"), + username: "viewer", + authVersion: 1, + expiresAt: "2030-01-01T00:00:00.000Z", + }); + const workspaceStore = new MemoryWorkspaceStore({ + uid: () => "workspace-uid", + }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const logs = { + async collect() { + return [{ type: "log", message: "hello" }]; + }, + } as unknown as LogService; + const app = createApp({ store, workspaceStore, logs }); + const request = (path: string, init: RequestInit = {}) => + new Request(`https://kuber.astrxl.dev${path}`, init); + const authenticated = (path: string, init: RequestInit = {}) => + request(path, { + ...init, + headers: { + authorization: "Bearer token", + ...Object.fromEntries(new Headers(init.headers)), + }, + }); + + const success = await app(request("/api/v2/health")); + const error = await app(request("/api/v2/me")); + const preflight = await app( + request("/api/v2/health", { + method: "OPTIONS", + headers: { origin: "https://kuber.astrxl.dev" }, + }), + ); + const stream = await app(authenticated("/api/v2/workspaces/demo/logs")); + const empty = await app(authenticated("/api/v2/logout", { method: "POST" })); + + for (const response of [success, error, preflight, empty, stream]) { + expect(response.headers.get(KUBER_VERSION_HEADER)).toBe(KUBER_VERSION); + } + expect([ + success.status, + error.status, + preflight.status, + empty.status, + stream.status, + ]).toEqual([200, 401, 204, 204, 200]); + expect(stream.headers.get("content-type")).toContain("application/x-ndjson"); + expect(await stream.text()).toContain("hello"); +}); + +test("includes the version on exec upgrade errors outside the ordinary API router", () => { + const response = execProblem(new Error("upgrade failed"), "request-id"); + expect(response.status).toBe(500); + expect(response.headers.get(KUBER_VERSION_HEADER)).toBe(KUBER_VERSION); +});