diff --git a/command/main.ts b/command/main.ts index f0da5a1..25ec5a1 100644 --- a/command/main.ts +++ b/command/main.ts @@ -22,7 +22,7 @@ import { trust } from "./trust"; export const main = defineCommand({ meta: { name: "kuber", - version: "2.3.2", + version: "2.3.3", description: "Docker Compose -> K8s translation layer", }, args: { diff --git a/command/up.ts b/command/up.ts index ad8b657..2b6efea 100644 --- a/command/up.ts +++ b/command/up.ts @@ -1,6 +1,7 @@ import { defineCommand } from "citty"; import { Listr, type ListrTaskWrapper } from "listr2"; import { randomUUID } from "node:crypto"; +import { isDeepStrictEqual } from "node:util"; import type { ComposeSpecification } from "../schema/docker.d"; import type { KuberResource } from "../types"; import { @@ -29,6 +30,7 @@ import { type Workspace = { metadata: { name: string; uid: string; resourceVersion: string }; + spec?: { source?: unknown; config?: unknown }; }; type ResourceIdentity = { @@ -195,6 +197,27 @@ function workspaceInput( }; } +function workspaceMatchesInput( + workspace: Workspace, + input: ReturnType, +): boolean { + return ( + isDeepStrictEqual(workspace.spec?.source, input.source) && + isDeepStrictEqual(workspace.spec?.config, input.config) + ); +} + +function isWorkspaceUpdateAmbiguous(error: unknown): boolean { + if (error instanceof KuberApiError) return error.status === 409; + if (error instanceof DOMException && error.name === "AbortError") + return false; + if (error instanceof Error && error.name === "AbortError") return false; + return ( + error instanceof TypeError || + (error instanceof Error && error.name === "TimeoutError") + ); +} + export async function ensureWorkspace( project: string, compose: ComposeSpecification, @@ -203,21 +226,42 @@ export async function ensureWorkspace( ): Promise { const path = `/workspaces/${encodeURIComponent(project)}`; const input = workspaceInput(compose, snapshot); + const update = (workspace: Workspace) => + request(path, { + method: "PUT", + headers: { "if-match": `"${workspace.metadata.resourceVersion}"` }, + json: input, + }); + const recover = async (originalError: unknown): Promise => { + try { + const fresh = await request(path); + if (workspaceMatchesInput(fresh, input)) return fresh; + return await update(fresh); + } catch { + throw originalError; + } + }; let current: Workspace; try { current = await request(path); } catch (error) { if (!(error instanceof KuberApiError) || error.status !== 404) throw error; - return request("/workspaces", { - method: "POST", - json: { id: project, ...input }, - }); + try { + return await request("/workspaces", { + method: "POST", + json: { id: project, ...input }, + }); + } catch (createError) { + if (!isWorkspaceUpdateAmbiguous(createError)) throw createError; + return recover(createError); + } + } + try { + return await update(current); + } catch (updateError) { + if (!isWorkspaceUpdateAmbiguous(updateError)) throw updateError; + return recover(updateError); } - return request(path, { - method: "PUT", - headers: { "if-match": `"${current.metadata.resourceVersion}"` }, - json: input, - }); } function operationEnvironment( diff --git a/package.json b/package.json index ec5d556..a5b69c8 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.3.2", + "version": "2.3.3", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/server/app.ts b/server/app.ts index 8a72db8..90c0710 100644 --- a/server/app.ts +++ b/server/app.ts @@ -57,13 +57,14 @@ import { WorkspaceConflictError, WorkspaceNotFoundError, WorkspaceValidationError, + WORKSPACE_ID_PATTERN, workspaceEtag, type CreateWorkspaceInput, type UpdateWorkspaceInput, type Workspace, type WorkspaceStore, } from "./workspace-store"; -import { redactString } from "./redact"; +import { REDACTED, redactString } from "./redact"; const API_PREFIX = "/api/v2"; const RUNTIME_SESSION_MS = 24 * 60 * 60 * 1000; @@ -107,6 +108,33 @@ export type AppOptions = { hashPassword?: (password: string) => Promise; now?: () => number; requestId?: () => string; + logger?: AppLogger; +}; + +export type UnknownFailureLog = { + event: "request.failed"; + requestId: string; + method: string; + pathname: string; + workspaceId?: string; + status: number; + code: string; + errorName: string; + message: string; + kubernetesStatus?: { status?: string; reason?: string; code?: number }; +}; + +export type RequestErrorLog = Omit & { + event: string; + operationId?: string; +}; + +export interface AppLogger { + error(entry: RequestErrorLog): void; +} + +const defaultAppLogger: AppLogger = { + error: (entry) => console.error(JSON.stringify(entry)), }; type LoginFailures = { count: number; resetAt: number }; @@ -236,6 +264,7 @@ export function createApp( Bun.password.hash(password, { algorithm: "argon2id" })); const now = options.now ?? Date.now; const makeRequestId = options.requestId ?? randomUUID; + const logger = options.logger ?? defaultAppLogger; const bodyLimit = options.jsonBodyLimit ?? DEFAULT_JSON_LIMIT; const allowedOrigins = new Set(options.allowedOrigins ?? []); const hasOriginConfiguration = options.allowedOrigins !== undefined; @@ -282,6 +311,70 @@ export function createApp( return result; } + function workspaceIdFromPath(pathname: string): string | undefined { + const match = new RegExp(`^${API_PREFIX}/workspaces/([^/]+)`).exec( + pathname, + ); + if (!match) return; + try { + const workspaceId = decodeURIComponent(match[1]!); + return WORKSPACE_ID_PATTERN.test(workspaceId) ? workspaceId : undefined; + } catch { + return; + } + } + + function kubernetesStatus( + error: unknown, + ): UnknownFailureLog["kubernetesStatus"] | undefined { + if (!isRecord(error)) return; + const candidate = isRecord(error.status) + ? error.status + : isRecord(error.body) + ? error.body + : undefined; + if (!candidate) return; + const status = typeof candidate.status === "string" ? REDACTED : undefined; + const reason = typeof candidate.reason === "string" ? REDACTED : undefined; + const code = + typeof candidate.code === "number" ? candidate.code : undefined; + return status || reason || code !== undefined + ? { status, reason, code } + : undefined; + } + + function logRequestError( + request: Request, + error: unknown, + { + event, + code, + message, + workspaceId, + operationId, + }: Pick & { + workspaceId?: string; + operationId?: string; + }, + ): void { + const pathname = new URL(request.url).pathname; + logger.error({ + event, + requestId: requestIds.get(request) ?? makeRequestId(), + method: request.method, + pathname, + ...(workspaceId && { workspaceId }), + ...(operationId && { operationId }), + status: 500, + code, + errorName: error instanceof Error ? error.name : "UnknownError", + message: redactString(message), + ...(kubernetesStatus(error) && { + kubernetesStatus: kubernetesStatus(error), + }), + }); + } + async function readJson( request: Request, limit = bodyLimit, @@ -911,10 +1004,13 @@ export function createApp( operation.metadata.name, ); } catch (error) { - console.error( - "Failed to append successful operation audit event", - error, - ); + logRequestError(request, error, { + event: "operation.audit.failed", + code: "AUDIT_APPEND_FAILED", + message: "Successful operation audit event could not be appended", + workspaceId: workspace.metadata.name, + operationId: operation.metadata.name, + }); } return response(operationBody(completed, result), 200, { location: `${API_PREFIX}/operations/${operation.metadata.name}`, @@ -948,10 +1044,13 @@ export function createApp( operation.metadata.name, ); } catch (auditError) { - console.error( - "Failed to append failed operation audit event", - auditError, - ); + logRequestError(request, auditError, { + event: "operation.audit.failed", + code: "AUDIT_APPEND_FAILED", + message: "Failed operation audit event could not be appended", + workspaceId: workspace.metadata.name, + operationId: operation.metadata.name, + }); } throw new HttpError( failureCode === "WORKSPACE_LEASE_LOST" ? 409 : 500, @@ -968,16 +1067,25 @@ export function createApp( try { await lease?.release(); } catch (error) { - console.error("Failed to release workspace operation lease", error); + logRequestError(request, error, { + event: "operation.lease_release.failed", + code: "LEASE_RELEASE_FAILED", + message: "Workspace operation lease could not be released", + workspaceId: workspace.metadata.name, + operationId: operation.metadata.name, + }); } } }; if (request.headers.get("prefer") === "respond-async") { void executeClaimed().catch((error) => { - console.error( - `Background operation '${operation.metadata.name}' failed`, - error, - ); + logRequestError(request, error, { + event: "operation.background.failed", + code: "BACKGROUND_OPERATION_FAILED", + message: "Background operation could not be completed", + workspaceId: workspace.metadata.name, + operationId: operation.metadata.name, + }); }); return response( { @@ -2122,7 +2230,6 @@ export function createApp( ); if (error instanceof Error && /already exists/i.test(error.message)) return new HttpError(409, "Conflict", "ALREADY_EXISTS", error.message); - console.error(error); return new HttpError( 500, "Internal server error", @@ -2195,7 +2302,15 @@ export function createApp( result = await handleAuthenticated(request, url, identity); } } catch (error) { - result = problem(normalizeError(error), requestId); + const normalized = normalizeError(error); + if (normalized.code === "INTERNAL_ERROR") + logRequestError(request, error, { + event: "request.failed", + code: normalized.code, + message: normalized.message, + workspaceId: workspaceIdFromPath(url.pathname), + }); + result = problem(normalized, requestId); } result.headers.set("x-request-id", requestId); if (origin && originAllowed) { @@ -2220,13 +2335,14 @@ export type ExecConnection = { export type ExecAuthOptions = Pick< AppOptions, - "store" | "workspaceStore" | "auditStore" | "now" + "store" | "workspaceStore" | "auditStore" | "now" | "logger" | "requestId" >; export async function authorizeExecConnection( options: ExecAuthOptions, request: Request, workspaceId: string, + requestId?: string, ): Promise { const now = options.now ?? Date.now; const identity = await authenticateRequest(options, request, now); @@ -2261,7 +2377,22 @@ export async function authorizeExecConnection( workspaceId, }); } catch (error) { - console.error("Failed to append denied exec audit event", error); + const suppliedRequestId = request.headers.get("x-request-id"); + (options.logger ?? defaultAppLogger).error({ + event: "exec.audit.failed", + requestId: + requestId ?? + (suppliedRequestId && /^[\x21-\x7e]{1,128}$/.test(suppliedRequestId) + ? suppliedRequestId + : (options.requestId ?? randomUUID)()), + method: request.method, + pathname: new URL(request.url).pathname, + ...(WORKSPACE_ID_PATTERN.test(workspaceId) && { workspaceId }), + status: 500, + code: "AUDIT_APPEND_FAILED", + errorName: error instanceof Error ? error.name : "UnknownError", + message: "Denied exec audit event could not be appended", + }); } } throw new HttpError( @@ -2289,7 +2420,6 @@ export function execProblem(error: unknown, requestId = ""): Response { let httpError: HttpError; if (error instanceof HttpError) httpError = error; else { - console.error(error); httpError = new HttpError( 500, "Internal server error", @@ -2571,11 +2701,33 @@ export async function handleExecUpgrade( workspaceId: string, upgrade: (request: Request, data: ExecConnection) => boolean, ): Promise { + const suppliedRequestId = request.headers.get("x-request-id"); + const requestId = + suppliedRequestId && /^[\x21-\x7e]{1,128}$/.test(suppliedRequestId) + ? suppliedRequestId + : (options.requestId ?? randomUUID)(); let connection: ExecConnection; try { - connection = await authorizeExecConnection(options, request, workspaceId); + connection = await authorizeExecConnection( + options, + request, + workspaceId, + requestId, + ); } catch (error) { - return execProblem(error); + if (!(error instanceof HttpError)) + (options.logger ?? defaultAppLogger).error({ + event: "request.failed", + requestId, + method: request.method, + pathname: new URL(request.url).pathname, + ...(WORKSPACE_ID_PATTERN.test(workspaceId) && { workspaceId }), + status: 500, + code: "INTERNAL_ERROR", + errorName: error instanceof Error ? error.name : "UnknownError", + message: "The request could not be completed", + }); + return execProblem(error, requestId); } if (!upgrade(request, connection)) return execProblem( @@ -2585,6 +2737,7 @@ export async function handleExecUpgrade( "UPGRADE_FAILED", "WebSocket upgrade failed", ), + requestId, ); return; } diff --git a/tests/command/up-api.test.ts b/tests/command/up-api.test.ts index 7f0d7cc..8d17634 100644 --- a/tests/command/up-api.test.ts +++ b/tests/command/up-api.test.ts @@ -158,6 +158,163 @@ describe("up API pipeline", () => { expect(new Headers(calls[1]?.init?.headers).get("if-match")).toBe('"7"'); }); + test("retries one conflict using a fresh ETag", async () => { + const calls: Array<{ path: string; init?: ApiRequestInit }> = []; + let gets = 0; + const request: ApiRequester = async ( + path: string, + init?: ApiRequestInit, + ) => { + calls.push({ path, init }); + if (!init) { + gets += 1; + return { + metadata: { + name: "shop", + uid: "uid", + resourceVersion: gets === 1 ? "7" : "8", + }, + spec: { source: { uri: "cas://old", digest: "old" }, config: {} }, + } as T; + } + if (calls.length === 2) + throw new KuberApiError("conflict", 409, { + title: "Conflict", + status: 409, + code: "WORKSPACE_CONFLICT", + }); + return { + metadata: { name: "shop", uid: "uid", resourceVersion: "9" }, + } as T; + }; + + await ensureWorkspace("shop", { services: {} }, snapshot, request); + + expect(calls.map(({ path }) => path)).toEqual([ + "/workspaces/shop", + "/workspaces/shop", + "/workspaces/shop", + "/workspaces/shop", + ]); + expect(new Headers(calls[3]?.init?.headers).get("if-match")).toBe('"8"'); + }); + + test("accepts a fresh workspace that already matches after a conflict", async () => { + let gets = 0; + let puts = 0; + const desired = { + source: { uri: `cas://${snapshot.digest}`, digest: snapshot.digest }, + config: { compose: { services: {} } }, + }; + const request: ApiRequester = async ( + _path: string, + init?: ApiRequestInit, + ) => { + if (!init) { + gets += 1; + return { + metadata: { name: "shop", uid: "uid", resourceVersion: "8" }, + spec: + gets === 1 + ? { ...desired, source: { uri: "cas://old", digest: "old" } } + : desired, + } as T; + } + puts += 1; + throw new KuberApiError("conflict", 409, { + title: "Conflict", + status: 409, + code: "WORKSPACE_CONFLICT", + }); + }; + + await ensureWorkspace("shop", { services: {} }, snapshot, request); + + expect({ gets, puts }).toEqual({ gets: 2, puts: 1 }); + }); + + test("recovers an ambiguous transport failure after the update committed", async () => { + let gets = 0; + let puts = 0; + const desired = { + source: { uri: `cas://${snapshot.digest}`, digest: snapshot.digest }, + config: { compose: { services: {} } }, + }; + const request: ApiRequester = async ( + _path: string, + init?: ApiRequestInit, + ) => { + if (!init) { + gets += 1; + return { + metadata: { name: "shop", uid: "uid", resourceVersion: "8" }, + spec: + gets === 1 + ? { ...desired, source: { uri: "cas://old", digest: "old" } } + : desired, + } as T; + } + puts += 1; + throw new TypeError("connection reset"); + }; + + await ensureWorkspace("shop", { services: {} }, snapshot, request); + + expect({ gets, puts }).toEqual({ gets: 2, puts: 1 }); + }); + + test("recovers an ambiguous transport failure with one fresh-ETag retry", async () => { + const etags: string[] = []; + let gets = 0; + const request: ApiRequester = async ( + _path: string, + init?: ApiRequestInit, + ) => { + if (!init) { + gets += 1; + return { + metadata: { + name: "shop", + uid: "uid", + resourceVersion: gets === 1 ? "7" : "8", + }, + spec: { source: { uri: "cas://old", digest: "old" }, config: {} }, + } as T; + } + etags.push(new Headers(init.headers).get("if-match")!); + if (etags.length === 1) throw new TypeError("connection reset"); + return { + metadata: { name: "shop", uid: "uid", resourceVersion: "9" }, + } as T; + }; + + await ensureWorkspace("shop", { services: {} }, snapshot, request); + + expect(etags).toEqual(['"7"', '"8"']); + }); + + test("does not retry a generic workspace update failure", async () => { + const failure = new KuberApiError("failed", 500, { + title: "Internal server error", + status: 500, + code: "INTERNAL_ERROR", + }); + let calls = 0; + const request: ApiRequester = async () => { + calls += 1; + if (calls === 1) + return { + metadata: { name: "shop", uid: "uid", resourceVersion: "7" }, + } as T; + throw failure; + }; + + await expect( + ensureWorkspace("shop", { services: {} }, snapshot, request), + ).rejects.toBe(failure); + expect(calls).toBe(2); + }); + test("plans, applies, runs the hook, waits, then deletes stale identities", async () => { const order: string[] = []; const desired = [ diff --git a/tests/server/app.test.ts b/tests/server/app.test.ts index 5622687..409e690 100644 --- a/tests/server/app.test.ts +++ b/tests/server/app.test.ts @@ -681,6 +681,83 @@ describe("kuber v2 HTTP routes", () => { expect(await operationStore.list("demo")).toHaveLength(1); }); + test("redacts provider status and reason in correlated generic failure logs", async () => { + const workspaceStore = new MemoryWorkspaceStore({ + uid: () => "workspace-uid", + }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const logs: unknown[] = []; + spyOn(workspaceStore, "update").mockRejectedValue( + Object.assign( + new Error( + 'provider failed password=top-secret config={"compose":"private"}', + ), + { + name: "KubernetesError", + body: { + status: "Failure", + reason: "InternalError", + code: 500, + message: "contains top-secret", + }, + }, + ), + ); + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + requestId: () => "request-123", + logger: { error: (entry) => logs.push(entry) }, + }); + + const result = await app( + request( + "/api/v2/workspaces/demo", + { + method: "PUT", + headers: { "if-match": '"1"' }, + body: JSON.stringify({ + source: { uri: "cas://secret-source", digest: "secret-digest" }, + config: { compose: "private-config", password: "top-secret" }, + }), + }, + "token", + ), + ); + + expect(result.status).toBe(500); + expect(await result.json()).toMatchObject({ + code: "INTERNAL_ERROR", + requestId: "request-123", + }); + expect(logs).toEqual([ + { + event: "request.failed", + requestId: "request-123", + method: "PUT", + pathname: "/api/v2/workspaces/demo", + workspaceId: "demo", + status: 500, + code: "INTERNAL_ERROR", + errorName: "KubernetesError", + message: "The request could not be completed", + kubernetesStatus: { + status: "[REDACTED]", + reason: "[REDACTED]", + code: 500, + }, + }, + ]); + expect(JSON.stringify(logs)).not.toContain("top-secret"); + expect(JSON.stringify(logs)).not.toContain("private-config"); + expect(JSON.stringify(logs)).not.toContain("secret-source"); + expect(JSON.stringify(logs)).not.toContain('"Failure"'); + expect(JSON.stringify(logs)).not.toContain('"InternalError"'); + }); + test("lists persisted operation events with capability and workspace scope checks", async () => { const store = await authenticatedStore("operator"); const workspaceStore = new MemoryWorkspaceStore({ @@ -801,7 +878,7 @@ describe("kuber v2 HTTP routes", () => { }); test("does not expose request internals or corrupt success when auditing fails", async () => { - const reported = spyOn(console, "error").mockImplementation(() => {}); + const logs: unknown[] = []; const workspaceStore = new MemoryWorkspaceStore({ uid: () => "workspace-uid", }); @@ -824,6 +901,8 @@ describe("kuber v2 HTTP routes", () => { }, list: async () => [], }, + requestId: () => "audit-request-123", + logger: { error: (entry) => logs.push(entry) }, }); const result = await app( request( @@ -843,8 +922,23 @@ describe("kuber v2 HTTP routes", () => { action: "workspace.stop", }); expect(body.operation.status.state).toBe("succeeded"); - expect(reported).toHaveBeenCalledTimes(1); - reported.mockRestore(); + expect(logs).toEqual([ + { + event: "operation.audit.failed", + requestId: "audit-request-123", + method: "POST", + pathname: "/api/v2/workspaces/demo/lifecycle", + workspaceId: "demo", + operationId: "operation-operation-uid", + status: 500, + code: "AUDIT_APPEND_FAILED", + errorName: "Error", + message: "Successful operation audit event could not be appended", + }, + ]); + expect(JSON.stringify(logs)).not.toContain("audit unavailable"); + expect(JSON.stringify(logs)).not.toContain("private-request"); + expect(JSON.stringify(logs)).not.toContain("do-not-store"); expect( await operationStore.get("operation-operation-uid"), ).not.toHaveProperty("spec.request");