import type { KubernetesObject } from "@kubernetes/client-node"; import { randomUUID } from "node:crypto"; import type { ComposeSpecification } from "../schema/docker.d"; import type { Operation as PublicOperation } from "../shared/api"; import type { BuildRequest, Sha256Digest } from "../shared/build-protocol"; import { createToken, hashToken, tokenHashesEqual, type AuthStore, type KuberUser, type ApiKeyRecord, type SessionRecord, } from "./auth"; import { hasCapability, isCapability, isRole, type Capability, type Role, } from "./authorization"; import type { AuditStore } from "./audit-store"; import { validateTrust, type TrustStore } from "./trust-store"; import { BuildConflictError, BuildNotFoundError, BuildValidationError, type BuildController, } from "./build-controller"; import { KubernetesLogError, type LogService } from "./log-service"; import { ExecService, type ExecClientFrame, type ExecServerFrame, } from "./exec-service"; import type { ExecClientWireFrame, ExecServerWireFrame, ExecStartFrame, } from "../lib/exec-api"; import type { ManagementService, OperationProgressEmitter, ResourceIdentity, } from "./management"; import { OperationConflictError, OperationNotFoundError, OperationValidationError, sanitizeOperationError, sanitizeOperationResult, type Operation, type OperationStore, type WorkspaceLeaseProvider, } from "./operation-store"; import { WorkspaceConflictError, WorkspaceNotFoundError, WorkspaceValidationError, WORKSPACE_ID_PATTERN, workspaceEtag, type CreateWorkspaceInput, type UpdateWorkspaceInput, type Workspace, type WorkspaceStore, } from "./workspace-store"; import { redactString } from "./redact"; import { logValue, logServerRequest, processLogger, safeLog, type ProcessLogEntry, } from "../lib/request-log"; import { MaintenanceBusyError, normalizeMaintenanceHost, type MaintenanceService, } from "./maintenance"; const API_PREFIX = "/api/v2"; const RUNTIME_SESSION_MS = 24 * 60 * 60 * 1000; const PERSISTENT_SESSION_MS = 30 * 24 * 60 * 60 * 1000; const LOGIN_WINDOW_MS = 5 * 60 * 1000; const MAX_LOGIN_FAILURES = 5; const DEFAULT_JSON_LIMIT = 1024 * 1024; const WORKSPACE_LEASE_TTL_MS = 30_000; const WORKSPACE_LEASE_RENEW_INTERVAL_MS = WORKSPACE_LEASE_TTL_MS / 3; const DEFAULT_API_KEY_MS = 90 * 24 * 60 * 60 * 1000; const MAX_API_KEY_MS = 365 * 24 * 60 * 60 * 1000; export interface ApiWorkspaceStore extends WorkspaceStore { delete?(id: string): Promise; adopt?( workspaceId: string, workspaceUid: string, ): Promise; adoptPlatform?(workspaceUid: string): Promise; } export type AppOptions = { store: AuthStore; workspaceStore?: ApiWorkspaceStore; operationStore?: OperationStore; auditStore?: AuditStore; trustStore?: TrustStore; management?: ManagementService; builds?: BuildController; logs?: LogService; execService?: ExecService; maintenance?: MaintenanceService; resolveImage?: ( project: string, service: string, ) => Promise<{ image: string; digest: Sha256Digest; reference: string }>; leases?: WorkspaceLeaseProvider; adoption?: WorkspaceAdoptionService; allowedOrigins?: readonly string[]; jsonBodyLimit?: number; verifyPassword?: (password: string, hash: string) => Promise; 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; stack?: string; kubernetesStatus?: { status?: string; reason?: string; code?: number }; }; export type RequestErrorLog = Omit & { event: string; operationId?: string; }; export interface AppLogger { error(entry: RequestErrorLog): void; log?(entry: ProcessLogEntry): void; } const defaultAppLogger: AppLogger = { error: (entry) => console.error(JSON.stringify(entry)), log: processLogger.log, }; type LoginFailures = { count: number; resetAt: number }; type Identity = { user: KuberUser; session?: SessionRecord; apiKey?: ApiKeyRecord; }; export type WorkspaceAdoptionResult = { workspaceId: string; workspaceUid: string; resourcesAdopted: number; }; export interface WorkspaceAdoptionService { adopt( workspaceId: string, workspaceUid: string, ): Promise; adoptPlatform(workspaceUid: string): Promise; } export class WorkspaceAdoptionError extends Error { readonly code = "ADOPTION_CONFLICT"; } export async function cleanupExpiredSessions( store: AuthStore, now = Date.now(), ): Promise { const sessions = await store.deleteExpiredSessions(now); return sessions + (await store.deleteExpiredApiKeys(now)); } class HttpError extends Error { constructor( readonly status: number, readonly title: string, readonly code: string, detail: string, readonly headers?: Record, readonly operationId?: string, ) { super(detail); } } function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function bearerToken(request: Request): string | undefined { return /^Bearer\s+(.+)$/i.exec( request.headers.get("authorization") ?? "", )?.[1]; } function publicUser(user: KuberUser) { return { username: user.username, roles: user.roles, disabled: Boolean(user.disabled), }; } function workspaceIdentity(workspace: Workspace) { return { project: workspace.metadata.name, uid: workspace.metadata.uid }; } export async function authenticateRequest( options: Pick, request: Request, now: () => number = options.now ?? Date.now, ): Promise { const token = bearerToken(request); if (!token) return; const tokenHash = hashToken(token); const session = await options.store.getSession(tokenHash); if (session) { if ( !tokenHashesEqual(tokenHash, session.tokenHash) || Date.parse(session.expiresAt) <= now() ) { await options.store.deleteSession(tokenHash); return; } const user = await options.store.getUser(session.username); if (!user || user.disabled || user.authVersion !== session.authVersion) return; return { user, session }; } const apiKey = await options.store.getApiKey(tokenHash); if ( !apiKey || !tokenHashesEqual(tokenHash, apiKey.tokenHash) || Date.parse(apiKey.expiresAt) <= now() ) return; const user = await options.store.getUser(apiKey.username); if (!user || user.disabled) return; return { user, apiKey }; } function pathPart(value: string): string { try { return decodeURIComponent(value); } catch { throw new HttpError( 400, "Bad request", "INVALID_PATH", "Invalid URL encoding", ); } } export function createApp( options: AppOptions, ): (request: Request) => Promise { const verifyPassword = options.verifyPassword ?? ((password: string, hash: string) => Bun.password.verify(password, hash)); const hashPassword = options.hashPassword ?? ((password: string) => 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; const loginFailures = new Map(); const requestIds = new WeakMap(); function loginKey(request: Request): string { return ( request.headers.get("cf-connecting-ip") ?? request.headers.get("x-forwarded-for")?.split(",")[0]?.trim() ?? "unknown" ); } function response( value: unknown, status = 200, headers?: Headers | Record, ): Response { return Response.json(value, { status, headers: { "cache-control": "no-store", ...Object.fromEntries(new Headers(headers)), }, }); } function problem(error: HttpError, requestId: string): Response { const result = response( { type: `https://kuber.astrxl.dev/problems/${error.code.toLowerCase()}`, title: error.title, status: error.status, detail: redactString(error.message), code: error.code, requestId, ...(error.operationId && { operationId: error.operationId }), }, error.status, error.headers, ); result.headers.set("content-type", "application/problem+json"); 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" ? candidate.status : undefined; const reason = typeof candidate.reason === "string" ? candidate.reason : undefined; const code = typeof candidate.code === "number" ? candidate.code : undefined; return status || reason || code !== undefined ? { status, reason, code } : undefined; } function kubernetesError( error: unknown, ): { status: number; detail: string } | undefined { if (!isRecord(error) || typeof error.statusCode !== "number") return; const status = error.statusCode; if (!Number.isInteger(status) || status < 100 || status > 599) return; const body = isRecord(error.body) ? error.body : undefined; if ( !body || body.kind !== "Status" || body.apiVersion !== "v1" || body.status !== "Failure" || typeof body.reason !== "string" || typeof body.message !== "string" || !Number.isInteger(body.code) || body.code !== status ) return; return { status, detail: body.message }; } function errorDiagnostics(error: unknown): { errorName: string; message: string; stack?: string; } { if (error instanceof Error) return { errorName: error.name, message: error.message, ...(error.stack && { stack: error.stack }), }; if (isRecord(error)) return { errorName: typeof error.name === "string" ? error.name : "UnknownError", message: typeof error.message === "string" ? error.message : "Unknown error", ...(typeof error.stack === "string" && { stack: error.stack }), }; return { errorName: "UnknownError", message: typeof error === "string" ? error : "Unknown error", }; } function logRequestError( request: Request, error: unknown, { event, code, message, workspaceId, operationId, }: Pick & { workspaceId?: string; operationId?: string; }, ): void { const pathname = new URL(request.url).pathname; const diagnostics = errorDiagnostics(error); const unknownFailure = event === "request.failed" && code === "INTERNAL_ERROR"; logger.error({ event, requestId: requestIds.get(request) ?? makeRequestId(), method: request.method, pathname, ...(workspaceId && { workspaceId }), ...(operationId && { operationId }), status: 500, code, errorName: diagnostics.errorName, message: unknownFailure ? diagnostics.message : redactString(message), ...(unknownFailure && diagnostics.stack && { stack: diagnostics.stack }), ...(unknownFailure && kubernetesStatus(error) && { kubernetesStatus: kubernetesStatus(error), }), }); } function logRequestBody(request: Request, body: Uint8Array): void { safeLog( { log: (entry) => (logger.log ?? processLogger.log)(entry) }, { event: "kuber.server.request.body", type: "server.request", requestId: requestIds.get(request) ?? makeRequestId(), method: request.method, url: request.url, pathname: new URL(request.url).pathname, headers: Object.fromEntries(request.headers), body: logValue(body), }, ); } async function readJson( request: Request, limit = bodyLimit, ): Promise> { const length = Number(request.headers.get("content-length")); if (Number.isFinite(length) && length > limit) throw new HttpError( 413, "Payload too large", "BODY_TOO_LARGE", `JSON body exceeds ${limit} bytes`, ); const contentType = request.headers.get("content-type"); if ( contentType && !/^application\/(?:[\w.+-]+\+)?json(?:\s*;|$)/i.test(contentType) ) throw new HttpError( 415, "Unsupported media type", "UNSUPPORTED_MEDIA_TYPE", "Request body must be JSON", ); const reader = request.body?.getReader(); const chunks: Uint8Array[] = []; let bytes = 0; if (reader) { while (true) { const { done, value } = await reader.read(); if (done) break; bytes += value.byteLength; if (bytes > limit) { await reader.cancel(); throw new HttpError( 413, "Payload too large", "BODY_TOO_LARGE", `JSON body exceeds ${limit} bytes`, ); } chunks.push(value); } } const payload = new Uint8Array(bytes); let offset = 0; for (const chunk of chunks) { payload.set(chunk, offset); offset += chunk.byteLength; } const text = new TextDecoder().decode(payload); logRequestBody(request, payload); let value: unknown; try { value = text ? JSON.parse(text) : {}; } catch { throw new HttpError( 400, "Invalid JSON", "INVALID_JSON", "Request body is not valid JSON", ); } if (!isRecord(value)) throw new HttpError( 400, "Invalid request", "INVALID_BODY", "JSON body must be an object", ); return value; } async function readBytes( request: Request, limit: number, ): Promise { const length = Number(request.headers.get("content-length")); if (Number.isFinite(length) && length > limit) throw new HttpError( 413, "Payload too large", "BODY_TOO_LARGE", `Request body exceeds ${limit} bytes`, ); const reader = request.body?.getReader(); const chunks: Uint8Array[] = []; let size = 0; if (reader) { for (;;) { const { done, value } = await reader.read(); if (done) break; size += value.byteLength; if (size > limit) { await reader.cancel(); throw new HttpError( 413, "Payload too large", "BODY_TOO_LARGE", `Request body exceeds ${limit} bytes`, ); } chunks.push(value); } } const result = new Uint8Array(size); let offset = 0; for (const chunk of chunks) { result.set(chunk, offset); offset += chunk.byteLength; } logRequestBody(request, result); return result; } function requireBuilds(): BuildController { if (!options.builds) throw new HttpError( 503, "Service unavailable", "BUILDS_UNAVAILABLE", "Build service is not configured", ); return options.builds; } function nonNegativeInteger( value: string | null, name: string, ): number | undefined { if (value === null) return; const number = Number(value); if (!Number.isSafeInteger(number) || number < 0) throw new HttpError( 400, "Invalid request", "INVALID_QUERY", `${name} must be a non-negative integer`, ); return number; } function ndjson( events: AsyncIterable | Iterable, signal?: AbortSignal, ): Response { const iterator = Symbol.asyncIterator in Object(events) ? (events as AsyncIterable)[Symbol.asyncIterator]() : (async function* () { yield* events as Iterable; })(); const encoder = new TextEncoder(); let reading = false; return new Response( new ReadableStream({ async pull(controller) { if (reading) return; reading = true; try { const next = await iterator.next(); if (next.done) controller.close(); else controller.enqueue( encoder.encode(`${JSON.stringify(next.value)}\n`), ); } catch (error) { controller.error(error); } finally { reading = false; } }, async cancel(reason) { await iterator.return?.(reason); }, }), { headers: { "content-type": "application/x-ndjson; charset=utf-8", "cache-control": "no-store", connection: "keep-alive", ...(signal ? { "x-accel-buffering": "no" } : {}), }, }, ); } async function authenticate(request: Request): Promise { return authenticateRequest(options, request, now); } function actor(identity: Identity, request: Request) { return { username: identity.user.username, roles: identity.user.roles, ip: request.headers.get("cf-connecting-ip") ?? request.headers.get("x-forwarded-for")?.split(",")[0]?.trim(), userAgent: request.headers.get("user-agent") ?? undefined, }; } async function audit( identity: Identity, request: Request, action: string, outcome: "success" | "failure" | "denied", details?: unknown, workspaceId?: string, operationId?: string, ): Promise { if (!options.auditStore) return; await options.auditStore.append({ actor: actor(identity, request), action, outcome, details, workspaceId, operationId, }); } async function requireCapability( identity: Identity, request: Request, capability: Capability, ): Promise { if ( identity.apiKey ? identity.apiKey.capabilities.includes(capability) : hasCapability(identity.user.roles, capability) ) return; await audit(identity, request, `authorization.${capability}`, "denied"); throw new HttpError( 403, "Forbidden", "FORBIDDEN", `Capability '${capability}' is required`, ); } async function requireWorkspaceScope( identity: Identity, request: Request, workspace: string, ): Promise { if (!identity.apiKey?.workspace || identity.apiKey.workspace === workspace) return; await audit(identity, request, "authorization.workspace", "denied", { workspace, }); throw new HttpError( 403, "Forbidden", "FORBIDDEN", "This API key is restricted to a different workspace", ); } async function requireApiKeyDelegation( identity: Identity, request: Request, username: string, capabilities: readonly Capability[], workspace: string | undefined, ): Promise { const parent = identity.apiKey; if (!parent) return; let reason: "target_user" | "capabilities" | "workspace" | undefined; if (username !== parent.username) reason = "target_user"; else if ( !capabilities.every((capability) => parent.capabilities.includes(capability), ) ) reason = "capabilities"; else if (parent.workspace !== undefined && workspace !== parent.workspace) reason = "workspace"; if (!reason) return; await audit(identity, request, "api_key.create", "denied", { reason, username, capabilities, ...(workspace && { workspace }), }); throw new HttpError( 403, "Forbidden", "API_KEY_DELEGATION_FORBIDDEN", "API key children must use the caller's user, capabilities, and workspace scope", ); } function requireWorkspaceStore(): ApiWorkspaceStore { if (!options.workspaceStore) throw new HttpError( 503, "Service unavailable", "WORKSPACE_STORE_UNAVAILABLE", "Workspace storage is not configured", ); return options.workspaceStore; } function requireManagement(): ManagementService { if (!options.management) throw new HttpError( 503, "Service unavailable", "MANAGEMENT_UNAVAILABLE", "Kubernetes management is not configured", ); return options.management; } function requireTrustStore(): TrustStore { if (!options.trustStore) throw new HttpError( 503, "Service unavailable", "TRUST_STORE_UNAVAILABLE", "Namespace trust service is unavailable. Contact your Kuber administrator.", ); return options.trustStore; } async function requireTrust( request: Request, project: string, ): Promise { const fingerprint = request.headers.get("x-kuber-trust-fingerprint"); const suppliedProject = request.headers.get("x-kuber-trust-project"); if (!fingerprint || suppliedProject !== project) throw new HttpError( 428, "Trusted workspace required", "TRUST_REQUIRED", "This workspace reconciliation requires trust. Run kuber trust and retry.", ); if (!(await requireTrustStore().has(project, fingerprint))) throw new HttpError( 403, "Trusted workspace rejected", "TRUST_REQUIRED", "This directory is not registered for this namespace.", ); } function requireAdoption(): WorkspaceAdoptionService { const workspaceStore = options.workspaceStore; const adoption = options.adoption ?? (workspaceStore?.adopt && workspaceStore.adoptPlatform ? { adopt: workspaceStore.adopt.bind(workspaceStore), adoptPlatform: workspaceStore.adoptPlatform.bind(workspaceStore), } : undefined); if (!adoption) throw new HttpError( 503, "Service unavailable", "ADOPTION_UNAVAILABLE", "Workspace adoption is not configured", ); return adoption; } async function getWorkspace(id: string): Promise { const workspace = await requireWorkspaceStore().get(id); if (!workspace) throw new HttpError( 404, "Not found", "WORKSPACE_NOT_FOUND", `Workspace '${id}' not found`, ); return workspace; } function idempotencyKey(request: Request): string { const key = request.headers.get("idempotency-key")?.trim(); return key || requestIds.get(request) || makeRequestId(); } function publicOperation(operation: Operation): PublicOperation { return { apiVersion: operation.apiVersion, kind: operation.kind, metadata: operation.metadata, spec: { workspaceId: operation.spec.workspaceId, action: operation.spec.action, }, status: { ...(() => { const { events: _events, ...status } = operation.status; return status; })(), ...(operation.status.error && { error: sanitizeOperationError( operation.status.error, operation.spec.action, ), }), ...(operation.status.result !== undefined && { result: sanitizeOperationResult( operation.status.result, operation.spec.action, ), }), }, }; } function operationBody(operation: Operation, result: unknown) { const visible = publicOperation(operation); const safeResult = sanitizeOperationResult(result, operation.spec.action); return isRecord(safeResult) ? { ...safeResult, operationId: operation.metadata.name, operation: visible, } : Array.isArray(safeResult) ? { deployments: safeResult, operationId: operation.metadata.name, operation: visible, } : { result: safeResult, operationId: operation.metadata.name, operation: visible, }; } async function transitionOperationToFailure( operationId: string, error: NonNullable, ): Promise { try { return await options.operationStore!.transition(operationId, "failed", { error, }); } catch (transitionError) { if (!(transitionError instanceof OperationConflictError)) throw transitionError; const latest = await options.operationStore!.get(operationId); if (latest?.status.state === "failed") return latest; throw transitionError; } } function operationFailureError(operation: Operation): HttpError { const failure = operation.status.error ? sanitizeOperationError(operation.status.error, operation.spec.action) : { code: "OPERATION_FAILED", message: "Operation failed" }; return new HttpError( failure.code === "WORKSPACE_BUSY" || failure.code === "WORKSPACE_LEASE_LOST" ? 409 : 500, "Operation failed", failure.code, failure.message, undefined, operation.metadata.name, ); } async function runOperation( identity: Identity, request: Request, workspace: Workspace, action: string, input: unknown, execute: ( signal: AbortSignal, emit: OperationProgressEmitter, ) => Promise, ): Promise { if (!options.operationStore) throw new HttpError( 503, "Service unavailable", "OPERATION_STORE_UNAVAILABLE", "Operation storage is not configured", ); const { operation } = await options.operationStore.createOrReuse({ workspaceId: workspace.metadata.name, workspaceUid: workspace.metadata.uid, action, idempotencyKey: idempotencyKey(request), request: input, }); if (operation.status.state === "failed") throw operationFailureError(operation); if (operation.status.state === "cancelled") throw new HttpError( 409, "Operation cancelled", "OPERATION_CANCELLED", "The idempotent operation was cancelled", undefined, operation.metadata.name, ); if (operation.status.state === "running") return response( { operationId: operation.metadata.name, operation: publicOperation(operation), }, 202, { location: `${API_PREFIX}/operations/${operation.metadata.name}` }, ); if (operation.status.state === "succeeded") return response(operationBody(operation, operation.status.result), 200, { location: `${API_PREFIX}/operations/${operation.metadata.name}`, }); const lease = options.leases ? await options.leases.acquire( workspace.metadata.name, operation.metadata.name, ) : undefined; if (options.leases && !lease) { throw new HttpError( 409, "Conflict", "WORKSPACE_BUSY", "Another workspace operation is running", undefined, operation.metadata.name, ); } if (!lease && options.leases) throw new Error("Unreachable lease state"); let operationStarted = false; let leaseRenewalTimer: ReturnType | undefined; let leaseOwnershipLost = false; let renewalInFlight: Promise | undefined; const executionController = new AbortController(); const renewLease = async (): Promise => { if (!lease || leaseOwnershipLost) return false; if (renewalInFlight) return renewalInFlight; renewalInFlight = lease .renew(WORKSPACE_LEASE_TTL_MS) .catch(() => false) .then((renewed) => { if (!renewed) { leaseOwnershipLost = true; executionController.abort("Workspace lease ownership was lost"); } return renewed; }) .finally(() => { renewalInFlight = undefined; }); return renewalInFlight; }; const scheduleLeaseRenewal = () => { if (!lease || leaseOwnershipLost) return; leaseRenewalTimer = setTimeout(() => { void renewLease().finally(scheduleLeaseRenewal); }, WORKSPACE_LEASE_RENEW_INTERVAL_MS); }; const requireLeaseOwnership = async () => { if (lease && !(await renewLease())) throw new HttpError( 409, "Conflict", "WORKSPACE_LEASE_LOST", "Workspace operation lease ownership was lost", undefined, operation.metadata.name, ); }; const executeClaimed = async (): Promise => { try { const claimed = await options.operationStore!.claimExecution( operation.metadata.name, ); if (!claimed) { const latest = await options.operationStore!.get( operation.metadata.name, ); if (latest?.status.state === "succeeded") return response(operationBody(latest, latest.status.result), 200, { location: `${API_PREFIX}/operations/${operation.metadata.name}`, }); if (latest?.status.state === "failed") throw operationFailureError(latest); return response( { operationId: operation.metadata.name, operation: publicOperation(latest ?? operation), }, 202, { location: `${API_PREFIX}/operations/${operation.metadata.name}` }, ); } operationStarted = true; scheduleLeaseRenewal(); await requireLeaseOwnership(); const result = await execute( executionController.signal, async (event) => { await options.operationStore!.emit(operation.metadata.name, event); }, ); await requireLeaseOwnership(); const completed = await options.operationStore!.transition( operation.metadata.name, "succeeded", { result }, ); try { await audit( identity, request, action, "success", undefined, workspace.metadata.name, operation.metadata.name, ); } catch (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}`, }); } 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 failed = await transitionOperationToFailure( operation.metadata.name, { code: leaseLost ? "WORKSPACE_LEASE_LOST" : "OPERATION_FAILED", message: leaseLost ? "Workspace operation lease ownership was lost" : message, }, ); const failure = failed.status.error; const failureCode = failure?.code ?? "OPERATION_FAILED"; const failureMessage = failure?.message ?? message; try { await audit( identity, request, action, "failure", { error: failureMessage }, workspace.metadata.name, operation.metadata.name, ); } catch (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, failureCode === "WORKSPACE_LEASE_LOST" ? "Conflict" : "Operation failed", failureCode, failureMessage, undefined, operation.metadata.name, ); } finally { if (leaseRenewalTimer !== undefined) clearTimeout(leaseRenewalTimer); try { await lease?.release(); } catch (error) { logRequestError(request, error, { event: "operation.lease_release.failed", code: "LEASE_RELEASE_FAILED", message: "Workspace operation lease could not be released", workspaceId: workspace.metadata.name, operationId: operation.metadata.name, }); } } }; if (request.headers.get("prefer") === "respond-async") { void executeClaimed().catch((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( { operationId: operation.metadata.name, operation: publicOperation(operation), }, 202, { location: `${API_PREFIX}/operations/${operation.metadata.name}` }, ); } return executeClaimed(); } async function handleLogin(request: Request): Promise { const key = loginKey(request); const previous = loginFailures.get(key); if ( previous && previous.resetAt > now() && previous.count >= MAX_LOGIN_FAILURES ) { const retryAfter = Math.ceil((previous.resetAt - now()) / 1000); throw new HttpError( 429, "Too many requests", "LOGIN_RATE_LIMITED", "Too many failed login attempts", { "retry-after": String(retryAfter) }, ); } if (previous && previous.resetAt <= now()) loginFailures.delete(key); const body = await readJson(request, 16_384); const username = typeof body.username === "string" ? body.username.trim() : ""; const password = typeof body.password === "string" ? body.password : ""; if (!username || !password) throw new HttpError( 400, "Invalid login request", "INVALID_LOGIN", "Username and password are required", ); const user = await options.store.getUser(username); if ( !user || user.disabled || !(await verifyPassword(password, user.passwordHash)) ) { const failures = loginFailures.get(key); loginFailures.set(key, { count: (failures?.count ?? 0) + 1, resetAt: failures?.resetAt ?? now() + LOGIN_WINDOW_MS, }); throw new HttpError( 401, "Unauthorized", "INVALID_CREDENTIALS", "Invalid username or password", ); } loginFailures.delete(key); const token = createToken(); const expiresAt = new Date( now() + (body.persistent === true ? PERSISTENT_SESSION_MS : RUNTIME_SESSION_MS), ).toISOString(); await options.store.putSession({ tokenHash: hashToken(token), username: user.username, authVersion: user.authVersion, expiresAt, }); return response({ token, expiresAt, user: { username: user.username, roles: user.roles }, }); } async function handleAuthenticated( request: Request, url: URL, identity: Identity, ): Promise { const path = url.pathname; if (request.method === "GET" && path === `${API_PREFIX}/me`) return response({ username: identity.user.username, roles: identity.user.roles, }); if (request.method === "POST" && path === `${API_PREFIX}/logout`) { if (identity.session) await options.store.deleteSession(identity.session.tokenHash); return new Response(null, { status: 204 }); } const trustMatch = new RegExp( `^${API_PREFIX}/workspaces/([^/]+)/trust$`, ).exec(path); if (trustMatch && options.trustStore) { const project = pathPart(trustMatch[1]!); await requireWorkspaceScope(identity, request, project); if (request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); return response({ fingerprints: await options.trustStore.list(project), }); } if (request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (typeof body.fingerprint !== "string") throw new HttpError( 400, "Invalid trust request", "TRUST_INVALID", "fingerprint is required", ); try { validateTrust(project, body.fingerprint); } catch (error) { throw new HttpError( 400, "Invalid trust request", "TRUST_INVALID", error instanceof Error ? error.message : "Invalid trust request", ); } await options.trustStore.grant(project, body.fingerprint); await audit( identity, request, "trust.grant", "success", undefined, project, ); return new Response(null, { status: 204 }); } if (request.method === "DELETE") { await requireCapability(identity, request, "kubernetes:write"); const fingerprint = url.searchParams.get("fingerprint"); if (!fingerprint) throw new HttpError( 400, "Invalid trust request", "TRUST_INVALID", "fingerprint is required", ); if (!(await options.trustStore.revoke(project, fingerprint))) throw new HttpError( 404, "Not found", "TRUST_REGISTRATION_MISSING", "Trust registration was not found", ); await audit( identity, request, "trust.revoke", "success", undefined, project, ); return new Response(null, { status: 204 }); } } if (path === `${API_PREFIX}/builds` && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (typeof body.project !== "string" || !body.project) throw new HttpError( 400, "Invalid build request", "BUILD_INVALID", "project is required", ); await requireWorkspaceScope(identity, request, body.project); return response( await requireBuilds().submitBuild(body as unknown as BuildRequest), 202, ); } if ( path === `${API_PREFIX}/snapshots/negotiate` && request.method === "POST" ) { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (typeof body.project !== "string" || !body.project) throw new HttpError( 400, "Invalid snapshot request", "BUILD_INVALID", "project is required", ); await requireWorkspaceScope(identity, request, body.project); return response( await requireBuilds().negotiateSnapshot(body.workspace as Sha256Digest), ); } if (path === `${API_PREFIX}/images/resolve` && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (typeof body.project !== "string" || typeof body.service !== "string") throw new HttpError( 400, "Invalid request", "BUILD_INVALID", "project and service are required", ); await requireWorkspaceScope(identity, request, body.project); if (!options.resolveImage) throw new HttpError( 503, "Service unavailable", "BUILDS_UNAVAILABLE", "Image resolution is not configured", ); return response(await options.resolveImage(body.project, body.service)); } let buildMatch = new RegExp( `^${API_PREFIX}/builds/([^/]+)(?:/(events|reconcile|cancel|cleanup|result))?$`, ).exec(path); if (buildMatch) { const id = pathPart(buildMatch[1]!); const action = buildMatch[2]; await requireCapability(identity, request, "kubernetes:write"); const builds = requireBuilds(); await requireWorkspaceScope( identity, request, await builds.getBuildProject(id), ); if (!action && request.method === "GET") return response(await builds.getBuildStatus(id)); if (action === "events" && request.method === "GET") 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") return response(await builds.cancelBuild(id)); if (action === "cleanup" && request.method === "DELETE") { await builds.cleanupBuild(id); return new Response(null, { status: 204 }); } if (action === "result" && request.method === "GET") return response(await builds.getBuildResult(id)); } const blobMatch = new RegExp( `^${API_PREFIX}/blobs/(sha256%3A|sha256:)([a-f0-9]{64})/uploads(?:/(complete))?$`, "i", ).exec(path); if (blobMatch) { await requireCapability(identity, request, "kubernetes:write"); const project = url.searchParams.get("project"); if (!project) throw new HttpError( 400, "Invalid upload request", "BUILD_INVALID", "project is required", ); await requireWorkspaceScope(identity, request, project); const digest = `sha256:${blobMatch[2]!.toLowerCase()}` as Sha256Digest; if (blobMatch[3] === "complete" && request.method === "POST") return response(await requireBuilds().completeBlobUpload(digest)); if (request.method === "POST") { const body = await readJson(request); return response( await requireBuilds().beginBlobUpload(digest, Number(body.size)), 201, ); } if (request.method === "PATCH") { const offset = nonNegativeInteger( request.headers.get("upload-offset"), "Upload-Offset", ); if (offset === undefined) throw new HttpError( 400, "Invalid request", "BUILD_INVALID", "Upload-Offset header is required", ); return response( await requireBuilds().uploadBlobChunk( digest, offset, await readBytes(request, 8 * 1024 * 1024), ), ); } } if (path === `${API_PREFIX}/users`) { if (request.method === "GET") { await requireCapability(identity, request, "users:read"); return response({ items: (await options.store.listUsers()).map(publicUser), }); } if (request.method === "POST") { await requireCapability(identity, request, "users:write"); const body = await readJson(request); if ( typeof body.username !== "string" || body.username !== body.username.trim() || !body.username || typeof body.password !== "string" || !body.password || !Array.isArray(body.roles) || body.roles.length === 0 || !body.roles.every(isRole) ) throw new HttpError( 400, "Invalid user", "USER_INVALID", "Username, password, and roles are required", ); const user = await options.store.createUser({ username: body.username, passwordHash: await hashPassword(body.password), roles: body.roles as Role[], }); await audit(identity, request, "user.create", "success", { username: user.username, }); return response(publicUser(user), 201, { location: `${API_PREFIX}/users/${encodeURIComponent(user.username)}`, }); } } let match = new RegExp( `^${API_PREFIX}/users/([^/]+)(?:/(sessions/revoke))?$`, ).exec(path); const keyMatch = new RegExp( `^${API_PREFIX}/users/([^/]+)/keys(?:/([^/]+))?$`, ).exec(path); if (keyMatch) { const username = pathPart(keyMatch[1]!); const keyId = keyMatch[2] && pathPart(keyMatch[2]); if (!keyId && request.method === "GET") { await requireCapability(identity, request, "users:read"); return response({ items: (await options.store.listApiKeys(username)).map((key) => ({ id: key.id, username: key.username, capabilities: key.capabilities, ...(key.workspace && { workspace: key.workspace }), expiresAt: key.expiresAt, disabled: Boolean(key.disabled), })), }); } if (!keyId && request.method === "POST") { await requireCapability(identity, request, "users:write"); const body = await readJson(request); const expiresAt = body.expiresAt === undefined ? new Date(now() + DEFAULT_API_KEY_MS).toISOString() : typeof body.expiresAt === "string" ? body.expiresAt : ""; const expires = new Date(expiresAt); if ( !Array.isArray(body.capabilities) || body.capabilities.length === 0 || new Set(body.capabilities).size !== body.capabilities.length || !body.capabilities.every(isCapability) || (body.workspace !== undefined && (typeof body.workspace !== "string" || !/^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?$/.test(body.workspace) || body.workspace.length > 63)) || !Number.isFinite(expires.getTime()) || expires.toISOString() !== expiresAt || expires.getTime() <= now() || expires.getTime() > now() + MAX_API_KEY_MS ) throw new HttpError( 400, "Invalid API key", "API_KEY_INVALID", "Capabilities and an expiry no more than 365 days away are required", ); await requireApiKeyDelegation( identity, request, username, body.capabilities as Capability[], typeof body.workspace === "string" && body.workspace ? body.workspace : undefined, ); const token = createToken(); const key: ApiKeyRecord = { id: randomUUID(), tokenHash: hashToken(token), username, capabilities: body.capabilities, ...(typeof body.workspace === "string" && body.workspace && { workspace: body.workspace }), expiresAt, }; if (!(await options.store.getUser(username))) throw new HttpError( 404, "Not found", "USER_NOT_FOUND", "User not found", ); await options.store.createApiKey(key); await audit(identity, request, "api_key.create", "success", { username, keyId: key.id, capabilities: key.capabilities, ...(key.workspace && { workspace: key.workspace }), expiresAt: key.expiresAt, }); return response( { id: key.id, username, capabilities: key.capabilities, ...(key.workspace && { workspace: key.workspace }), expiresAt: key.expiresAt, disabled: false, token, }, 201, ); } if (keyId && request.method === "DELETE") { await requireCapability(identity, request, "users:write"); if (!(await options.store.revokeApiKey(username, keyId))) throw new HttpError( 404, "Not found", "API_KEY_NOT_FOUND", "API key not found", ); await audit(identity, request, "api_key.revoke", "success", { username, keyId, }); return new Response(null, { status: 204 }); } } if (match) { const username = pathPart(match[1]!); if (match[2]) { if (request.method !== "POST") throw new HttpError( 405, "Method not allowed", "METHOD_NOT_ALLOWED", "This endpoint only accepts POST", { allow: "POST" }, ); await requireCapability(identity, request, "sessions:revoke"); const revoked = await options.store.revokeUserSessions(username); await audit(identity, request, "sessions.revoke", "success", { username, revoked, }); return response({ username, revoked }); } if (request.method === "GET") { await requireCapability(identity, request, "users:read"); const user = await options.store.getUser(username); if (!user) throw new HttpError( 404, "Not found", "USER_NOT_FOUND", `User '${username}' not found`, ); return response(publicUser(user)); } if (request.method === "PATCH" || request.method === "PUT") { await requireCapability(identity, request, "users:write"); const body = await readJson(request); if ( (body.roles !== undefined && (!Array.isArray(body.roles) || body.roles.length === 0 || !body.roles.every(isRole))) || (body.password !== undefined && (typeof body.password !== "string" || !body.password)) || (body.disabled !== undefined && typeof body.disabled !== "boolean") ) throw new HttpError( 400, "Invalid user", "USER_INVALID", "Invalid user update fields", ); const updated = await options.store.updateUser(username, { ...(typeof body.password === "string" && body.password && { passwordHash: await hashPassword(body.password), }), ...(Array.isArray(body.roles) && { roles: body.roles as Role[] }), ...(typeof body.disabled === "boolean" && { disabled: body.disabled, }), }); if (!updated) throw new HttpError( 404, "Not found", "USER_NOT_FOUND", `User '${username}' not found`, ); const revoked = await options.store.revokeUserSessions(username); await audit(identity, request, "user.update", "success", { username, revoked, }); return response(publicUser(updated)); } if (request.method === "DELETE") { await requireCapability(identity, request, "users:write"); if (!(await options.store.deleteUser(username))) throw new HttpError( 404, "Not found", "USER_NOT_FOUND", `User '${username}' not found`, ); await audit(identity, request, "user.delete", "success", { username }); return new Response(null, { status: 204 }); } } if (path === `${API_PREFIX}/workspaces`) { if (request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); const workspaces = await requireWorkspaceStore().list(); return response({ items: identity.apiKey?.workspace ? workspaces.filter( (workspace) => workspace.metadata.name === identity.apiKey?.workspace, ) : workspaces, }); } if (request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (typeof body.id !== "string" || !body.id) throw new HttpError( 400, "Invalid workspace", "WORKSPACE_INVALID", "Workspace id is required", ); await requireWorkspaceScope(identity, request, body.id); const created = await requireWorkspaceStore().create( body as unknown as CreateWorkspaceInput, ); await audit( identity, request, "workspace.create", "success", undefined, created.metadata.name, ); return response(created, 201, { etag: workspaceEtag(created), location: `${API_PREFIX}/workspaces/${created.metadata.name}`, }); } } if ( path === `${API_PREFIX}/platform/kuber-system/adopt` && request.method === "POST" ) { await requireCapability(identity, request, "platform:adopt"); if (identity.apiKey && identity.apiKey.workspace !== "kuber-system") { await audit(identity, request, "authorization.workspace", "denied", { workspace: "kuber-system", }); throw new HttpError( 403, "Forbidden", "FORBIDDEN", "Platform adoption requires an API key scoped to kuber-system", ); } const body = await readJson(request); if (typeof body.workspaceUid !== "string" || !body.workspaceUid.trim()) throw new HttpError( 400, "Invalid request", "ADOPTION_INVALID", "workspaceUid is required", ); const adopted = await requireAdoption().adoptPlatform(body.workspaceUid); await audit( identity, request, "platform.adopt", "success", undefined, "kuber-system", ); return response(adopted); } match = new RegExp(`^${API_PREFIX}/maintenance/([^/]+)$`).exec(path); if (match) { const host = pathPart(match[1]!); try { normalizeMaintenanceHost(host); } catch (error) { throw new HttpError( 400, "Invalid hostname", "MAINTENANCE_HOST_INVALID", error instanceof Error ? error.message : "Invalid hostname", ); } if (!options.maintenance) throw new HttpError( 503, "Service unavailable", "MAINTENANCE_UNAVAILABLE", "Maintenance service is not configured", ); if (identity.apiKey?.workspace) { await audit(identity, request, "authorization.workspace", "denied", { workspace: identity.apiKey.workspace, }); throw new HttpError( 403, "Forbidden", "MAINTENANCE_GLOBAL_SCOPE_REQUIRED", "Maintenance requires an unscoped API key or an operator session", ); } if (request.method === "GET") { await requireCapability(identity, request, "kubernetes:write"); return response(await options.maintenance.status(host)); } if (request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (typeof body.enabled !== "boolean") throw new HttpError( 400, "Invalid request", "MAINTENANCE_ENABLED_REQUIRED", "enabled must be a boolean", ); const result = await options.maintenance.set(host, body.enabled); await audit(identity, request, "maintenance.toggle", "success", result); return response(result); } } match = new RegExp(`^${API_PREFIX}/workspaces/([^/]+)(?:/(.*))?$`).exec( path, ); if (match) { const id = pathPart(match[1]!); const subpath = match[2]; await requireWorkspaceScope(identity, request, id); if (!subpath) { if (request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); const workspace = await getWorkspace(id); return response(workspace, 200, { etag: workspaceEtag(workspace) }); } if (request.method === "PUT" || request.method === "PATCH") { await requireCapability(identity, request, "kubernetes:write"); const ifMatch = request.headers.get("if-match"); if (!ifMatch) throw new HttpError( 428, "Precondition required", "IF_MATCH_REQUIRED", "If-Match header is required", ); const updated = await requireWorkspaceStore().update( id, (await readJson(request)) as unknown as UpdateWorkspaceInput, ifMatch, ); await audit( identity, request, "workspace.update", "success", undefined, id, ); return response(updated, 200, { etag: workspaceEtag(updated) }); } if (request.method === "DELETE") { await requireCapability(identity, request, "kubernetes:write"); const workspace = await getWorkspace(id); return runOperation( identity, request, workspace, "workspace.delete", { full: true }, async (signal) => { const result = options.management ? await options.management.down( workspaceIdentity(workspace), true, { signal }, ) : undefined; if (signal.aborted) throw new Error("Workspace operation execution was cancelled"); if (!requireWorkspaceStore().delete) throw new Error("Workspace store does not support deletion"); await requireWorkspaceStore().delete!(id); return result ?? { deleted: true }; }, ); } } const workspace = await getWorkspace(id); if (subpath === "adopt" && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); let adopted: WorkspaceAdoptionResult; try { adopted = await requireAdoption().adopt( workspace.metadata.name, workspace.metadata.uid, ); } catch (error) { try { await audit( identity, request, "workspace.adopt", "failure", { route: path, requestId: requestIds.get(request) ?? makeRequestId(), }, id, ); } catch (auditError) { logRequestError(request, auditError, { event: "workspace.adopt.audit.failed", code: "AUDIT_APPEND_FAILED", message: "Failed workspace adoption audit event could not be appended", workspaceId: id, }); } throw error; } await audit( identity, request, "workspace.adopt", "success", undefined, id, ); return response(adopted); } if (subpath === "status" && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); return response( await requireManagement().graphStatus(workspaceIdentity(workspace), { includeIdle: url.searchParams.get("includeIdle") === "true", }), ); } if (subpath === "logs" && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); if (!options.logs) throw new HttpError( 503, "Service unavailable", "LOGS_UNAVAILABLE", "Log service is not configured", ); const service = url.searchParams.get("service")?.trim(); const logOptions = { namespace: workspace.metadata.name, target: service ? ({ kind: "service", name: service } as const) : ({ kind: "managed-deployments" } as const), tailLines: nonNegativeInteger( url.searchParams.get("tailLines"), "tailLines", ), sinceSeconds: nonNegativeInteger( url.searchParams.get("sinceSeconds"), "sinceSeconds", ), sinceTime: url.searchParams.get("sinceTime") ?? undefined, timestamps: url.searchParams.get("timestamps") === "true", }; if (url.searchParams.get("follow") !== "true") return ndjson(await options.logs.collect(logOptions, request.signal)); return ndjson( options.logs.follow({ ...logOptions, signal: request.signal }), request.signal, ); } if (subpath === "revisions" && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); return response({ items: await requireWorkspaceStore().listRevisions(id), }); } const revisionMatch = /^revisions\/(\d+)$/.exec(subpath ?? ""); if (revisionMatch && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); const revision = await requireWorkspaceStore().getRevision( id, Number(revisionMatch[1]), ); if (!revision) throw new HttpError( 404, "Not found", "REVISION_NOT_FOUND", "Workspace revision not found", ); return response(revision); } if ( ["lifecycle", "stop", "restart", "rollback", "down"].includes( subpath ?? "", ) && request.method === "POST" ) { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); const lifecycleAction = subpath === "lifecycle" && (body.action === "stop" || body.action === "restart") ? body.action : subpath; if (subpath === "lifecycle" && lifecycleAction === "lifecycle") throw new HttpError( 400, "Invalid lifecycle action", "LIFECYCLE_INVALID", "action must be stop or restart", ); return runOperation( identity, request, workspace, `workspace.${lifecycleAction}`, body, async (signal) => { const management = requireManagement(); const names = Array.isArray(body.services) ? (body.services as string[]) : undefined; if (lifecycleAction === "stop") return management.stop(workspaceIdentity(workspace), names, { signal, }); if (lifecycleAction === "restart") return management.restart(workspaceIdentity(workspace), names, { signal, }); if (lifecycleAction === "rollback") return management.rollback( workspaceIdentity(workspace), names, typeof body.timeoutMs === "number" ? body.timeoutMs : undefined, { signal }, ); return management.down( workspaceIdentity(workspace), body.full === true, { signal }, ); }, ); } if (subpath === "resources/plan" && request.method === "POST") { await requireCapability(identity, request, "kubernetes:read"); const body = await readJson(request); if (!Array.isArray(body.resources)) throw new HttpError( 400, "Invalid resources", "RESOURCES_INVALID", "resources must be an array", ); return response( await requireManagement().planResources( workspaceIdentity(workspace), body.resources as KubernetesObject[], ), ); } if ( ["resources/apply", "resources/wait", "resources/delete"].includes( subpath ?? "", ) && request.method === "POST" ) { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); if (subpath === "resources/apply") await requireTrust(request, id); return runOperation( identity, request, workspace, subpath!, body, async (signal, emit) => { const management = requireManagement(); if (subpath === "resources/apply") { if (!Array.isArray(body.resources)) throw new Error("resources must be an array"); return management.applyResources( workspaceIdentity(workspace), body.resources as KubernetesObject[], { signal, emit }, ); } if (subpath === "resources/wait") { if (!Array.isArray(body.deployments)) throw new Error("deployments must be an array"); await management.waitForResources( workspaceIdentity(workspace), body.deployments as string[], typeof body.timeoutMs === "number" ? body.timeoutMs : undefined, { signal, emit }, ); return { ready: true }; } if (!Array.isArray(body.resources)) throw new Error("resources must be an array"); await management.deleteResources( workspaceIdentity(workspace), body.resources as unknown as ResourceIdentity[], { signal, emit }, ); return { deleted: body.resources.length }; }, ); } if (subpath === "databases" && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); const compose = (body.compose ?? body) as ComposeSpecification; return runOperation( identity, request, workspace, "databases.reconcile", body, (signal) => requireManagement().reconcileDatabases( workspaceIdentity(workspace), compose, { signal }, ), ); } const databaseCredentials = /^databases\/([^/]+)\/credentials$/.exec( subpath ?? "", ); if (databaseCredentials && request.method === "GET") { await requireCapability(identity, request, "kubernetes:write"); return response( await requireManagement().getDatabaseCredentials( workspaceIdentity(workspace), pathPart(databaseCredentials[1]!), ), ); } if (subpath === "databases/credentials" && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); const claim = isRecord(body.claim) ? body.claim : body; if (typeof claim.username !== "string") throw new HttpError( 400, "Invalid database claim", "DATABASE_CLAIM_INVALID", "claim.username is required", ); return response( await requireManagement().getDatabaseCredentials( workspaceIdentity(workspace), claim.username, ), ); } if ( subpath === "databases/credentials-metadata" && request.method === "POST" ) { await requireCapability(identity, request, "kubernetes:read"); const body = await readJson(request); return response({ items: requireManagement().databaseCredentialsMetadata( (body.compose ?? body) as ComposeSpecification, ), }); } if (subpath === "storage" && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); const compose = (body.compose ?? body) as ComposeSpecification; return runOperation( identity, request, workspace, "storage.reconcile", body, (signal) => requireManagement().reconcileStorage( workspaceIdentity(workspace), compose, { signal }, ), ); } if (subpath === "storage/credentials" && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); const body = await readJson(request); const claim = isRecord(body.claim) ? body.claim : body; return response( await requireManagement().getStorageCredentials( workspaceIdentity(workspace), claim as unknown as { service: string; key: string; bucket: string; }, ), ); } if ( subpath === "storage/credentials-metadata" && request.method === "POST" ) { await requireCapability(identity, request, "kubernetes:read"); const body = await readJson(request); return response({ items: requireManagement().storageCredentialsMetadata( (body.compose ?? body) as ComposeSpecification, ), }); } } if (path === `${API_PREFIX}/operations` && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); const requestedWorkspace = url.searchParams.get("workspaceId") ?? undefined; if ( identity.apiKey?.workspace && requestedWorkspace && requestedWorkspace !== identity.apiKey.workspace ) await requireWorkspaceScope(identity, request, requestedWorkspace); if (!options.operationStore) throw new HttpError( 503, "Service unavailable", "OPERATION_STORE_UNAVAILABLE", "Operation storage is not configured", ); return response({ items: ( await options.operationStore.list( identity.apiKey?.workspace ?? requestedWorkspace, ) ).map(publicOperation), }); } match = new RegExp(`^${API_PREFIX}/operations/([^/]+)/events$`).exec(path); if (match && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); const operation = await options.operationStore?.get(pathPart(match[1]!)); if (!operation) throw new HttpError( 404, "Not found", "OPERATION_NOT_FOUND", "Operation not found", ); await requireWorkspaceScope( identity, request, operation.spec.workspaceId, ); const afterValue = url.searchParams.get("after") ?? "0"; const after = Number(afterValue); if (!Number.isSafeInteger(after) || after < 0) throw new HttpError( 400, "Invalid cursor", "OPERATION_EVENTS_CURSOR_INVALID", "after must be a non-negative integer", ); const events = await options.operationStore!.events( operation.metadata.name, after, ); return response({ ...events, ...(events.items.length > 0 && { nextCursor: events.items.at(-1)!.sequence, }), }); } match = new RegExp(`^${API_PREFIX}/operations/([^/]+)$`).exec(path); if (match && request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); const operation = await options.operationStore?.get(pathPart(match[1]!)); if (!operation) throw new HttpError( 404, "Not found", "OPERATION_NOT_FOUND", "Operation not found", ); await requireWorkspaceScope( identity, request, operation.spec.workspaceId, ); return response(publicOperation(operation)); } if (path === `${API_PREFIX}/audit` && request.method === "GET") { await requireCapability(identity, request, "users:read"); if (!options.auditStore) throw new HttpError( 503, "Service unavailable", "AUDIT_STORE_UNAVAILABLE", "Audit storage is not configured", ); const requestedWorkspace = url.searchParams.get("workspaceId") ?? undefined; if ( identity.apiKey?.workspace && requestedWorkspace && requestedWorkspace !== identity.apiKey.workspace ) await requireWorkspaceScope(identity, request, requestedWorkspace); return response({ items: await options.auditStore.list( identity.apiKey?.workspace ?? requestedWorkspace, ), }); } throw new HttpError( 404, "Not found", "NOT_FOUND", "The requested API endpoint does not exist", ); } function normalizeError(error: unknown): HttpError { if (error instanceof HttpError) return error; const kubernetes = kubernetesError(error); if (kubernetes) return new HttpError( kubernetes.status, "Kubernetes error", "KUBERNETES_ERROR", kubernetes.detail, ); if (error instanceof WorkspaceNotFoundError) return new HttpError(404, "Not found", error.code, error.message); if (error instanceof OperationNotFoundError) return new HttpError(404, "Not found", error.code, error.message); if (error instanceof BuildNotFoundError) return new HttpError(404, "Not found", error.code, error.message); if ( error instanceof WorkspaceConflictError || error instanceof OperationConflictError || error instanceof BuildConflictError || error instanceof WorkspaceAdoptionError || error instanceof MaintenanceBusyError ) return new HttpError( 409, "Conflict", error.code, error.message, error instanceof MaintenanceBusyError ? { "retry-after": "1" } : undefined, ); if ( error instanceof WorkspaceValidationError || error instanceof OperationValidationError || error instanceof BuildValidationError ) return new HttpError(400, "Invalid request", error.code, error.message); if (error instanceof KubernetesLogError) return new HttpError( 502, "Kubernetes log error", "LOGS_FAILED", error.message, ); if (error instanceof Error && /already exists/i.test(error.message)) return new HttpError(409, "Conflict", "ALREADY_EXISTS", error.message); return new HttpError( 500, "Internal server error", "INTERNAL_ERROR", "The request could not be completed", ); } return async (request) => { const suppliedRequestId = request.headers.get("x-request-id"); const requestId = suppliedRequestId && /^[\x21-\x7e]{1,128}$/.test(suppliedRequestId) ? suppliedRequestId : makeRequestId(); return logServerRequest( request, requestId, async () => { requestIds.set(request, requestId); const url = new URL(request.url); const origin = request.headers.get("origin"); const originAllowed = !origin || (hasOriginConfiguration ? allowedOrigins.has(origin) : origin === url.origin); let result: Response; try { if (!originAllowed) throw new HttpError( 403, "Forbidden", "ORIGIN_NOT_ALLOWED", "Request origin is not allowed", ); if (request.method === "OPTIONS") { if (!origin) throw new HttpError( 400, "Bad request", "ORIGIN_REQUIRED", "Origin header is required for preflight", ); result = new Response(null, { status: 204, headers: { "access-control-allow-methods": "GET, POST, PUT, PATCH, DELETE, OPTIONS", "access-control-allow-headers": "Authorization, Content-Type, Idempotency-Key, If-Match, Upload-Offset, X-Request-Id", "access-control-max-age": "600", }, }); } else if ( request.method === "GET" && url.pathname === `${API_PREFIX}/health` ) { result = response({ status: "ok" }); } else if ( request.method === "POST" && url.pathname === `${API_PREFIX}/login` ) { result = await handleLogin(request); } else { const identity = await authenticate(request); if (!identity) throw new HttpError( 401, "Unauthorized", "UNAUTHORIZED", "A valid kuber login is required", { "www-authenticate": "Bearer" }, ); result = await handleAuthenticated(request, url, identity); } } catch (error) { const normalized = normalizeError(error); if (normalized.code === "INTERNAL_ERROR") logRequestError(request, error, { event: "request.failed", code: normalized.code, message: normalized.message, workspaceId: workspaceIdFromPath(url.pathname), }); result = problem(normalized, requestId); } result.headers.set("x-request-id", requestId); if (origin && originAllowed) { result.headers.set("access-control-allow-origin", origin); result.headers.set("vary", "Origin"); } return result; }, { log: (entry) => (logger.log ?? processLogger.log)(entry) }, ); }; } const EXEC_UPGRADE_PATH = new RegExp(`^${API_PREFIX}/workspaces/([^/]+)/exec$`); export function execUpgradeMatch(url: URL): string | undefined { const match = EXEC_UPGRADE_PATH.exec(url.pathname); return match?.[1]; } export type ExecConnection = { identity: Identity; workspace: Workspace; }; export type ExecAuthOptions = Pick< AppOptions, "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); if (!identity) throw new HttpError( 401, "Unauthorized", "UNAUTHORIZED", "A valid kuber login is required", { "www-authenticate": "Bearer" }, ); if ( identity.apiKey ? !identity.apiKey.capabilities.includes("kubernetes:exec") || (identity.apiKey.workspace !== undefined && identity.apiKey.workspace !== workspaceId) : !hasCapability(identity.user.roles, "kubernetes:exec") ) { if (options.auditStore) { try { await options.auditStore.append({ actor: { username: identity.user.username, roles: identity.user.roles, ip: request.headers.get("cf-connecting-ip") ?? request.headers.get("x-forwarded-for")?.split(",")[0]?.trim(), userAgent: request.headers.get("user-agent") ?? undefined, }, action: "authorization.kubernetes:exec", outcome: "denied", workspaceId, }); } catch (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( 403, "Forbidden", "FORBIDDEN", "Capability 'kubernetes:exec' is required", ); } const workspaceStore = options.workspaceStore; const workspace = workspaceStore ? await workspaceStore.get(pathPart(workspaceId)) : undefined; if (!workspace) throw new HttpError( 404, "Not found", "WORKSPACE_NOT_FOUND", `Workspace '${workspaceId}' not found`, ); return { identity, workspace }; } export function execProblem(error: unknown, requestId = ""): Response { let httpError: HttpError; if (error instanceof HttpError) httpError = error; else { httpError = new HttpError( 500, "Internal server error", "INTERNAL_ERROR", "The request could not be completed", ); } const result = Response.json( { type: `https://kuber.astrxl.dev/problems/${httpError.code.toLowerCase()}`, title: httpError.title, status: httpError.status, detail: httpError.message, code: httpError.code, requestId, ...(httpError.operationId && { operationId: httpError.operationId }), }, { status: httpError.status, headers: { "cache-control": "no-store", ...httpError.headers, }, }, ); result.headers.set("content-type", "application/problem+json"); return result; } export function execWireExit( server: ExecServerFrame, ): ExecServerWireFrame | undefined { if (server.type !== "exit" && server.type !== "error") return; if (server.type === "error") return { type: "error", code: server.code, message: server.message }; return { type: "exit", exitCode: server.exitCode, ...(server.reason !== undefined && { reason: server.reason }), ...(server.message !== undefined && { message: server.message }), }; } function isStartFrame(value: ExecClientWireFrame): value is ExecStartFrame { return (value as ExecStartFrame).type === "start"; } function base64ToBytes(value: string): Uint8Array | undefined { const bytes = Uint8Array.from(Buffer.from(value, "base64")); return bytes; } export type ExecWebSocketLink = { sendText(data: string): void; close(code?: number, reason?: string): void; }; export type ExecLinkOptions = { maxFrameBytes?: number; }; const DEFAULT_EXEC_MAX_FRAME_BYTES = 64 * 1024; export class WireExecSession { private readonly controller = new AbortController(); private session?: Awaited>; private readonly maxFrameBytes: number; private started = false; private done: Promise = Promise.resolve(); private closed = false; constructor( private readonly socket: ExecWebSocketLink, private readonly execService: ExecService, private readonly connection: ExecConnection, options: ExecLinkOptions = {}, ) { this.maxFrameBytes = options.maxFrameBytes ?? DEFAULT_EXEC_MAX_FRAME_BYTES; } private send(frame: ExecServerWireFrame) { if (this.closed) return; this.socket.sendText(JSON.stringify(frame)); } async receive(raw: unknown): Promise { if (typeof raw !== "string") { this.fail("EXEC_INVALID", "Exec frames must be JSON text", 1003); return; } if (new TextEncoder().encode(raw).byteLength > this.maxFrameBytes) { this.fail("EXEC_INVALID", "Exec frame exceeds the size limit", 1009); return; } let frame: ExecClientWireFrame; try { frame = JSON.parse(raw) as ExecClientWireFrame; } catch { this.fail("EXEC_INVALID", "Exec frame is not valid JSON", 1003); return; } if (!frame || typeof frame !== "object" || !("type" in frame)) { this.fail("EXEC_INVALID", "Invalid exec frame", 1003); return; } if (!this.started) { if (!isStartFrame(frame)) { this.fail("EXEC_INVALID", "First exec frame must be start", 1003); return; } await this.start(frame); return; } if (!this.session) return; try { await this.session.send(this.toClientFrame(frame)); } catch (error) { if (!this.controller.signal.aborted) { this.send({ type: "error", code: "EXEC_INVALID", message: error instanceof Error ? error.message : "Invalid exec frame", }); } } } private toClientFrame(frame: ExecClientWireFrame): ExecClientFrame { switch (frame.type) { case "stdin": { if (frame.encoding !== "base64" || typeof frame.data !== "string") throw new Error("Invalid stdin frame"); const data = base64ToBytes(frame.data); if (!data) throw new Error("Invalid stdin frame"); return { type: "stdin", data, ...(frame.eof ? { eof: true } : {}), }; } case "resize": { if ( !Number.isSafeInteger(frame.columns) || !Number.isSafeInteger(frame.rows) || frame.columns < 1 || frame.rows < 1 || frame.columns > 65_535 || frame.rows > 65_535 ) throw new Error("Terminal dimensions are invalid"); return { type: "resize", columns: frame.columns, rows: frame.rows }; } case "close": return { type: "close" }; default: throw new Error("Unknown exec frame"); } } private fail(code: string, message: string, closeCode: number) { if (this.closed) return; this.send({ type: "error", code, message }); this.close(closeCode, message); } close(code = 1000, reason = "closed") { if (this.closed) return; this.closed = true; this.controller.abort(); void this.closeSession(); this.socket.close(code, reason); } private async closeSession() { try { await this.session?.close(); } catch { // session already closed } this.session = undefined; } private async start(frame: ExecStartFrame) { if (frame.version !== 1) { this.fail( "EXEC_INVALID", `Unsupported exec protocol version ${frame.version}`, 1003, ); return; } if ( typeof frame.deployment !== "string" || !frame.deployment || !Array.isArray(frame.command) || frame.command.length === 0 ) { this.fail("EXEC_INVALID", "deployment and command are required", 1003); return; } const workspace = this.connection.workspace; try { this.session = await this.execService.openInteractive({ workspace: { project: workspace.metadata.name, uid: workspace.metadata.uid, }, deployment: frame.deployment, container: frame.container, command: frame.command, tty: frame.tty ?? true, signal: this.controller.signal, }); } catch (error) { if (this.controller.signal.aborted) return; this.send({ type: "error", code: error instanceof Error && "code" in error && typeof (error as { code?: unknown }).code === "string" ? (error as { code: string }).code : "EXEC_FAILED", message: error instanceof Error ? error.message : "Exec failed to start", }); this.close(1011, "exec failed"); return; } this.started = true; this.done = this.pump(); } private async pump() { const session = this.session; if (!session) return; try { for await (const server of session) { if (this.closed || this.controller.signal.aborted) return; if (server.type === "stdout" || server.type === "stderr") { this.send({ type: server.type, data: Buffer.from(server.data).toString("base64"), encoding: "base64", }); } else { const wire = execWireExit(server); if (wire) { this.send(wire); this.close(1000, "process finished"); return; } } } } catch (error) { if (!this.controller.signal.aborted && !this.closed) { const code = error instanceof Error && "code" in error && typeof (error as { code?: unknown }).code === "string" ? (error as { code: string }).code : "EXEC_FAILED"; this.send({ type: "error", code, message: error instanceof Error ? error.message : "Exec session failed", }); this.close(1011, "exec failed"); } } } } export async function handleExecUpgrade( options: AppOptions, request: Request, 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, requestId, ); } catch (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( new HttpError( 400, "Bad request", "UPGRADE_FAILED", "WebSocket upgrade failed", ), requestId, ); return; }